-
Notifications
You must be signed in to change notification settings - Fork 25.6k
Expose merge events and their memory usage estimate #126667
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
Merged
Changes from 10 commits
Commits
Show all changes
17 commits
Select commit
Hold shift + click to select a range
ca095be
draft
pxsalehi 643979c
Expose merge events and their memory usage estimate
pxsalehi a724f42
Merge remote-tracking branch 'upstream/main' into ps250410-exposeMerg…
pxsalehi ffea088
Merge remote-tracking branch 'upstream/main' into ps250410-exposeMerg…
pxsalehi 74610ad
tests
pxsalehi 0c389e2
Merge remote-tracking branch 'upstream/main' into ps250410-exposeMerg…
pxsalehi 8fbb8b5
tiny change
pxsalehi 649a5bf
Merge remote-tracking branch 'upstream/main' into ps250410-exposeMerg…
pxsalehi ffea307
Merge remote-tracking branch 'upstream/main' into ps250410-exposeMerg…
pxsalehi f86c86b
comment
pxsalehi 3724a21
Ensure merge complete/abort listeners are not called before merge que…
pxsalehi b236f12
Merge remote-tracking branch 'upstream/main' into ps250410-exposeMerg…
pxsalehi 74f3b6b
rename
pxsalehi 5395c0a
fix test
pxsalehi ba7ebb5
call listener before adding to queue
pxsalehi ae030b5
Merge remote-tracking branch 'upstream/main' into ps250410-exposeMerg…
pxsalehi 2508ce6
Merge remote-tracking branch 'upstream/main' into ps250410-exposeMerg…
pxsalehi File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
26 changes: 26 additions & 0 deletions
26
server/src/main/java/org/elasticsearch/index/engine/MergeEventListener.java
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,26 @@ | ||
| /* | ||
| * Copyright Elasticsearch B.V. and/or licensed to Elasticsearch B.V. under one | ||
| * or more contributor license agreements. Licensed under the "Elastic License | ||
| * 2.0", the "GNU Affero General Public License v3.0 only", and the "Server Side | ||
| * Public License v 1"; you may not use this file except in compliance with, at | ||
| * your election, the "Elastic License 2.0", the "GNU Affero General Public | ||
| * License v3.0 only", or the "Server Side Public License, v 1". | ||
| */ | ||
|
|
||
| package org.elasticsearch.index.engine; | ||
|
|
||
| import org.elasticsearch.index.merge.OnGoingMerge; | ||
|
|
||
| public interface MergeEventListener { | ||
|
|
||
| /** | ||
| * | ||
| * @param merge | ||
| * @param estimateMergeMemoryBytes estimate of the memory needed to perform a merge | ||
| */ | ||
| void onMergeQueued(OnGoingMerge merge, long estimateMergeMemoryBytes); | ||
|
|
||
| void onMergeCompleted(OnGoingMerge merge); | ||
|
|
||
| void onMergeAborted(OnGoingMerge merge); | ||
| } |
21 changes: 21 additions & 0 deletions
21
server/src/main/java/org/elasticsearch/index/engine/MergeMemoryEstimateProvider.java
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,21 @@ | ||
| /* | ||
| * Copyright Elasticsearch B.V. and/or licensed to Elasticsearch B.V. under one | ||
| * or more contributor license agreements. Licensed under the "Elastic License | ||
| * 2.0", the "GNU Affero General Public License v3.0 only", and the "Server Side | ||
| * Public License v 1"; you may not use this file except in compliance with, at | ||
| * your election, the "Elastic License 2.0", the "GNU Affero General Public | ||
| * License v3.0 only", or the "Server Side Public License, v 1". | ||
| */ | ||
|
|
||
| package org.elasticsearch.index.engine; | ||
|
|
||
| import org.apache.lucene.index.MergePolicy; | ||
|
|
||
| @FunctionalInterface | ||
| public interface MergeMemoryEstimateProvider { | ||
|
|
||
| /** | ||
| * Returns an estimate of the memory needed to perform a merge | ||
| */ | ||
| long estimateMergeMemoryBytes(MergePolicy.OneMerge merge); | ||
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -65,11 +65,13 @@ public class ThreadPoolMergeScheduler extends MergeScheduler implements Elastics | |
| private final AtomicLong doneMergeTaskCount = new AtomicLong(); | ||
| private final CountDownLatch closedWithNoRunningMerges = new CountDownLatch(1); | ||
| private volatile boolean closed = false; | ||
| private final MergeMemoryEstimateProvider mergeMemoryEstimateProvider; | ||
|
|
||
| public ThreadPoolMergeScheduler( | ||
| ShardId shardId, | ||
| IndexSettings indexSettings, | ||
| ThreadPoolMergeExecutorService threadPoolMergeExecutorService | ||
| ThreadPoolMergeExecutorService threadPoolMergeExecutorService, | ||
| MergeMemoryEstimateProvider mergeMemoryEstimateProvider | ||
| ) { | ||
| this.shardId = shardId; | ||
| this.config = indexSettings.getMergeSchedulerConfig(); | ||
|
|
@@ -81,6 +83,7 @@ public ThreadPoolMergeScheduler( | |
| : Double.POSITIVE_INFINITY | ||
| ); | ||
| this.threadPoolMergeExecutorService = threadPoolMergeExecutorService; | ||
| this.mergeMemoryEstimateProvider = mergeMemoryEstimateProvider; | ||
| } | ||
|
|
||
| @Override | ||
|
|
@@ -176,11 +179,13 @@ MergeTask newMergeTask(MergeSource mergeSource, MergePolicy.OneMerge merge, Merg | |
| // forced merges, as well as merges triggered when closing a shard, always run un-IO-throttled | ||
| boolean isAutoThrottle = mergeTrigger != MergeTrigger.CLOSING && merge.getStoreMergeInfo().mergeMaxNumSegments() == -1; | ||
| // IO throttling cannot be toggled for existing merge tasks, only new merge tasks pick up the updated IO throttling setting | ||
| long estimateMergeMemoryBytes = mergeMemoryEstimateProvider.estimateMergeMemoryBytes(merge); | ||
| return new MergeTask( | ||
| mergeSource, | ||
| merge, | ||
| isAutoThrottle && config.isAutoThrottle(), | ||
| "Lucene Merge Task #" + submittedMergeTaskCount.incrementAndGet() + " for shard " + shardId | ||
| "Lucene Merge Task #" + submittedMergeTaskCount.incrementAndGet() + " for shard " + shardId, | ||
| estimateMergeMemoryBytes | ||
| ); | ||
| } | ||
|
|
||
|
|
@@ -312,14 +317,22 @@ class MergeTask implements Runnable { | |
| private final OnGoingMerge onGoingMerge; | ||
| private final MergeRateLimiter rateLimiter; | ||
| private final boolean supportsIOThrottling; | ||
|
|
||
| MergeTask(MergeSource mergeSource, MergePolicy.OneMerge merge, boolean supportsIOThrottling, String name) { | ||
| private final long estimateMergeMemoryBytes; | ||
|
|
||
| MergeTask( | ||
| MergeSource mergeSource, | ||
| MergePolicy.OneMerge merge, | ||
| boolean supportsIOThrottling, | ||
| String name, | ||
| long estimateMergeMemoryBytes | ||
| ) { | ||
| this.name = name; | ||
| this.mergeStartTimeNS = new AtomicLong(); | ||
| this.mergeSource = mergeSource; | ||
| this.onGoingMerge = new OnGoingMerge(merge); | ||
| this.rateLimiter = new MergeRateLimiter(merge.getMergeProgress()); | ||
| this.supportsIOThrottling = supportsIOThrottling; | ||
| this.estimateMergeMemoryBytes = estimateMergeMemoryBytes; | ||
| } | ||
|
|
||
| Schedule schedule() { | ||
|
|
@@ -449,6 +462,14 @@ long estimatedMergeSize() { | |
| return onGoingMerge.getMerge().getStoreMergeInfo().estimatedMergeBytes(); | ||
| } | ||
|
|
||
| public long getEstimateMergeMemoryBytes() { | ||
|
||
| return estimateMergeMemoryBytes; | ||
| } | ||
|
|
||
| public OnGoingMerge getOnGoingMerge() { | ||
| return onGoingMerge; | ||
| } | ||
|
|
||
| @Override | ||
| public String toString() { | ||
| return name + (onGoingMerge.getMerge().isAborted() ? " (aborted)" : ""); | ||
|
|
||
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Moved this out of stateless.