Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -500,13 +500,7 @@ public DataFormatAwareEngine(EngineConfig engineConfig) {
refreshLock.unlock();
}
}
},
this::activateThrottling,
this::deactivateThrottling,
shardId,
engineConfig.getIndexSettings(),
engineConfig.getThreadPool()
);
}, shardId, engineConfig.getIndexSettings(), engineConfig.getThreadPool());
success = true;
logger.trace("created new DataFormatBasedEngine");
} catch (IOException | TranslogCorruptedException e) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -48,11 +48,8 @@ public class MergeScheduler {
private final MergeHandler mergeHandler;
private final BiConsumer<MergeResult, OneMerge> applyMergeChanges;
private final Runnable onMergeFailureCleanup;
private final Runnable activateThrottling;
private final Runnable deactivateThrottling;
private final ThreadPool threadPool;
private final AtomicInteger activeMerges = new AtomicInteger(0);
private final AtomicBoolean isThrottling = new AtomicBoolean(false);
private final AtomicBoolean isShutdown = new AtomicBoolean(false);
private final Semaphore forceMergeLock = new Semaphore(1);
private final AtomicBoolean frozen = new AtomicBoolean(false);
Expand All @@ -78,8 +75,6 @@ public class MergeScheduler {
* @param mergeHandler the handler that selects and executes merges
* @param applyMergeChanges callback to apply merge results (e.g., update the catalog)
* @param onMergeFailureCleanup callback invoked when a merge fails and cleanup is performed
* @param activateThrottling callback to activate indexing throttle when merge pressure is high
* @param deactivateThrottling callback to deactivate indexing throttle when merge pressure subsides
* @param shardId the shard this scheduler is associated with
* @param indexSettings the index settings providing merge scheduler configuration
* @param threadPool the OpenSearch thread pool for executing merge tasks
Expand All @@ -88,17 +83,13 @@ public MergeScheduler(
MergeHandler mergeHandler,
BiConsumer<MergeResult, OneMerge> applyMergeChanges,
Runnable onMergeFailureCleanup,
Runnable activateThrottling,
Runnable deactivateThrottling,
ShardId shardId,
IndexSettings indexSettings,
ThreadPool threadPool
) {
this.mergeHandler = mergeHandler;
this.applyMergeChanges = applyMergeChanges;
this.onMergeFailureCleanup = onMergeFailureCleanup;
this.activateThrottling = activateThrottling;
this.deactivateThrottling = deactivateThrottling;
this.threadPool = threadPool;
logger = Loggers.getLogger(getClass(), shardId);
this.indexSettings = indexSettings;
Expand Down Expand Up @@ -147,7 +138,6 @@ public void triggerMerges() {
if (!isFrozen()) {
mergeHandler.findAndRegisterMerges();
}
evaluateThrottle();
executeMerge();
}

Expand Down Expand Up @@ -353,7 +343,6 @@ private void submitMergeTask(OneMerge oneMerge) {
// uncaught exception on the merge thread pool.
} finally {
activeMerges.decrementAndGet();
evaluateThrottle();
// Fire all drain listeners if all merges completed and none pending
if (isFrozen() && activeMerges.get() == 0 && !mergeHandler.hasPendingMerges() && !onDrainedListeners.isEmpty()) {
List<Runnable> listeners = List.copyOf(onDrainedListeners);
Expand Down Expand Up @@ -405,27 +394,4 @@ private void runMerge(OneMerge oneMerge) throws IOException {
mergeStatsTracker.afterMerge(tookMS, totalNumDocs, totalSizeInBytes);
}
}

private synchronized void evaluateThrottle() {
int numMergesInFlight = activeMerges.get() + mergeHandler.getPendingMergeCount();
if (numMergesInFlight > maxMergeCount) {
if (isThrottling.getAndSet(true) == false) {
logger.info("now throttling indexing: numMergesInFlight={}, maxMergeCount={}", numMergesInFlight, maxMergeCount);
try {
activateThrottling.run();
} catch (Exception e) {
logger.warn("exception in activateThrottling callback", e);
}
}
} else if (numMergesInFlight < maxMergeCount) {
if (isThrottling.getAndSet(false)) {
logger.info("stop throttling indexing: numMergesInFlight={}, maxMergeCount={}", numMergesInFlight, maxMergeCount);
try {
deactivateThrottling.run();
} catch (Exception e) {
logger.warn("exception in deactivateThrottling callback", e);
}
}
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -67,16 +67,7 @@ public void testOnDrained_AlreadyDrained_FiresListenerImmediately() {

IndexSettings indexSettings = IndexSettingsModule.newIndexSettings("test", Settings.EMPTY);
ShardId testShardId = new ShardId(indexSettings.getIndex(), 0);
MergeScheduler scheduler = new MergeScheduler(
mockHandler,
(result, merge) -> {},
() -> {},
() -> {},
() -> {},
testShardId,
indexSettings,
threadPool
);
MergeScheduler scheduler = new MergeScheduler(mockHandler, (result, merge) -> {}, () -> {}, testShardId, indexSettings, threadPool);

AtomicBoolean listenerCalled = new AtomicBoolean(false);
scheduler.onDrained(() -> listenerCalled.set(true));
Expand All @@ -93,16 +84,7 @@ public void testOnDrained_MergesPending_RegistersListener() {

IndexSettings indexSettings = IndexSettingsModule.newIndexSettings("test", Settings.EMPTY);
ShardId testShardId = new ShardId(indexSettings.getIndex(), 0);
MergeScheduler scheduler = new MergeScheduler(
mockHandler,
(result, merge) -> {},
() -> {},
() -> {},
() -> {},
testShardId,
indexSettings,
threadPool
);
MergeScheduler scheduler = new MergeScheduler(mockHandler, (result, merge) -> {}, () -> {}, testShardId, indexSettings, threadPool);

AtomicBoolean listenerCalled = new AtomicBoolean(false);
scheduler.onDrained(() -> listenerCalled.set(true));
Expand All @@ -119,16 +101,7 @@ public void testOnDrained_MultipleListeners_AllFire() {

IndexSettings indexSettings = IndexSettingsModule.newIndexSettings("test", Settings.EMPTY);
ShardId testShardId = new ShardId(indexSettings.getIndex(), 0);
MergeScheduler scheduler = new MergeScheduler(
mockHandler,
(result, merge) -> {},
() -> {},
() -> {},
() -> {},
testShardId,
indexSettings,
threadPool
);
MergeScheduler scheduler = new MergeScheduler(mockHandler, (result, merge) -> {}, () -> {}, testShardId, indexSettings, threadPool);

AtomicInteger callCount = new AtomicInteger(0);

Expand Down Expand Up @@ -204,16 +177,7 @@ public void testOnDrained_ListenersFire_WhenMergesGoFromNToZero() throws Excepti

IndexSettings indexSettings = IndexSettingsModule.newIndexSettings("test", Settings.EMPTY);
ShardId testShardId = new ShardId(indexSettings.getIndex(), 0);
MergeScheduler scheduler = new MergeScheduler(
mockHandler,
(result, merge) -> {},
() -> {},
() -> {},
() -> {},
testShardId,
indexSettings,
threadPool
);
MergeScheduler scheduler = new MergeScheduler(mockHandler, (result, merge) -> {}, () -> {}, testShardId, indexSettings, threadPool);

CountDownLatch latch = new CountDownLatch(3);
AtomicInteger callCount = new AtomicInteger(0);
Expand Down Expand Up @@ -259,16 +223,7 @@ public void testOnDrained_DoubleCheckRace_ListenerFiresImmediately() throws Exce

IndexSettings indexSettings = IndexSettingsModule.newIndexSettings("test", Settings.EMPTY);
ShardId testShardId = new ShardId(indexSettings.getIndex(), 0);
MergeScheduler scheduler = new MergeScheduler(
mockHandler,
(result, merge) -> {},
() -> {},
() -> {},
() -> {},
testShardId,
indexSettings,
threadPool
);
MergeScheduler scheduler = new MergeScheduler(mockHandler, (result, merge) -> {}, () -> {}, testShardId, indexSettings, threadPool);

// First call to hasPendingMerges returns true (first check fails, goes to add),
// second call returns false (double-check succeeds, fires the listener)
Expand All @@ -293,16 +248,7 @@ public void testOnDrained_ListenerExceptionIsolation_OtherListenersStillFire() t

IndexSettings indexSettings = IndexSettingsModule.newIndexSettings("test", Settings.EMPTY);
ShardId testShardId = new ShardId(indexSettings.getIndex(), 0);
MergeScheduler scheduler = new MergeScheduler(
mockHandler,
(result, merge) -> {},
() -> {},
() -> {},
() -> {},
testShardId,
indexSettings,
threadPool
);
MergeScheduler scheduler = new MergeScheduler(mockHandler, (result, merge) -> {}, () -> {}, testShardId, indexSettings, threadPool);

// Register listeners individually via the double-check path.
// Each onDrained call is independent — if one listener throws during its own
Expand Down Expand Up @@ -339,16 +285,7 @@ public void testHasPendingMerges_DelegatesToMergeHandler() {

IndexSettings indexSettings = IndexSettingsModule.newIndexSettings("test", Settings.EMPTY);
ShardId testShardId = new ShardId(indexSettings.getIndex(), 0);
MergeScheduler scheduler = new MergeScheduler(
mockHandler,
(result, merge) -> {},
() -> {},
() -> {},
() -> {},
testShardId,
indexSettings,
threadPool
);
MergeScheduler scheduler = new MergeScheduler(mockHandler, (result, merge) -> {}, () -> {}, testShardId, indexSettings, threadPool);

assertTrue("hasPendingMerges should delegate to handler when handler reports true", scheduler.hasPendingMerges());

Expand All @@ -365,16 +302,7 @@ public void testGetActiveMergeCount_InitiallyZero() {

IndexSettings indexSettings = IndexSettingsModule.newIndexSettings("test", Settings.EMPTY);
ShardId testShardId = new ShardId(indexSettings.getIndex(), 0);
MergeScheduler scheduler = new MergeScheduler(
mockHandler,
(result, merge) -> {},
() -> {},
() -> {},
() -> {},
testShardId,
indexSettings,
threadPool
);
MergeScheduler scheduler = new MergeScheduler(mockHandler, (result, merge) -> {}, () -> {}, testShardId, indexSettings, threadPool);

assertEquals("Active merge count should be 0 initially", 0, scheduler.getActiveMergeCount());
}
Expand Down Expand Up @@ -408,16 +336,7 @@ public void testSubmitMergeTask_FinallyBlock_FiresListenersWhenLastMergeComplete

IndexSettings indexSettings = IndexSettingsModule.newIndexSettings("test", Settings.EMPTY);
ShardId testShardId = new ShardId(indexSettings.getIndex(), 0);
MergeScheduler scheduler = new MergeScheduler(
mockHandler,
(result, merge) -> {},
() -> {},
() -> {},
() -> {},
testShardId,
indexSettings,
threadPool
);
MergeScheduler scheduler = new MergeScheduler(mockHandler, (result, merge) -> {}, () -> {}, testShardId, indexSettings, threadPool);

// Register an onDrained listener BEFORE triggering merges.
// Reset hasPendingMerges stub for the full flow:
Expand Down Expand Up @@ -497,7 +416,7 @@ private MergeScheduler newIdleScheduler(MergeHandler mockHandler) {
when(mockHandler.hasPendingMerges()).thenReturn(false);
IndexSettings indexSettings = IndexSettingsModule.newIndexSettings("test", Settings.EMPTY);
ShardId testShardId = new ShardId(indexSettings.getIndex(), 0);
return new MergeScheduler(mockHandler, (result, merge) -> {}, () -> {}, () -> {}, () -> {}, testShardId, indexSettings, threadPool);
return new MergeScheduler(mockHandler, (result, merge) -> {}, () -> {}, testShardId, indexSettings, threadPool);
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,6 @@
import org.opensearch.core.index.shard.ShardId;
import org.opensearch.index.IndexModule;
import org.opensearch.index.IndexSettings;
import org.opensearch.index.MergeSchedulerConfig;
import org.opensearch.index.engine.dataformat.MergeResult;
import org.opensearch.index.engine.exec.Segment;
import org.opensearch.test.IndexSettingsModule;
Expand All @@ -26,9 +25,7 @@
import java.io.IOException;
import java.util.List;
import java.util.Map;
import java.util.concurrent.atomic.AtomicInteger;

import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyInt;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
Expand Down Expand Up @@ -70,16 +67,7 @@ private IndexSettings indexSettings(IndexModule.TieringState tieringState) {
}

private MergeScheduler newScheduler(MergeHandler mergeHandler, IndexModule.TieringState tieringState) {
return new MergeScheduler(
mergeHandler,
(result, merge) -> {},
() -> {},
() -> {},
() -> {},
shardId,
indexSettings(tieringState),
threadPool
);
return new MergeScheduler(mergeHandler, (result, merge) -> {}, () -> {}, shardId, indexSettings(tieringState), threadPool);
}

public void testFreezeBlocksTriggerMerges() {
Expand Down Expand Up @@ -169,8 +157,6 @@ public void testForceMergeAbortsRemainingMergesOnShutdown() throws Exception {
mergeHandler,
(result, merge) -> { schedulerRef.get().shutdown(); },
() -> {},
() -> {},
() -> {},
shardId,
indexSettings(IndexModule.TieringState.HOT),
threadPool
Expand All @@ -188,82 +174,4 @@ public void testForceMergeAbortsRemainingMergesOnShutdown() throws Exception {
verify(mergeHandler).doMerge(merge1);
verify(mergeHandler, never()).doMerge(merge2);
}

public void testThrottlingActivatesWhenMergesExceedMaxCount() throws Exception {
AtomicInteger activateCount = new AtomicInteger();
AtomicInteger deactivateCount = new AtomicInteger();

IndexSettings idxSettings = IndexSettingsModule.newIndexSettings(
"test",
Settings.builder()
.put(IndexMetadata.SETTING_VERSION_CREATED, Version.CURRENT)
.put(MergeSchedulerConfig.MAX_THREAD_COUNT_SETTING.getKey(), "1")
.put(MergeSchedulerConfig.MAX_MERGE_COUNT_SETTING.getKey(), "2")
.build()
);

MergeHandler mergeHandler = mock(MergeHandler.class);
OneMerge merge1 = mock(OneMerge.class);
OneMerge merge2 = mock(OneMerge.class);
OneMerge merge3 = mock(OneMerge.class);
when(merge1.getSegmentsToMerge()).thenReturn(List.of());
when(merge2.getSegmentsToMerge()).thenReturn(List.of());
when(merge3.getSegmentsToMerge()).thenReturn(List.of());

when(mergeHandler.getPendingMergeCount()).thenReturn(3, 3, 2, 1, 0);
when(mergeHandler.hasPendingMerges()).thenReturn(true, true, true, false);
when(mergeHandler.getNextMerge()).thenReturn(merge1).thenReturn(merge2).thenReturn(merge3).thenReturn(null);
when(mergeHandler.doMerge(any())).thenReturn(new MergeResult(Map.of()));

MergeScheduler scheduler = new MergeScheduler(
mergeHandler,
(result, merge) -> {},
() -> {},
activateCount::incrementAndGet,
deactivateCount::incrementAndGet,
shardId,
idxSettings,
threadPool
);

scheduler.triggerMerges();
assertBusy(() -> assertTrue("throttle should have activated", activateCount.get() > 0));
assertBusy(() -> assertTrue("throttle should have deactivated", deactivateCount.get() > 0));
}

public void testThrottlingNotActivatedWhenMergesWithinLimit() throws Exception {
AtomicInteger activateCount = new AtomicInteger();

IndexSettings idxSettings = IndexSettingsModule.newIndexSettings(
"test",
Settings.builder()
.put(IndexMetadata.SETTING_VERSION_CREATED, Version.CURRENT)
.put(MergeSchedulerConfig.MAX_THREAD_COUNT_SETTING.getKey(), "1")
.put(MergeSchedulerConfig.MAX_MERGE_COUNT_SETTING.getKey(), "6")
.build()
);

MergeHandler mergeHandler = mock(MergeHandler.class);
OneMerge merge1 = mock(OneMerge.class);
when(merge1.getSegmentsToMerge()).thenReturn(List.of());
when(mergeHandler.getPendingMergeCount()).thenReturn(1, 0);
when(mergeHandler.hasPendingMerges()).thenReturn(true, false);
when(mergeHandler.getNextMerge()).thenReturn(merge1).thenReturn(null);
when(mergeHandler.doMerge(any())).thenReturn(new MergeResult(Map.of()));

MergeScheduler scheduler = new MergeScheduler(
mergeHandler,
(result, merge) -> {},
() -> {},
activateCount::incrementAndGet,
() -> {},
shardId,
idxSettings,
threadPool
);

scheduler.triggerMerges();
Thread.sleep(200);
assertEquals("throttle should not activate when merges within limit", 0, activateCount.get());
}
}
Loading
Loading