From 9bf510767ea8b641da0de94c0a874e3373d6f27a Mon Sep 17 00:00:00 2001 From: Shourya Dutta Biswas <114977491+shourya035@users.noreply.github.com> Date: Fri, 21 Aug 2026 15:47:42 +0530 Subject: [PATCH] Restore size-based remote translog flushes Track remote translog bytes independently of retained generations. Add a dynamic, disabled-by-default cluster setting for opt-in rollout. Limit tracking to Lucene-based engines and bypass the new machinery while disabled. Reset tracking only after successful index commits. Signed-off-by: Shourya Dutta Biswas <114977491+shourya035@users.noreply.github.com> --- .../common/settings/ClusterSettings.java | 1 + .../index/engine/InternalEngine.java | 19 ++- .../translog/InternalTranslogManager.java | 75 ++++++++- .../index/translog/RemoteFsTranslog.java | 6 + .../indices/RemoteStoreSettings.java | 26 +++ .../InternalTranslogManagerTests.java | 156 ++++++++++++++++++ ...RemoteStoreSettingsDynamicUpdateTests.java | 16 ++ 7 files changed, 296 insertions(+), 3 deletions(-) diff --git a/server/src/main/java/org/opensearch/common/settings/ClusterSettings.java b/server/src/main/java/org/opensearch/common/settings/ClusterSettings.java index 77108dedb4245..6877b7e3f4014 100644 --- a/server/src/main/java/org/opensearch/common/settings/ClusterSettings.java +++ b/server/src/main/java/org/opensearch/common/settings/ClusterSettings.java @@ -867,6 +867,7 @@ public void apply(Settings value, Settings current, Settings previous) { RemoteStoreSettings.CLUSTER_REMOTE_STORE_PATH_TYPE_SETTING, RemoteStoreSettings.CLUSTER_REMOTE_STORE_PATH_HASH_ALGORITHM_SETTING, RemoteStoreSettings.CLUSTER_REMOTE_MAX_TRANSLOG_READERS, + RemoteStoreSettings.CLUSTER_REMOTE_TRANSLOG_TRACK_BYTES_SINCE_LAST_COMMIT_SETTING, RemoteStoreSettings.CLUSTER_REMOTE_STORE_TRANSLOG_METADATA, RemoteStoreSettings.CLUSTER_REMOTE_STORE_PINNED_TIMESTAMP_SCHEDULER_INTERVAL, RemoteStoreSettings.CLUSTER_REMOTE_STORE_PINNED_TIMESTAMP_LOOKBACK_INTERVAL, diff --git a/server/src/main/java/org/opensearch/index/engine/InternalEngine.java b/server/src/main/java/org/opensearch/index/engine/InternalEngine.java index 997161ffb60c2..f10555246d667 100644 --- a/server/src/main/java/org/opensearch/index/engine/InternalEngine.java +++ b/server/src/main/java/org/opensearch/index/engine/InternalEngine.java @@ -409,7 +409,8 @@ protected TranslogManager createTranslogManager( this::ensureOpen, engineConfig.getTranslogFactory(), engineConfig.getStartedPrimarySupplier(), - TranslogOperationHelper.create(engineConfig) + TranslogOperationHelper.create(engineConfig), + true ); } @@ -2222,7 +2223,21 @@ protected void commitIndexWriter(final DocumentIndexWriter writer, final String refresh("commit", SearcherScope.INTERNAL, true); } - writer.commit(); + final InternalTranslogManager internalTranslogManager = translogManager instanceof InternalTranslogManager manager + ? manager + : null; + if (internalTranslogManager != null && internalTranslogManager.isTranslogBytesTrackingEnabled()) { + internalTranslogManager.startIndexCommit(); + boolean commitSuccessful = false; + try { + writer.commit(); + commitSuccessful = true; + } finally { + internalTranslogManager.finishIndexCommit(commitSuccessful); + } + } else { + writer.commit(); + } } catch (final Exception ex) { try { failEngine("lucene commit failed", ex); diff --git a/server/src/main/java/org/opensearch/index/translog/InternalTranslogManager.java b/server/src/main/java/org/opensearch/index/translog/InternalTranslogManager.java index 031fba0403cc7..b51fa2e9d3167 100644 --- a/server/src/main/java/org/opensearch/index/translog/InternalTranslogManager.java +++ b/server/src/main/java/org/opensearch/index/translog/InternalTranslogManager.java @@ -46,6 +46,9 @@ public class InternalTranslogManager implements TranslogManager { private final TranslogEventListener translogEventListener; private final Supplier localCheckpointTrackerSupplier; private final Logger logger; + private final TranslogBytesTracker translogBytesTracker = new TranslogBytesTracker(); + private final boolean supportsTranslogBytesTracking; + private TranslogBytesTracker.CommitSnapshot commitSnapshot; public InternalTranslogManager( TranslogConfig translogConfig, @@ -61,6 +64,40 @@ public InternalTranslogManager( TranslogFactory translogFactory, BooleanSupplier startedPrimarySupplier, TranslogOperationHelper translogOperationHelper + ) throws IOException { + this( + translogConfig, + primaryTermSupplier, + globalCheckpointSupplier, + translogDeletionPolicy, + shardId, + readLock, + localCheckpointTrackerSupplier, + translogUUID, + translogEventListener, + engineLifeCycleAware, + translogFactory, + startedPrimarySupplier, + translogOperationHelper, + false + ); + } + + public InternalTranslogManager( + TranslogConfig translogConfig, + LongSupplier primaryTermSupplier, + LongSupplier globalCheckpointSupplier, + TranslogDeletionPolicy translogDeletionPolicy, + ShardId shardId, + ReleasableLock readLock, + Supplier localCheckpointTrackerSupplier, + String translogUUID, + TranslogEventListener translogEventListener, + LifecycleAware engineLifeCycleAware, + TranslogFactory translogFactory, + BooleanSupplier startedPrimarySupplier, + TranslogOperationHelper translogOperationHelper, + boolean supportsTranslogBytesTracking ) throws IOException { this.shardId = shardId; this.readLock = readLock; @@ -76,6 +113,7 @@ public InternalTranslogManager( }, translogUUID, translogFactory, startedPrimarySupplier, translogOperationHelper); assert translog.getGeneration() != null; this.translog = translog; + this.supportsTranslogBytesTracking = supportsTranslogBytesTracking && translog instanceof RemoteFsTranslog; assert pendingTranslogRecovery.get() == false : "translog recovery can't be pending before we set it"; // don't allow commits until we are done with recovering pendingTranslogRecovery.set(true); @@ -337,7 +375,39 @@ public Translog.Operation readOperation(Translog.Location location) throws IOExc */ @Override public Translog.Location add(Translog.Operation operation) throws IOException { - return translog.add(operation); + Translog.Location location = translog.add(operation); + if (isTranslogBytesTrackingEnabled()) { + translogBytesTracker.addBytes(location.size); + } + return location; + } + + /** + * Returns whether byte tracking is supported by the engine and enabled for the remote translog. + */ + public boolean isTranslogBytesTrackingEnabled() { + return supportsTranslogBytesTracking && ((RemoteFsTranslog) translog).isTranslogBytesTrackingEnabled(); + } + + /** + * Captures the remote translog bytes that are eligible to be cleared by the next index commit. + */ + public void startIndexCommit() { + assert commitSnapshot == null : "an index commit is already in progress"; + commitSnapshot = translogBytesTracker.startCommit(); + } + + /** + * Completes byte tracking for an index commit. + * + * @param successful whether the index commit completed successfully + */ + public void finishIndexCommit(boolean successful) { + assert commitSnapshot != null : "index commit was not started"; + if (successful) { + translogBytesTracker.completeCommit(commitSnapshot); + } + commitSnapshot = null; } /** @@ -451,6 +521,9 @@ public boolean shouldPeriodicallyFlush(long localCheckpointOfLastCommit, long fl if (translog.shouldFlush()) { return true; } + if (isTranslogBytesTrackingEnabled()) { + return translogBytesTracker.getBytesSinceLastCommit() >= flushThreshold; + } // This is the minimum seqNo that is referred in translog and considered for calculating translog size long minTranslogRefSeqNo = translog.getMinUnreferencedSeqNoInSegments(localCheckpointOfLastCommit + 1); final long minReferencedTranslogGeneration = translog.getMinGenerationForSeqNo(minTranslogRefSeqNo).translogFileGeneration; diff --git a/server/src/main/java/org/opensearch/index/translog/RemoteFsTranslog.java b/server/src/main/java/org/opensearch/index/translog/RemoteFsTranslog.java index 2f99903403168..7f469bb9b4d91 100644 --- a/server/src/main/java/org/opensearch/index/translog/RemoteFsTranslog.java +++ b/server/src/main/java/org/opensearch/index/translog/RemoteFsTranslog.java @@ -73,6 +73,7 @@ public class RemoteFsTranslog extends Translog { protected final FileTransferTracker fileTransferTracker; protected final BooleanSupplier startedPrimarySupplier; private final RemoteTranslogTransferTracker remoteTranslogTransferTracker; + private final RemoteStoreSettings remoteStoreSettings; private volatile long maxRemoteTranslogGenerationUploaded; private volatile long minSeqNoToKeep; @@ -128,6 +129,7 @@ public RemoteFsTranslog( logger = Loggers.getLogger(getClass(), shardId); this.startedPrimarySupplier = startedPrimarySupplier; this.remoteTranslogTransferTracker = remoteTranslogTransferTracker; + this.remoteStoreSettings = remoteStoreSettings; fileTransferTracker = new FileTransferTracker(shardId, remoteTranslogTransferTracker); isTranslogMetadataEnabled = indexSettings().isTranslogMetadataEnabled(); this.isServerSideEncryptionEnabled = isServerSideEncryptionEnabled; @@ -795,6 +797,10 @@ protected boolean shouldFlush() { return readers.size() >= maxRemoteTlogReaders; } + boolean isTranslogBytesTrackingEnabled() { + return remoteStoreSettings.isTranslogBytesTrackingEnabled(); + } + private CryptoMetadata resolveCryptoMetadata() { IndexMetadata indexMetadata = indexSettings.getIndexMetadata(); if (indexMetadata == null) { diff --git a/server/src/main/java/org/opensearch/indices/RemoteStoreSettings.java b/server/src/main/java/org/opensearch/indices/RemoteStoreSettings.java index d6c6692fb641c..693c760efec91 100644 --- a/server/src/main/java/org/opensearch/indices/RemoteStoreSettings.java +++ b/server/src/main/java/org/opensearch/indices/RemoteStoreSettings.java @@ -126,6 +126,17 @@ public class RemoteStoreSettings { Property.NodeScope ); + /** + * Controls whether remote translog bytes written since the last index commit are used to evaluate + * {@code index.translog.flush_threshold_size}. + */ + public static final Setting CLUSTER_REMOTE_TRANSLOG_TRACK_BYTES_SINCE_LAST_COMMIT_SETTING = Setting.boolSetting( + "cluster.remote_store.translog.track_bytes_since_last_commit.enabled", + false, + Property.Dynamic, + Property.NodeScope + ); + /** * Controls timeout value while uploading segment files to remote segment store */ @@ -227,6 +238,7 @@ public class RemoteStoreSettings { private volatile RemoteStoreEnums.PathType pathType; private volatile RemoteStoreEnums.PathHashAlgorithm pathHashAlgorithm; private volatile int maxRemoteTranslogReaders; + private volatile boolean translogBytesTrackingEnabled; private volatile boolean isClusterServerSideEncryptionRepoEnabled; private volatile boolean isTranslogMetadataEnabled; private static volatile boolean isPinnedTimestampsEnabled; @@ -267,6 +279,12 @@ public RemoteStoreSettings(Settings settings, ClusterSettings clusterSettings) { maxRemoteTranslogReaders = CLUSTER_REMOTE_MAX_TRANSLOG_READERS.get(settings); clusterSettings.addSettingsUpdateConsumer(CLUSTER_REMOTE_MAX_TRANSLOG_READERS, this::setMaxRemoteTranslogReaders); + translogBytesTrackingEnabled = CLUSTER_REMOTE_TRANSLOG_TRACK_BYTES_SINCE_LAST_COMMIT_SETTING.get(settings); + clusterSettings.addSettingsUpdateConsumer( + CLUSTER_REMOTE_TRANSLOG_TRACK_BYTES_SINCE_LAST_COMMIT_SETTING, + this::setTranslogBytesTrackingEnabled + ); + clusterRemoteSegmentTransferTimeout = CLUSTER_REMOTE_SEGMENT_TRANSFER_TIMEOUT_SETTING.get(settings); clusterSettings.addSettingsUpdateConsumer( CLUSTER_REMOTE_SEGMENT_TRANSFER_TIMEOUT_SETTING, @@ -356,6 +374,14 @@ private void setMaxRemoteTranslogReaders(int maxRemoteTranslogReaders) { this.maxRemoteTranslogReaders = maxRemoteTranslogReaders; } + public boolean isTranslogBytesTrackingEnabled() { + return translogBytesTrackingEnabled; + } + + private void setTranslogBytesTrackingEnabled(boolean translogBytesTrackingEnabled) { + this.translogBytesTrackingEnabled = translogBytesTrackingEnabled; + } + public boolean isClusterServerSideEncryptionEnabled() { return isClusterServerSideEncryptionRepoEnabled; } diff --git a/server/src/test/java/org/opensearch/index/translog/InternalTranslogManagerTests.java b/server/src/test/java/org/opensearch/index/translog/InternalTranslogManagerTests.java index 237bbe1b7648c..e7636c4823a67 100644 --- a/server/src/test/java/org/opensearch/index/translog/InternalTranslogManagerTests.java +++ b/server/src/test/java/org/opensearch/index/translog/InternalTranslogManagerTests.java @@ -23,10 +23,15 @@ import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.locks.ReentrantReadWriteLock; +import java.util.function.BooleanSupplier; +import java.util.function.LongConsumer; +import java.util.function.LongSupplier; import static org.opensearch.index.seqno.SequenceNumbers.NO_OPS_PERFORMED; import static org.opensearch.index.translog.TranslogDeletionPolicies.createTranslogDeletionPolicy; import static org.hamcrest.Matchers.equalTo; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; public class InternalTranslogManagerTests extends TranslogManagerTestCase { @@ -362,4 +367,155 @@ public void testAcquireHistoryRetentionLock() throws IOException { translogManager.close(); } } + + public void testRemoteTranslogBytesControlPeriodicFlush() throws IOException { + RemoteFsTranslog remoteTranslog = mockRemoteTranslog(true); + when(remoteTranslog.add(org.mockito.ArgumentMatchers.any())).thenReturn(new Translog.Location(1, 0, 100)); + + try (InternalTranslogManager translogManager = createTranslogManager(remoteTranslog, true)) { + translogManager.add(mock(Translog.Operation.class)); + + assertFalse(translogManager.shouldPeriodicallyFlush(0, 101)); + assertTrue(translogManager.shouldPeriodicallyFlush(0, 100)); + + when(remoteTranslog.isTranslogBytesTrackingEnabled()).thenReturn(false); + assertFalse(translogManager.shouldPeriodicallyFlush(0, 100)); + } + } + + public void testRemoteTranslogBytesResetOnlyAfterSuccessfulCommit() throws IOException { + RemoteFsTranslog remoteTranslog = mockRemoteTranslog(true); + when(remoteTranslog.add(org.mockito.ArgumentMatchers.any())).thenReturn( + new Translog.Location(1, 0, 100), + new Translog.Location(1, 100, 20), + new Translog.Location(1, 120, 30) + ); + + try (InternalTranslogManager translogManager = createTranslogManager(remoteTranslog, true)) { + Translog.Operation operation = mock(Translog.Operation.class); + translogManager.add(operation); + translogManager.startIndexCommit(); + translogManager.add(operation); + translogManager.finishIndexCommit(true); + + assertFalse(translogManager.shouldPeriodicallyFlush(0, 21)); + assertTrue(translogManager.shouldPeriodicallyFlush(0, 20)); + + translogManager.add(operation); + translogManager.startIndexCommit(); + translogManager.finishIndexCommit(false); + + assertTrue(translogManager.shouldPeriodicallyFlush(0, 50)); + } + } + + public void testMaxRemoteTranslogReadersRemainsIndependentFlushTrigger() throws IOException { + RemoteFsTranslog remoteTranslog = mockRemoteTranslog(false); + when(remoteTranslog.shouldFlush()).thenReturn(true); + + try (InternalTranslogManager translogManager = createTranslogManager(remoteTranslog)) { + assertTrue(translogManager.shouldPeriodicallyFlush(0, Long.MAX_VALUE)); + } + } + + public void testRemoteTranslogBytesAreNotTrackedWhenDisabled() throws IOException { + RemoteFsTranslog remoteTranslog = mockRemoteTranslog(false); + when(remoteTranslog.add(org.mockito.ArgumentMatchers.any())).thenReturn( + new Translog.Location(1, 0, 100), + new Translog.Location(1, 100, 20) + ); + + try (InternalTranslogManager translogManager = createTranslogManager(remoteTranslog, true)) { + Translog.Operation operation = mock(Translog.Operation.class); + translogManager.add(operation); + + when(remoteTranslog.isTranslogBytesTrackingEnabled()).thenReturn(true); + assertFalse(translogManager.shouldPeriodicallyFlush(0, 100)); + + translogManager.add(operation); + assertTrue(translogManager.shouldPeriodicallyFlush(0, 20)); + } + } + + public void testRemoteTranslogBytesAreNotTrackedForUnsupportedEngine() throws IOException { + RemoteFsTranslog remoteTranslog = mockRemoteTranslog(true); + when(remoteTranslog.add(org.mockito.ArgumentMatchers.any())).thenReturn(new Translog.Location(1, 0, 100)); + + try (InternalTranslogManager translogManager = createTranslogManager(remoteTranslog)) { + translogManager.add(mock(Translog.Operation.class)); + + assertFalse(translogManager.isTranslogBytesTrackingEnabled()); + assertFalse(translogManager.shouldPeriodicallyFlush(0, 100)); + } + } + + private RemoteFsTranslog mockRemoteTranslog(boolean bytesTrackingEnabled) { + RemoteFsTranslog remoteTranslog = mock(RemoteFsTranslog.class); + Translog.TranslogGeneration generation = new Translog.TranslogGeneration(translogUUID, 1); + when(remoteTranslog.getGeneration()).thenReturn(generation); + when(remoteTranslog.isTranslogBytesTrackingEnabled()).thenReturn(bytesTrackingEnabled); + when(remoteTranslog.getMinUnreferencedSeqNoInSegments(org.mockito.ArgumentMatchers.anyLong())).thenReturn(0L); + when(remoteTranslog.getMinGenerationForSeqNo(org.mockito.ArgumentMatchers.anyLong())).thenReturn(generation); + when(remoteTranslog.sizeInBytesByMinGen(org.mockito.ArgumentMatchers.anyLong())).thenReturn(0L); + return remoteTranslog; + } + + private InternalTranslogManager createTranslogManager(Translog translog) throws IOException { + return createTranslogManager(translog, false); + } + + private InternalTranslogManager createTranslogManager(Translog translog, boolean supportsTranslogBytesTracking) throws IOException { + LocalCheckpointTracker tracker = new LocalCheckpointTracker(NO_OPS_PERFORMED, NO_OPS_PERFORMED); + return new InternalTranslogManager( + new TranslogConfig(shardId, primaryTranslogDir, INDEX_SETTINGS, BigArrays.NON_RECYCLING_INSTANCE, "", false), + primaryTerm, + () -> NO_OPS_PERFORMED, + createTranslogDeletionPolicy(INDEX_SETTINGS), + shardId, + new ReleasableLock(new ReentrantReadWriteLock().readLock()), + () -> tracker, + translogUUID, + TranslogEventListener.NOOP_TRANSLOG_EVENT_LISTENER, + () -> {}, + new StubTranslogFactory(translog), + () -> true, + TranslogOperationHelper.DEFAULT, + supportsTranslogBytesTracking + ); + } + + private static class StubTranslogFactory implements TranslogFactory { + private final Translog translog; + + StubTranslogFactory(Translog translog) { + this.translog = translog; + } + + @Override + public Translog newTranslog( + TranslogConfig config, + String translogUUID, + TranslogDeletionPolicy deletionPolicy, + LongSupplier globalCheckpointSupplier, + LongSupplier primaryTermSupplier, + LongConsumer persistedSequenceNumberConsumer, + BooleanSupplier startedPrimarySupplier + ) { + return translog; + } + + @Override + public Translog newTranslog( + TranslogConfig config, + String translogUUID, + TranslogDeletionPolicy deletionPolicy, + LongSupplier globalCheckpointSupplier, + LongSupplier primaryTermSupplier, + LongConsumer persistedSequenceNumberConsumer, + BooleanSupplier startedPrimarySupplier, + TranslogOperationHelper translogOperationHelper + ) { + return translog; + } + } } diff --git a/server/src/test/java/org/opensearch/indices/RemoteStoreSettingsDynamicUpdateTests.java b/server/src/test/java/org/opensearch/indices/RemoteStoreSettingsDynamicUpdateTests.java index 3e2fbc15408af..e4a58ed8f428a 100644 --- a/server/src/test/java/org/opensearch/indices/RemoteStoreSettingsDynamicUpdateTests.java +++ b/server/src/test/java/org/opensearch/indices/RemoteStoreSettingsDynamicUpdateTests.java @@ -128,6 +128,22 @@ public void testDisableMaxRemoteReferencedTranslogFiles() { assertEquals(-1, remoteStoreSettings.getMaxRemoteTranslogReaders()); } + public void testTranslogBytesTracking() { + assertFalse(remoteStoreSettings.isTranslogBytesTrackingEnabled()); + + clusterSettings.applySettings( + Settings.builder().put(RemoteStoreSettings.CLUSTER_REMOTE_TRANSLOG_TRACK_BYTES_SINCE_LAST_COMMIT_SETTING.getKey(), true).build() + ); + assertTrue(remoteStoreSettings.isTranslogBytesTrackingEnabled()); + + clusterSettings.applySettings( + Settings.builder() + .put(RemoteStoreSettings.CLUSTER_REMOTE_TRANSLOG_TRACK_BYTES_SINCE_LAST_COMMIT_SETTING.getKey(), false) + .build() + ); + assertFalse(remoteStoreSettings.isTranslogBytesTrackingEnabled()); + } + public void testUploadedSegmentsCleanupThreshold() { // Test default value assertEquals(1000, remoteStoreSettings.getUploadedSegmentsCleanupThreshold());