Skip to content
Closed
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 @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -409,7 +409,8 @@ protected TranslogManager createTranslogManager(
this::ensureOpen,
engineConfig.getTranslogFactory(),
engineConfig.getStartedPrimarySupplier(),
TranslogOperationHelper.create(engineConfig)
TranslogOperationHelper.create(engineConfig),
true
);
}

Expand Down Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,9 @@ public class InternalTranslogManager implements TranslogManager {
private final TranslogEventListener translogEventListener;
private final Supplier<LocalCheckpointTracker> localCheckpointTrackerSupplier;
private final Logger logger;
private final TranslogBytesTracker translogBytesTracker = new TranslogBytesTracker();
private final boolean supportsTranslogBytesTracking;
private TranslogBytesTracker.CommitSnapshot commitSnapshot;

public InternalTranslogManager(
TranslogConfig translogConfig,
Expand All @@ -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<LocalCheckpointTracker> localCheckpointTrackerSupplier,
String translogUUID,
TranslogEventListener translogEventListener,
LifecycleAware engineLifeCycleAware,
TranslogFactory translogFactory,
BooleanSupplier startedPrimarySupplier,
TranslogOperationHelper translogOperationHelper,
boolean supportsTranslogBytesTracking
) throws IOException {
this.shardId = shardId;
this.readLock = readLock;
Expand All @@ -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);
Expand Down Expand Up @@ -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;
}

/**
Expand Down Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<Boolean> 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
*/
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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;
}
Expand Down
Loading