diff --git a/config/checkstyle/checkstyle.xml b/config/checkstyle/checkstyle.xml
index 3a88f90de8c..31b81beebcb 100644
--- a/config/checkstyle/checkstyle.xml
+++ b/config/checkstyle/checkstyle.xml
@@ -122,7 +122,7 @@
-
+
diff --git a/driver-core/src/main/com/mongodb/MongoClientSettings.java b/driver-core/src/main/com/mongodb/MongoClientSettings.java
index 858c1c83754..3cc4e9faba8 100644
--- a/driver-core/src/main/com/mongodb/MongoClientSettings.java
+++ b/driver-core/src/main/com/mongodb/MongoClientSettings.java
@@ -32,6 +32,7 @@
import com.mongodb.connection.SslSettings;
import com.mongodb.connection.TransportSettings;
import com.mongodb.event.CommandListener;
+import com.mongodb.internal.operation.CommandOperationHelper;
import com.mongodb.lang.Nullable;
import com.mongodb.observability.ObservabilitySettings;
import com.mongodb.spi.dns.DnsClient;
@@ -503,9 +504,11 @@ public Builder retryReads(final boolean retryReads) {
* the {@value MongoException#SYSTEM_OVERLOADED_ERROR_LABEL} and {@value MongoException#RETRYABLE_ERROR_LABEL} labels.
* Such errors are referred to as retryable overload errors.
*
- * Default is {@code null}, implies the value 2 and the above retry behavior. The implied value and behavior may change in
+ * Default is {@code null}, implies the value {@value CommandOperationHelper#DEFAULT_MAX_ADAPTIVE_RETRIES} and the above retry behavior.
+ * The implied value and behavior may change in
* the future in a minor version.
- * This means, there is no guarantee that not setting a value is equivalent to setting the value 2.
+ * This means, there is no guarantee that not setting a value is equivalent to setting the value
+ * {@value CommandOperationHelper#DEFAULT_MAX_ADAPTIVE_RETRIES}.
* The value 0 results in not retrying the attempts failed due to retryable overload errors.
*
*
@@ -961,7 +964,6 @@ public boolean getRetryReads() {
*/
@Beta(Reason.CLIENT)
@Nullable
- // TODO-BACKPRESSURE Valentin Use the `maxAdaptiveRetries` setting when retrying.
public Integer getMaxAdaptiveRetries() {
return maxAdaptiveRetries;
}
diff --git a/driver-core/src/main/com/mongodb/internal/operation/AbortTransactionOperation.java b/driver-core/src/main/com/mongodb/internal/operation/AbortTransactionOperation.java
index 9d655cd3d7a..1b3d6345839 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/AbortTransactionOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/AbortTransactionOperation.java
@@ -36,8 +36,8 @@ public class AbortTransactionOperation extends TransactionOperation {
private static final String COMMAND_NAME = "abortTransaction";
private BsonDocument recoveryToken;
- public AbortTransactionOperation(final WriteConcern writeConcern) {
- super(writeConcern);
+ public AbortTransactionOperation(final WriteConcern writeConcern, @Nullable final Integer maxAdaptiveRetriesSetting) {
+ super(writeConcern, maxAdaptiveRetriesSetting);
}
public AbortTransactionOperation recoveryToken(@Nullable final BsonDocument recoveryToken) {
diff --git a/driver-core/src/main/com/mongodb/internal/operation/AggregateOperation.java b/driver-core/src/main/com/mongodb/internal/operation/AggregateOperation.java
index 1f2a6ba518f..fa6abd6a00f 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/AggregateOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/AggregateOperation.java
@@ -45,13 +45,14 @@
public class AggregateOperation implements ReadOperationExplainable {
private final AggregateOperationImpl wrapped;
- public AggregateOperation(final MongoNamespace namespace, final List pipeline, final Decoder decoder) {
- this(namespace, pipeline, decoder, AggregationLevel.COLLECTION);
+ public AggregateOperation(final MongoNamespace namespace, final List pipeline, final Decoder decoder,
+ @Nullable final Integer maxAdaptiveRetriesSetting) {
+ this(namespace, pipeline, decoder, AggregationLevel.COLLECTION, maxAdaptiveRetriesSetting);
}
public AggregateOperation(final MongoNamespace namespace, final List pipeline, final Decoder decoder,
- final AggregationLevel aggregationLevel) {
- this.wrapped = new AggregateOperationImpl<>(namespace, pipeline, decoder, aggregationLevel);
+ final AggregationLevel aggregationLevel, @Nullable final Integer maxAdaptiveRetriesSetting) {
+ this.wrapped = new AggregateOperationImpl<>(namespace, pipeline, decoder, aggregationLevel, maxAdaptiveRetriesSetting);
}
public List getPipeline() {
diff --git a/driver-core/src/main/com/mongodb/internal/operation/AggregateOperationImpl.java b/driver-core/src/main/com/mongodb/internal/operation/AggregateOperationImpl.java
index 940073d4238..ae1cf007429 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/AggregateOperationImpl.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/AggregateOperationImpl.java
@@ -65,6 +65,8 @@ class AggregateOperationImpl implements ReadOperationCursor {
private final PipelineCreator pipelineCreator;
private boolean retryReads;
+ @Nullable
+ private final Integer maxAdaptiveRetriesSetting;
private Boolean allowDiskUse;
private Integer batchSize;
private Collation collation;
@@ -75,21 +77,24 @@ class AggregateOperationImpl implements ReadOperationCursor {
private CursorType cursorType;
AggregateOperationImpl(final MongoNamespace namespace,
- final List pipeline, final Decoder decoder, final AggregationLevel aggregationLevel) {
+ final List pipeline, final Decoder decoder, final AggregationLevel aggregationLevel,
+ @Nullable final Integer maxAdaptiveRetriesSetting) {
this(namespace, pipeline, decoder,
defaultAggregateTarget(notNull("aggregationLevel", aggregationLevel),
notNull("namespace", namespace).getCollectionName()),
- defaultPipelineCreator(pipeline));
+ defaultPipelineCreator(pipeline), maxAdaptiveRetriesSetting);
}
AggregateOperationImpl(final MongoNamespace namespace,
final List pipeline, final Decoder decoder, final AggregateTarget aggregateTarget,
- final PipelineCreator pipelineCreator) {
+ final PipelineCreator pipelineCreator,
+ @Nullable final Integer maxAdaptiveRetriesSetting) {
this.namespace = notNull("namespace", namespace);
this.pipeline = notNull("pipeline", pipeline);
this.decoder = notNull("decoder", decoder);
this.aggregateTarget = notNull("aggregateTarget", aggregateTarget);
this.pipelineCreator = notNull("pipelineCreator", pipelineCreator);
+ this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting;
}
List getPipeline() {
@@ -196,7 +201,7 @@ public MongoNamespace getNamespace() {
public BatchCursor execute(final ReadBinding binding, final OperationContext operationContext) {
return executeRetryableRead(binding, applyTimeoutModeToOperationContext(timeoutMode, operationContext), namespace.getDatabaseName(),
getCommandCreator(), CommandResultDocumentCodec.create(decoder, FIELD_NAMES_WITH_RESULT),
- transformer(), retryReads);
+ transformer(), retryReads, maxAdaptiveRetriesSetting);
}
@Override
@@ -204,7 +209,7 @@ public void executeAsync(final AsyncReadBinding binding, final OperationContext
SingleResultCallback> errHandlingCallback = errorHandlingCallback(callback, LOGGER);
executeRetryableReadAsync(binding, applyTimeoutModeToOperationContext(timeoutMode, operationContext), namespace.getDatabaseName(),
getCommandCreator(), CommandResultDocumentCodec.create(decoder, FIELD_NAMES_WITH_RESULT),
- asyncTransformer(), retryReads,
+ asyncTransformer(), retryReads, maxAdaptiveRetriesSetting,
errHandlingCallback);
}
diff --git a/driver-core/src/main/com/mongodb/internal/operation/AggregateToCollectionOperation.java b/driver-core/src/main/com/mongodb/internal/operation/AggregateToCollectionOperation.java
index 146a6680d37..69332c409b4 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/AggregateToCollectionOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/AggregateToCollectionOperation.java
@@ -178,7 +178,8 @@ public Void execute(final ReadBinding binding, final OperationContext operationC
getCommandCreator(),
new BsonDocumentCodec(),
transformer(),
- false);
+ false,
+ null);
}
@Override
@@ -194,6 +195,7 @@ public void executeAsync(final AsyncReadBinding binding, final OperationContext
new BsonDocumentCodec(),
asyncTransformer(),
false,
+ null,
callback);
}
diff --git a/driver-core/src/main/com/mongodb/internal/operation/AsyncOperationHelper.java b/driver-core/src/main/com/mongodb/internal/operation/AsyncOperationHelper.java
index d327ad663c9..cf0729b8555 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/AsyncOperationHelper.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/AsyncOperationHelper.java
@@ -176,9 +176,10 @@ static void executeRetryableReadAsync(
final Decoder decoder,
final CommandReadTransformerAsync transformer,
final boolean retryReadsSetting,
+ @Nullable final Integer maxAdaptiveRetriesSetting,
final SingleResultCallback callback) {
executeRetryableReadAsync(binding, operationContext, binding::getReadConnectionSource, database, commandCreator,
- decoder, transformer, retryReadsSetting, callback);
+ decoder, transformer, retryReadsSetting, maxAdaptiveRetriesSetting, callback);
}
static void executeRetryableReadAsync(
@@ -190,9 +191,10 @@ static void executeRetryableReadAsync(
final Decoder decoder,
final CommandReadTransformerAsync transformer,
final boolean retryReadsSetting,
+ @Nullable final Integer maxAdaptiveRetriesSetting,
final SingleResultCallback callback) {
RetryControl retryControl = createSpecRetryControl(
- new SpecRetryPolicy.IndividualPolicies(retryReadsSetting).includeRead(operationContext),
+ new SpecRetryPolicy.IndividualPolicies(retryReadsSetting).includeRead(operationContext).includeOverload(maxAdaptiveRetriesSetting),
operationContext);
binding.retain();
AsyncCallbackSupplier asyncRead = decorateWithRetriesAsync(retryControl, operationContext,
@@ -257,12 +259,13 @@ static void executeRetryableWriteAsync(
final CommandWriteTransformerAsync transformer,
final Function retryCommandModifier,
final boolean effectiveRetryWritesSetting,
+ @Nullable final Integer maxAdaptiveRetriesSetting,
final SingleResultCallback callback) {
beginAsync().thenSupply(c -> {
binding.retain();
MutableValue command = new MutableValue<>();
RetryControl retryControl = createSpecRetryControl(
- new SpecRetryPolicy.IndividualPolicies(effectiveRetryWritesSetting).includeWrite(),
+ new SpecRetryPolicy.IndividualPolicies(effectiveRetryWritesSetting).includeWrite().includeOverload(maxAdaptiveRetriesSetting),
operationContext);
AsyncCallbackSupplier retryingWrite = decorateWithRetriesAsync(retryControl, operationContext, supplierCallback -> {
beginAsync().thenSupply(withSourceAndConnectionCallback -> {
diff --git a/driver-core/src/main/com/mongodb/internal/operation/BaseFindAndModifyOperation.java b/driver-core/src/main/com/mongodb/internal/operation/BaseFindAndModifyOperation.java
index fdb24752b62..9e291c159be 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/BaseFindAndModifyOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/BaseFindAndModifyOperation.java
@@ -51,6 +51,8 @@ public abstract class BaseFindAndModifyOperation implements WriteOperation
private final MongoNamespace namespace;
private final WriteConcern writeConcern;
private final boolean retryWrites;
+ @Nullable
+ private final Integer maxAdaptiveRetriesSetting;
private final Decoder decoder;
private BsonDocument filter;
@@ -62,11 +64,14 @@ public abstract class BaseFindAndModifyOperation implements WriteOperation
private BsonValue comment;
private BsonDocument variables;
- protected BaseFindAndModifyOperation(final MongoNamespace namespace, final WriteConcern writeConcern, final boolean retryWrites,
+ protected BaseFindAndModifyOperation(final MongoNamespace namespace, final WriteConcern writeConcern,
+ final boolean retryWrites,
+ @Nullable final Integer maxAdaptiveRetriesSetting,
final Decoder decoder) {
this.namespace = notNull("namespace", namespace);
this.writeConcern = notNull("writeConcern", writeConcern);
this.retryWrites = retryWrites;
+ this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting;
this.decoder = notNull("decoder", decoder);
}
@@ -83,7 +88,8 @@ public T execute(final WriteBinding binding, final OperationContext operationCon
getCommandCreator(),
FindAndModifyHelper.transformer(),
cmd -> cmd,
- retryWrites);
+ retryWrites,
+ maxAdaptiveRetriesSetting);
}
@Override
@@ -91,7 +97,7 @@ public void executeAsync(final AsyncWriteBinding binding, final OperationContext
executeRetryableWriteAsync(binding, operationContext, getDatabaseName(), null, getFieldNameValidator(),
CommandResultDocumentCodec.create(getDecoder(), "value"),
getCommandCreator(),
- FindAndModifyHelper.asyncTransformer(), cmd -> cmd, retryWrites, callback);
+ FindAndModifyHelper.asyncTransformer(), cmd -> cmd, retryWrites, maxAdaptiveRetriesSetting, callback);
}
@Override
diff --git a/driver-core/src/main/com/mongodb/internal/operation/ChangeStreamOperation.java b/driver-core/src/main/com/mongodb/internal/operation/ChangeStreamOperation.java
index ee45e885ceb..c94dccebf50 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/ChangeStreamOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/ChangeStreamOperation.java
@@ -69,14 +69,15 @@ public class ChangeStreamOperation implements ReadOperationCursor {
@VisibleForTesting(otherwise = VisibleForTesting.AccessModifier.PRIVATE)
public ChangeStreamOperation(final MongoNamespace namespace, final FullDocument fullDocument,
final FullDocumentBeforeChange fullDocumentBeforeChange, final List pipeline, final Decoder decoder) {
- this(namespace, fullDocument, fullDocumentBeforeChange, pipeline, decoder, ChangeStreamLevel.COLLECTION);
+ this(namespace, fullDocument, fullDocumentBeforeChange, pipeline, decoder, ChangeStreamLevel.COLLECTION, null);
}
public ChangeStreamOperation(final MongoNamespace namespace, final FullDocument fullDocument,
final FullDocumentBeforeChange fullDocumentBeforeChange, final List pipeline, final Decoder decoder,
- final ChangeStreamLevel changeStreamLevel) {
+ final ChangeStreamLevel changeStreamLevel,
+ @Nullable final Integer maxAdaptiveRetriesSetting) {
this.wrapped = new AggregateOperationImpl<>(namespace, pipeline, RAW_BSON_DOCUMENT_CODEC, getAggregateTarget(),
- getPipelineCreator()).cursorType(CursorType.TailableAwait);
+ getPipelineCreator(), maxAdaptiveRetriesSetting).cursorType(CursorType.TailableAwait);
this.fullDocument = notNull("fullDocument", fullDocument);
this.fullDocumentBeforeChange = notNull("fullDocumentBeforeChange", fullDocumentBeforeChange);
this.decoder = notNull("decoder", decoder);
diff --git a/driver-core/src/main/com/mongodb/internal/operation/ClientBulkWriteOperation.java b/driver-core/src/main/com/mongodb/internal/operation/ClientBulkWriteOperation.java
index 4cfd0282c8e..ae9cd51b828 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/ClientBulkWriteOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/ClientBulkWriteOperation.java
@@ -156,6 +156,8 @@ public final class ClientBulkWriteOperation implements WriteOperation retryControl = createSpecRetryControl(
- new SpecRetryPolicy.IndividualPolicies(retryWritesSetting).includeWrite(),
+ new SpecRetryPolicy.IndividualPolicies(retryWritesSetting).includeWrite().includeOverload(maxAdaptiveRetriesSetting),
operationContext);
BatchEncoder batchEncoder = new BatchEncoder();
@@ -334,7 +338,7 @@ private void executeBatchAsync(
assertFalse(unexecutedModels.isEmpty());
SessionContext sessionContext = operationContext.getSessionContext();
RetryControl retryControl = createSpecRetryControl(
- new SpecRetryPolicy.IndividualPolicies(retryWritesSetting).includeWrite(),
+ new SpecRetryPolicy.IndividualPolicies(retryWritesSetting).includeWrite().includeOverload(maxAdaptiveRetriesSetting),
operationContext);
BatchEncoder batchEncoder = new BatchEncoder();
diff --git a/driver-core/src/main/com/mongodb/internal/operation/CommandOperationHelper.java b/driver-core/src/main/com/mongodb/internal/operation/CommandOperationHelper.java
index 6efe10b9e8c..f7d4c30b694 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/CommandOperationHelper.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/CommandOperationHelper.java
@@ -17,6 +17,7 @@
package com.mongodb.internal.operation;
import com.mongodb.MongoClientException;
+import com.mongodb.MongoClientSettings;
import com.mongodb.MongoCommandException;
import com.mongodb.MongoConnectionPoolClearedException;
import com.mongodb.MongoException;
@@ -42,6 +43,11 @@
@SuppressWarnings("overloads")
public final class CommandOperationHelper {
+ /**
+ * The value used when {@link MongoClientSettings#getMaxAdaptiveRetries()} is {@code null}.
+ */
+ public static final int DEFAULT_MAX_ADAPTIVE_RETRIES = 2;
+
static WriteConcern validateAndGetEffectiveWriteConcern(final WriteConcern writeConcernSetting, final SessionContext sessionContext)
throws MongoClientException {
boolean activeTransaction = sessionContext.hasActiveTransaction();
diff --git a/driver-core/src/main/com/mongodb/internal/operation/CommandReadOperation.java b/driver-core/src/main/com/mongodb/internal/operation/CommandReadOperation.java
index 6d51b13fd04..c593afbc60a 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/CommandReadOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/CommandReadOperation.java
@@ -68,14 +68,14 @@ public MongoNamespace getNamespace() {
@Override
public T execute(final ReadBinding binding, final OperationContext operationContext) {
return executeRetryableRead(binding, operationContext, databaseName, commandCreator, decoder,
- transformer(), false);
+ transformer(), false, null);
}
@Override
public void executeAsync(final AsyncReadBinding binding, final OperationContext operationContext,
final SingleResultCallback callback) {
executeRetryableReadAsync(binding, operationContext, databaseName, commandCreator, decoder,
- asyncTransformer(), false, callback);
+ asyncTransformer(), false, null, callback);
}
private static CommandReadTransformer transformer() {
diff --git a/driver-core/src/main/com/mongodb/internal/operation/CommitTransactionOperation.java b/driver-core/src/main/com/mongodb/internal/operation/CommitTransactionOperation.java
index a7a60ba7206..e56f0e07a4f 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/CommitTransactionOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/CommitTransactionOperation.java
@@ -53,12 +53,8 @@ public class CommitTransactionOperation extends TransactionOperation {
private final boolean alreadyCommitted;
private BsonDocument recoveryToken;
- public CommitTransactionOperation(final WriteConcern writeConcern) {
- this(writeConcern, false);
- }
-
- public CommitTransactionOperation(final WriteConcern writeConcern, final boolean alreadyCommitted) {
- super(writeConcern);
+ public CommitTransactionOperation(final WriteConcern writeConcern, @Nullable final Integer maxAdaptiveRetriesSetting, final boolean alreadyCommitted) {
+ super(writeConcern, maxAdaptiveRetriesSetting);
this.alreadyCommitted = alreadyCommitted;
}
diff --git a/driver-core/src/main/com/mongodb/internal/operation/CountDocumentsOperation.java b/driver-core/src/main/com/mongodb/internal/operation/CountDocumentsOperation.java
index ed762dfbd29..435718c5801 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/CountDocumentsOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/CountDocumentsOperation.java
@@ -43,6 +43,8 @@ public class CountDocumentsOperation implements ReadOperationSimple {
private static final Decoder DECODER = new BsonDocumentCodec();
private final MongoNamespace namespace;
private boolean retryReads;
+ @Nullable
+ private final Integer maxAdaptiveRetriesSetting;
private BsonDocument filter;
private BsonValue hint;
private BsonValue comment;
@@ -50,8 +52,9 @@ public class CountDocumentsOperation implements ReadOperationSimple {
private long limit;
private Collation collation;
- public CountDocumentsOperation(final MongoNamespace namespace) {
+ public CountDocumentsOperation(final MongoNamespace namespace, @Nullable final Integer maxAdaptiveRetriesSetting) {
this.namespace = notNull("namespace", namespace);
+ this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting;
}
@Nullable
@@ -157,7 +160,7 @@ public void executeAsync(final AsyncReadBinding binding, final OperationContext
}
private AggregateOperation getAggregateOperation() {
- return new AggregateOperation<>(namespace, getPipeline(), DECODER)
+ return new AggregateOperation<>(namespace, getPipeline(), DECODER, maxAdaptiveRetriesSetting)
.retryReads(retryReads)
.collation(collation)
.comment(comment)
diff --git a/driver-core/src/main/com/mongodb/internal/operation/CountOperation.java b/driver-core/src/main/com/mongodb/internal/operation/CountOperation.java
index 4bddd08edf4..09ba8d5a894 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/CountOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/CountOperation.java
@@ -47,14 +47,17 @@ public class CountOperation implements ReadOperationSimple {
private static final Decoder DECODER = new BsonDocumentCodec();
private final MongoNamespace namespace;
private boolean retryReads;
+ @Nullable
+ private final Integer maxAdaptiveRetriesSetting;
private BsonDocument filter;
private BsonValue hint;
private long skip;
private long limit;
private Collation collation;
- public CountOperation(final MongoNamespace namespace) {
+ public CountOperation(final MongoNamespace namespace, @Nullable final Integer maxAdaptiveRetriesSetting) {
this.namespace = notNull("namespace", namespace);
+ this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting;
}
public BsonDocument getFilter() {
@@ -124,13 +127,13 @@ public MongoNamespace getNamespace() {
@Override
public Long execute(final ReadBinding binding, final OperationContext operationContext) {
return executeRetryableRead(binding, operationContext, namespace.getDatabaseName(),
- getCommandCreator(), DECODER, transformer(), retryReads);
+ getCommandCreator(), DECODER, transformer(), retryReads, maxAdaptiveRetriesSetting);
}
@Override
public void executeAsync(final AsyncReadBinding binding, final OperationContext operationContext, final SingleResultCallback callback) {
executeRetryableReadAsync(binding, operationContext, namespace.getDatabaseName(),
- getCommandCreator(), DECODER, asyncTransformer(), retryReads, callback);
+ getCommandCreator(), DECODER, asyncTransformer(), retryReads, maxAdaptiveRetriesSetting, callback);
}
private CommandReadTransformer transformer() {
diff --git a/driver-core/src/main/com/mongodb/internal/operation/DistinctOperation.java b/driver-core/src/main/com/mongodb/internal/operation/DistinctOperation.java
index 94d669c24f7..0196d32cd27 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/DistinctOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/DistinctOperation.java
@@ -53,15 +53,19 @@ public class DistinctOperation implements ReadOperationCursor {
private final String fieldName;
private final Decoder decoder;
private boolean retryReads;
+ @Nullable
+ private final Integer maxAdaptiveRetriesSetting;
private BsonDocument filter;
private Collation collation;
private BsonValue comment;
private BsonValue hint;
- public DistinctOperation(final MongoNamespace namespace, final String fieldName, final Decoder decoder) {
+ public DistinctOperation(final MongoNamespace namespace, final String fieldName, final Decoder decoder,
+ @Nullable final Integer maxAdaptiveRetriesSetting) {
this.namespace = notNull("namespace", namespace);
this.fieldName = notNull("fieldName", fieldName);
this.decoder = notNull("decoder", decoder);
+ this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting;
}
public BsonDocument getFilter() {
@@ -122,13 +126,14 @@ public MongoNamespace getNamespace() {
@Override
public BatchCursor execute(final ReadBinding binding, final OperationContext operationContext) {
return executeRetryableRead(binding, operationContext, namespace.getDatabaseName(), getCommandCreator(), createCommandDecoder(),
- singleBatchCursorTransformer(VALUES), retryReads);
+ singleBatchCursorTransformer(VALUES), retryReads, maxAdaptiveRetriesSetting);
}
@Override
public void executeAsync(final AsyncReadBinding binding, final OperationContext operationContext, final SingleResultCallback> callback) {
executeRetryableReadAsync(binding, operationContext, namespace.getDatabaseName(),
- getCommandCreator(), createCommandDecoder(), asyncSingleBatchCursorTransformer(VALUES), retryReads,
+ getCommandCreator(), createCommandDecoder(), asyncSingleBatchCursorTransformer(VALUES),
+ retryReads, maxAdaptiveRetriesSetting,
errorHandlingCallback(callback, LOGGER));
}
diff --git a/driver-core/src/main/com/mongodb/internal/operation/DropCollectionOperation.java b/driver-core/src/main/com/mongodb/internal/operation/DropCollectionOperation.java
index 3a6b0487fbe..af6b73e6da0 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/DropCollectionOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/DropCollectionOperation.java
@@ -228,7 +228,7 @@ private BsonDocument getCollectionEncryptedFields(final BsonDocument defaultEncr
}
private ListCollectionsOperation listCollectionOperation() {
- return new ListCollectionsOperation<>(namespace.getDatabaseName(), BSON_VALUE_CODEC)
+ return new ListCollectionsOperation<>(namespace.getDatabaseName(), BSON_VALUE_CODEC, null)
.filter(new BsonDocument("name", new BsonString(namespace.getCollectionName())))
.batchSize(1);
}
diff --git a/driver-core/src/main/com/mongodb/internal/operation/EstimatedDocumentCountOperation.java b/driver-core/src/main/com/mongodb/internal/operation/EstimatedDocumentCountOperation.java
index 3e8bc6a49e0..d44c2fa87c8 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/EstimatedDocumentCountOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/EstimatedDocumentCountOperation.java
@@ -50,10 +50,13 @@ public class EstimatedDocumentCountOperation implements ReadOperationSimple DECODER = new BsonDocumentCodec();
private final MongoNamespace namespace;
private boolean retryReads;
+ @Nullable
+ private final Integer maxAdaptiveRetriesSetting;
private BsonValue comment;
- public EstimatedDocumentCountOperation(final MongoNamespace namespace) {
+ public EstimatedDocumentCountOperation(final MongoNamespace namespace, @Nullable final Integer maxAdaptiveRetriesSetting) {
this.namespace = notNull("namespace", namespace);
+ this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting;
}
public EstimatedDocumentCountOperation retryReads(final boolean retryReads) {
@@ -86,7 +89,7 @@ public Long execute(final ReadBinding binding, final OperationContext operationC
try {
return executeRetryableRead(binding, operationContext, namespace.getDatabaseName(),
getCommandCreator(), CommandResultDocumentCodec.create(DECODER, singletonList("firstBatch")),
- transformer(), retryReads);
+ transformer(), retryReads, maxAdaptiveRetriesSetting);
} catch (MongoCommandException e) {
return assertNotNull(rethrowIfNotNamespaceError(e, 0L));
}
@@ -96,7 +99,7 @@ public Long execute(final ReadBinding binding, final OperationContext operationC
public void executeAsync(final AsyncReadBinding binding, final OperationContext operationContext, final SingleResultCallback callback) {
executeRetryableReadAsync(binding, operationContext, namespace.getDatabaseName(),
getCommandCreator(), CommandResultDocumentCodec.create(DECODER, singletonList("firstBatch")),
- asyncTransformer(), retryReads,
+ asyncTransformer(), retryReads, maxAdaptiveRetriesSetting,
(result, t) -> {
if (isNamespaceError(t)) {
callback.onResult(0L, null);
diff --git a/driver-core/src/main/com/mongodb/internal/operation/FindAndDeleteOperation.java b/driver-core/src/main/com/mongodb/internal/operation/FindAndDeleteOperation.java
index db9d61b1dd4..bd87b53c7c2 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/FindAndDeleteOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/FindAndDeleteOperation.java
@@ -35,9 +35,11 @@
*/
public class FindAndDeleteOperation extends BaseFindAndModifyOperation {
- public FindAndDeleteOperation(final MongoNamespace namespace, final WriteConcern writeConcern, final boolean retryWrites,
+ public FindAndDeleteOperation(final MongoNamespace namespace, final WriteConcern writeConcern,
+ final boolean retryWrites,
+ @Nullable final Integer maxAdaptiveRetriesSetting,
final Decoder decoder) {
- super(namespace, writeConcern, retryWrites, decoder);
+ super(namespace, writeConcern, retryWrites, maxAdaptiveRetriesSetting, decoder);
}
@Override
diff --git a/driver-core/src/main/com/mongodb/internal/operation/FindAndReplaceOperation.java b/driver-core/src/main/com/mongodb/internal/operation/FindAndReplaceOperation.java
index 7073260a4c7..d77fad067a5 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/FindAndReplaceOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/FindAndReplaceOperation.java
@@ -45,9 +45,11 @@ public class FindAndReplaceOperation extends BaseFindAndModifyOperation {
private boolean upsert;
private Boolean bypassDocumentValidation;
- public FindAndReplaceOperation(final MongoNamespace namespace, final WriteConcern writeConcern, final boolean retryWrites,
+ public FindAndReplaceOperation(final MongoNamespace namespace, final WriteConcern writeConcern,
+ final boolean retryWrites,
+ @Nullable final Integer maxAdaptiveRetriesSetting,
final Decoder decoder, final BsonDocument replacement) {
- super(namespace, writeConcern, retryWrites, decoder);
+ super(namespace, writeConcern, retryWrites, maxAdaptiveRetriesSetting, decoder);
this.replacement = notNull("replacement", replacement);
}
diff --git a/driver-core/src/main/com/mongodb/internal/operation/FindAndUpdateOperation.java b/driver-core/src/main/com/mongodb/internal/operation/FindAndUpdateOperation.java
index e83deba30f3..8f46ce4d516 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/FindAndUpdateOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/FindAndUpdateOperation.java
@@ -52,15 +52,20 @@ public class FindAndUpdateOperation extends BaseFindAndModifyOperation {
private List arrayFilters;
public FindAndUpdateOperation(final MongoNamespace namespace,
- final WriteConcern writeConcern, final boolean retryWrites, final Decoder decoder, final BsonDocument update) {
- super(namespace, writeConcern, retryWrites, decoder);
+ final WriteConcern writeConcern,
+ final boolean retryWrites,
+ @Nullable final Integer maxAdaptiveRetriesSetting,
+ final Decoder decoder, final BsonDocument update) {
+ super(namespace, writeConcern, retryWrites, maxAdaptiveRetriesSetting, decoder);
this.update = notNull("update", update);
this.updatePipeline = null;
}
- public FindAndUpdateOperation(final MongoNamespace namespace, final WriteConcern writeConcern, final boolean retryWrites,
+ public FindAndUpdateOperation(final MongoNamespace namespace, final WriteConcern writeConcern,
+ final boolean retryWrites,
+ @Nullable final Integer maxAdaptiveRetriesSetting,
final Decoder decoder, final List update) {
- super(namespace, writeConcern, retryWrites, decoder);
+ super(namespace, writeConcern, retryWrites, maxAdaptiveRetriesSetting, decoder);
this.updatePipeline = update;
this.update = null;
}
diff --git a/driver-core/src/main/com/mongodb/internal/operation/FindOperation.java b/driver-core/src/main/com/mongodb/internal/operation/FindOperation.java
index bad874ac02c..1202e45a877 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/FindOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/FindOperation.java
@@ -73,6 +73,8 @@ public class FindOperation implements ReadOperationExplainable {
private final MongoNamespace namespace;
private final Decoder decoder;
private boolean retryReads;
+ @Nullable
+ private final Integer maxAdaptiveRetriesSetting;
private BsonDocument filter;
private int batchSize;
private int limit;
@@ -93,9 +95,13 @@ public class FindOperation implements ReadOperationExplainable {
private Boolean allowDiskUse;
private TimeoutMode timeoutMode;
- public FindOperation(final MongoNamespace namespace, final Decoder decoder) {
+ public FindOperation(
+ final MongoNamespace namespace,
+ final Decoder decoder,
+ @Nullable final Integer maxAdaptiveRetriesSetting) {
this.namespace = notNull("namespace", namespace);
this.decoder = notNull("decoder", decoder);
+ this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting;
}
@Override
@@ -299,7 +305,7 @@ public BatchCursor execute(final ReadBinding binding, final OperationContext
OperationContext findOperationContext = getFindOperationContext(operationContext);
RetryControl retryControl = createSpecRetryControl(
- new SpecRetryPolicy.IndividualPolicies(retryReads).includeRead(findOperationContext),
+ new SpecRetryPolicy.IndividualPolicies(retryReads).includeRead(findOperationContext).includeOverload(maxAdaptiveRetriesSetting),
findOperationContext);
Supplier> read = decorateWithRetries(retryControl, findOperationContext, () ->
withSourceAndConnection(binding::getReadConnectionSource, false, findOperationContext,
@@ -327,7 +333,7 @@ public void executeAsync(final AsyncReadBinding binding, final OperationContext
OperationContext findOperationContext = getFindOperationContext(operationContext);
RetryControl retryControl = createSpecRetryControl(
- new SpecRetryPolicy.IndividualPolicies(retryReads).includeRead(findOperationContext),
+ new SpecRetryPolicy.IndividualPolicies(retryReads).includeRead(findOperationContext).includeOverload(maxAdaptiveRetriesSetting),
findOperationContext);
binding.retain();
AsyncCallbackSupplier> asyncRead = decorateWithRetriesAsync(
diff --git a/driver-core/src/main/com/mongodb/internal/operation/ListCollectionsOperation.java b/driver-core/src/main/com/mongodb/internal/operation/ListCollectionsOperation.java
index e597e86a04e..c91a62bee55 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/ListCollectionsOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/ListCollectionsOperation.java
@@ -76,6 +76,8 @@ public class ListCollectionsOperation implements ReadOperationCursor {
private final String databaseName;
private final Decoder decoder;
private boolean retryReads;
+ @Nullable
+ private final Integer maxAdaptiveRetriesSetting;
private BsonDocument filter;
private int batchSize;
private boolean nameOnly;
@@ -83,9 +85,13 @@ public class ListCollectionsOperation implements ReadOperationCursor {
private BsonValue comment;
private TimeoutMode timeoutMode = TimeoutMode.CURSOR_LIFETIME;
- public ListCollectionsOperation(final String databaseName, final Decoder decoder) {
+ public ListCollectionsOperation(
+ final String databaseName,
+ final Decoder decoder,
+ @Nullable final Integer maxAdaptiveRetriesSetting) {
this.databaseName = notNull("databaseName", databaseName);
this.decoder = notNull("decoder", decoder);
+ this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting;
}
@Override
@@ -175,7 +181,7 @@ public BatchCursor execute(final ReadBinding binding, final OperationContext
OperationContext listCollectionsOperationContext = applyTimeoutModeToOperationContext(timeoutMode, operationContext);
RetryControl retryControl = createSpecRetryControl(
- new SpecRetryPolicy.IndividualPolicies(retryReads).includeRead(listCollectionsOperationContext),
+ new SpecRetryPolicy.IndividualPolicies(retryReads).includeRead(listCollectionsOperationContext).includeOverload(maxAdaptiveRetriesSetting),
listCollectionsOperationContext);
Supplier> read = decorateWithRetries(retryControl, listCollectionsOperationContext, () ->
withSourceAndConnection(binding::getReadConnectionSource, false, listCollectionsOperationContext, (source, connection, operationContextWithMinRTT) -> {
@@ -197,7 +203,7 @@ public void executeAsync(final AsyncReadBinding binding, final OperationContext
OperationContext listCollectionsOperationContext = applyTimeoutModeToOperationContext(timeoutMode, operationContext);
RetryControl retryControl = createSpecRetryControl(
- new SpecRetryPolicy.IndividualPolicies(retryReads).includeRead(listCollectionsOperationContext),
+ new SpecRetryPolicy.IndividualPolicies(retryReads).includeRead(listCollectionsOperationContext).includeOverload(maxAdaptiveRetriesSetting),
listCollectionsOperationContext);
binding.retain();
AsyncCallbackSupplier> asyncRead = decorateWithRetriesAsync(
diff --git a/driver-core/src/main/com/mongodb/internal/operation/ListDatabasesOperation.java b/driver-core/src/main/com/mongodb/internal/operation/ListDatabasesOperation.java
index b3abbdbc85c..2e5206cf557 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/ListDatabasesOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/ListDatabasesOperation.java
@@ -50,13 +50,16 @@ public class ListDatabasesOperation implements ReadOperationCursor {
private static final String DATABASES = "databases";
private final Decoder decoder;
private boolean retryReads;
+ @Nullable
+ private final Integer maxAdaptiveRetriesSetting;
private BsonDocument filter;
private Boolean nameOnly;
private Boolean authorizedDatabasesOnly;
private BsonValue comment;
- public ListDatabasesOperation(final Decoder decoder) {
+ public ListDatabasesOperation(final Decoder decoder, @Nullable final Integer maxAdaptiveRetriesSetting) {
this.decoder = notNull("decoder", decoder);
+ this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting;
}
public ListDatabasesOperation filter(@Nullable final BsonDocument filter) {
@@ -119,14 +122,14 @@ public MongoNamespace getNamespace() {
public BatchCursor execute(final ReadBinding binding, final OperationContext operationContext) {
return executeRetryableRead(binding, operationContext, "admin", getCommandCreator(),
CommandResultDocumentCodec.create(decoder, DATABASES),
- singleBatchCursorTransformer(DATABASES), retryReads);
+ singleBatchCursorTransformer(DATABASES), retryReads, maxAdaptiveRetriesSetting);
}
@Override
public void executeAsync(final AsyncReadBinding binding, final OperationContext operationContext,
final SingleResultCallback> callback) {
executeRetryableReadAsync(binding, operationContext, "admin", getCommandCreator(), CommandResultDocumentCodec.create(decoder, DATABASES),
- asyncSingleBatchCursorTransformer(DATABASES), retryReads, errorHandlingCallback(callback, LOGGER));
+ asyncSingleBatchCursorTransformer(DATABASES), retryReads, maxAdaptiveRetriesSetting, errorHandlingCallback(callback, LOGGER));
}
private CommandCreator getCommandCreator() {
diff --git a/driver-core/src/main/com/mongodb/internal/operation/ListIndexesOperation.java b/driver-core/src/main/com/mongodb/internal/operation/ListIndexesOperation.java
index c340c1c2f72..b1a03f93fed 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/ListIndexesOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/ListIndexesOperation.java
@@ -69,13 +69,19 @@ public class ListIndexesOperation implements ReadOperationCursor {
private final MongoNamespace namespace;
private final Decoder decoder;
private boolean retryReads;
+ @Nullable
+ private final Integer maxAdaptiveRetriesSetting;
private int batchSize;
private BsonValue comment;
private TimeoutMode timeoutMode = TimeoutMode.CURSOR_LIFETIME;
- public ListIndexesOperation(final MongoNamespace namespace, final Decoder decoder) {
+ public ListIndexesOperation(
+ final MongoNamespace namespace,
+ final Decoder decoder,
+ @Nullable final Integer maxAdaptiveRetriesSetting) {
this.namespace = notNull("namespace", namespace);
this.decoder = notNull("decoder", decoder);
+ this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting;
}
public Integer getBatchSize() {
@@ -132,7 +138,7 @@ public BatchCursor execute(final ReadBinding binding, final OperationContext
OperationContext listIndexesOperationContext = applyTimeoutModeToOperationContext(timeoutMode, operationContext);
RetryControl retryControl = createSpecRetryControl(
- new SpecRetryPolicy.IndividualPolicies(retryReads).includeRead(listIndexesOperationContext),
+ new SpecRetryPolicy.IndividualPolicies(retryReads).includeRead(listIndexesOperationContext).includeOverload(maxAdaptiveRetriesSetting),
listIndexesOperationContext);
Supplier> read = decorateWithRetries(retryControl, listIndexesOperationContext, () ->
withSourceAndConnection(binding::getReadConnectionSource, false, listIndexesOperationContext, (source, connection, operationContextWithMinRTT) -> {
@@ -153,7 +159,7 @@ public void executeAsync(final AsyncReadBinding binding, final OperationContext
OperationContext listIndexesOperationContext = applyTimeoutModeToOperationContext(timeoutMode, operationContext);
RetryControl retryControl = createSpecRetryControl(
- new SpecRetryPolicy.IndividualPolicies(retryReads).includeRead(listIndexesOperationContext),
+ new SpecRetryPolicy.IndividualPolicies(retryReads).includeRead(listIndexesOperationContext).includeOverload(maxAdaptiveRetriesSetting),
listIndexesOperationContext);
binding.retain();
AsyncCallbackSupplier> asyncRead = decorateWithRetriesAsync(
diff --git a/driver-core/src/main/com/mongodb/internal/operation/ListSearchIndexesOperation.java b/driver-core/src/main/com/mongodb/internal/operation/ListSearchIndexesOperation.java
index 6d2e5278184..71ffdcfc3c8 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/ListSearchIndexesOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/ListSearchIndexesOperation.java
@@ -59,11 +59,15 @@ public final class ListSearchIndexesOperation implements ReadOperationExplain
@Nullable
private final String indexName;
private final boolean retryReads;
+ @Nullable
+ private final Integer maxAdaptiveRetriesSetting;
public ListSearchIndexesOperation(final MongoNamespace namespace, final Decoder decoder, @Nullable final String indexName,
@Nullable final Integer batchSize, @Nullable final Collation collation,
@Nullable final BsonValue comment,
- @Nullable final Boolean allowDiskUse, final boolean retryReads) {
+ @Nullable final Boolean allowDiskUse,
+ final boolean retryReads,
+ @Nullable final Integer maxAdaptiveRetriesSetting) {
this.namespace = namespace;
this.decoder = decoder;
this.allowDiskUse = allowDiskUse;
@@ -72,6 +76,7 @@ public ListSearchIndexesOperation(final MongoNamespace namespace, final Decoder<
this.comment = comment;
this.indexName = indexName;
this.retryReads = retryReads;
+ this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting;
}
@Override
@@ -119,7 +124,7 @@ public ReadOperationSimple asExplainableOperation(@Nullable final Explain
private AggregateOperation asAggregateOperation() {
BsonDocument searchDefinition = getSearchDefinition();
BsonDocument listSearchIndexesStage = new BsonDocument(STAGE_LIST_SEARCH_INDEXES, searchDefinition);
- return new AggregateOperation<>(namespace, singletonList(listSearchIndexesStage), decoder)
+ return new AggregateOperation<>(namespace, singletonList(listSearchIndexesStage), decoder, maxAdaptiveRetriesSetting)
.retryReads(retryReads)
.collation(collation)
.comment(comment)
diff --git a/driver-core/src/main/com/mongodb/internal/operation/MapReduceWithInlineResultsOperation.java b/driver-core/src/main/com/mongodb/internal/operation/MapReduceWithInlineResultsOperation.java
index 6e2511beebc..b43a619b203 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/MapReduceWithInlineResultsOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/MapReduceWithInlineResultsOperation.java
@@ -175,7 +175,7 @@ public String getCommandName() {
public MapReduceBatchCursor execute(final ReadBinding binding, final OperationContext operationContext) {
return executeRetryableRead(binding, operationContext, namespace.getDatabaseName(),
getCommandCreator(),
- CommandResultDocumentCodec.create(decoder, "results"), transformer(), false);
+ CommandResultDocumentCodec.create(decoder, "results"), transformer(), false, null);
}
@Override
@@ -184,7 +184,7 @@ public void executeAsync(final AsyncReadBinding binding, final OperationContext
SingleResultCallback> errHandlingCallback = errorHandlingCallback(callback, LOGGER);
executeRetryableReadAsync(binding, operationContext, namespace.getDatabaseName(),
getCommandCreator(), CommandResultDocumentCodec.create(decoder, "results"),
- asyncTransformer(), false, errHandlingCallback);
+ asyncTransformer(), false, null, errHandlingCallback);
}
public ReadOperationSimple asExplainableOperation(final ExplainVerbosity explainVerbosity) {
diff --git a/driver-core/src/main/com/mongodb/internal/operation/MixedBulkWriteOperation.java b/driver-core/src/main/com/mongodb/internal/operation/MixedBulkWriteOperation.java
index 80ab345fec7..4c771b14cac 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/MixedBulkWriteOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/MixedBulkWriteOperation.java
@@ -78,6 +78,8 @@ public class MixedBulkWriteOperation implements WriteOperation
private final List extends WriteRequest> writeRequests;
private final boolean ordered;
private final boolean retryWrites;
+ @Nullable
+ private final Integer maxAdaptiveRetriesSetting;
private final WriteConcern writeConcern;
private Boolean bypassDocumentValidation;
private String commandName;
@@ -85,7 +87,9 @@ public class MixedBulkWriteOperation implements WriteOperation
private BsonDocument variables;
public MixedBulkWriteOperation(final MongoNamespace namespace, final List extends WriteRequest> writeRequests,
- final boolean ordered, final WriteConcern writeConcern, final boolean retryWrites) {
+ final boolean ordered, final WriteConcern writeConcern,
+ final boolean retryWrites,
+ @Nullable final Integer maxAdaptiveRetriesSetting) {
notNull("writeRequests", writeRequests);
isTrueArgument("writeRequests is not an empty list", !writeRequests.isEmpty());
this.commandName = notNull("commandName", writeRequests.get(0).getType().toString().toLowerCase(Locale.ROOT));
@@ -94,6 +98,7 @@ public MixedBulkWriteOperation(final MongoNamespace namespace, final List exte
this.ordered = ordered;
this.writeConcern = notNull("writeConcern", writeConcern);
this.retryWrites = retryWrites;
+ this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting;
}
@Override
@@ -238,7 +243,7 @@ private BatchWithSourceAndConnection executeBatchReusingCon
MutableValue batch = new MutableValue<>(maybeBatch);
MutableValue sourceAndConnection = new MutableValue<>(maybeSourceAndConnection);
RetryControl retryControl = createSpecRetryControl(
- new SpecRetryPolicy.IndividualPolicies(retryWrites).includeWrite(),
+ new SpecRetryPolicy.IndividualPolicies(retryWrites).includeWrite().includeOverload(maxAdaptiveRetriesSetting),
operationContext);
Supplier> retryingBatchExecutor = decorateWithRetries(
retryControl,
@@ -291,7 +296,7 @@ private void executeBatchReusingConnectionAsync(
MutableValue batch = new MutableValue<>(maybeBatch);
MutableValue sourceAndConnection = new MutableValue<>(maybeSourceAndConnection);
RetryControl retryControl = createSpecRetryControl(
- new SpecRetryPolicy.IndividualPolicies(retryWrites).includeWrite(),
+ new SpecRetryPolicy.IndividualPolicies(retryWrites).includeWrite().includeOverload(maxAdaptiveRetriesSetting),
operationContext);
AsyncCallbackSupplier> retryingBatchExecutor = decorateWithRetriesAsync(
retryControl,
diff --git a/driver-core/src/main/com/mongodb/internal/operation/Operations.java b/driver-core/src/main/com/mongodb/internal/operation/Operations.java
index da0661220da..0ff5b31422f 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/Operations.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/Operations.java
@@ -105,23 +105,25 @@ public final class Operations {
private final WriteConcern writeConcern;
private final boolean retryWrites;
private final boolean retryReads;
+ @Nullable
+ private final Integer maxAdaptiveRetriesSetting;
private final TimeoutSettings timeoutSettings;
public Operations(final Class documentClass, final ReadPreference readPreference, final CodecRegistry codecRegistry,
- final boolean retryReads, final TimeoutSettings timeoutSettings) {
+ final boolean retryReads, @Nullable final Integer maxAdaptiveRetriesSetting, final TimeoutSettings timeoutSettings) {
this(null, documentClass, readPreference, codecRegistry, ReadConcern.DEFAULT, WriteConcern.ACKNOWLEDGED,
- true, retryReads, timeoutSettings);
+ true, retryReads, maxAdaptiveRetriesSetting, timeoutSettings);
}
public Operations(@Nullable final MongoNamespace namespace, final Class documentClass, final ReadPreference readPreference,
- final CodecRegistry codecRegistry, final boolean retryReads, final TimeoutSettings timeoutSettings) {
+ final CodecRegistry codecRegistry, final boolean retryReads, @Nullable final Integer maxAdaptiveRetriesSetting, final TimeoutSettings timeoutSettings) {
this(namespace, documentClass, readPreference, codecRegistry, ReadConcern.DEFAULT, WriteConcern.ACKNOWLEDGED,
- true, retryReads, timeoutSettings);
+ true, retryReads, maxAdaptiveRetriesSetting, timeoutSettings);
}
public Operations(@Nullable final MongoNamespace namespace, final Class documentClass, final ReadPreference readPreference,
final CodecRegistry codecRegistry, final ReadConcern readConcern, final WriteConcern writeConcern, final boolean retryWrites,
- final boolean retryReads, final TimeoutSettings timeoutSettings) {
+ final boolean retryReads, @Nullable final Integer maxAdaptiveRetriesSetting, final TimeoutSettings timeoutSettings) {
this.namespace = namespace;
this.documentClass = documentClass;
this.readPreference = readPreference;
@@ -129,6 +131,7 @@ public Operations(@Nullable final MongoNamespace namespace, final Class docum
this.readConcern = readConcern;
this.retryWrites = retryWrites;
this.retryReads = retryReads;
+ this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting;
this.timeoutSettings = timeoutSettings;
WriteConcern writeConcernToUse = writeConcern;
@@ -225,7 +228,8 @@ public TimeoutSettings createTimeoutSettings(final DropIndexOptions options) {
public ReadOperationSimple countDocuments(final Bson filter, final CountOptions options) {
CountDocumentsOperation operation = new CountDocumentsOperation(
- assertNotNull(namespace))
+ assertNotNull(namespace),
+ maxAdaptiveRetriesSetting)
.retryReads(retryReads)
.filter(toBsonDocument(filter))
.skip(options.getSkip())
@@ -242,7 +246,7 @@ public ReadOperationSimple countDocuments(final Bson filter, final CountOp
public ReadOperationSimple estimatedDocumentCount(final EstimatedDocumentCountOptions options) {
return new EstimatedDocumentCountOperation(
- assertNotNull(namespace))
+ assertNotNull(namespace), maxAdaptiveRetriesSetting)
.retryReads(retryReads)
.comment(options.getComment());
}
@@ -265,7 +269,8 @@ public ReadOperationExplainable find(final MongoNamespace findNamespace,
private FindOperation createFindOperation(final MongoNamespace findNamespace, @Nullable final Bson filter,
final Class resultClass, final FindOptions options) {
FindOperation operation = new FindOperation<>(
- findNamespace, codecRegistry.get(resultClass))
+ findNamespace, codecRegistry.get(resultClass),
+ maxAdaptiveRetriesSetting)
.retryReads(retryReads)
.filter(filter == null ? new BsonDocument() : filter.toBsonDocument(documentClass, codecRegistry))
.batchSize(options.getBatchSize())
@@ -297,7 +302,7 @@ private FindOperation createFindOperation(final MongoNamespace findNamesp
public ReadOperationCursor distinct(final String fieldName, @Nullable final Bson filter, final Class resultClass,
final Collation collation, final BsonValue comment, @Nullable final Bson hint, @Nullable final String hintString) {
DistinctOperation operation = new DistinctOperation<>(assertNotNull(namespace),
- fieldName, codecRegistry.get(resultClass))
+ fieldName, codecRegistry.get(resultClass), maxAdaptiveRetriesSetting)
.retryReads(retryReads)
.filter(filter == null ? null : filter.toBsonDocument(documentClass, codecRegistry))
.collation(collation)
@@ -316,7 +321,8 @@ public ReadOperationExplainable aggregate(final List extends Bson> pipe
final Collation collation, @Nullable final Bson hint, @Nullable final String hintString,
final BsonValue comment, final Bson variables, final Boolean allowDiskUse, final AggregationLevel aggregationLevel) {
return new AggregateOperation<>(assertNotNull(namespace),
- assertNotNull(toBsonDocumentList(pipeline)), codecRegistry.get(resultClass), aggregationLevel)
+ assertNotNull(toBsonDocumentList(pipeline)), codecRegistry.get(resultClass), aggregationLevel,
+ maxAdaptiveRetriesSetting)
.retryReads(retryReads)
.allowDiskUse(allowDiskUse)
.batchSize(batchSize)
@@ -392,7 +398,7 @@ public ReadOperationMapReduceCursor mapReduce(final String mapFunction,
public WriteOperation findOneAndDelete(final Bson filter, final FindOneAndDeleteOptions options) {
return new FindAndDeleteOperation<>(
- assertNotNull(namespace), writeConcern, retryWrites, getCodec())
+ assertNotNull(namespace), writeConcern, retryWrites, maxAdaptiveRetriesSetting, getCodec())
.filter(toBsonDocument(filter))
.projection(toBsonDocument(options.getProjection()))
.sort(toBsonDocument(options.getSort()))
@@ -406,7 +412,7 @@ public WriteOperation findOneAndDelete(final Bson filter, final FindOneAndDel
public WriteOperation findOneAndReplace(final Bson filter, final T replacement,
final FindOneAndReplaceOptions options) {
return new FindAndReplaceOperation<>(
- assertNotNull(namespace), writeConcern, retryWrites, getCodec(), documentToBsonDocument(replacement))
+ assertNotNull(namespace), writeConcern, retryWrites, maxAdaptiveRetriesSetting, getCodec(), documentToBsonDocument(replacement))
.filter(toBsonDocument(filter))
.projection(toBsonDocument(options.getProjection()))
.sort(toBsonDocument(options.getSort()))
@@ -422,7 +428,7 @@ public WriteOperation findOneAndReplace(final Bson filter, final T replacemen
public WriteOperation findOneAndUpdate(final Bson filter, final Bson update, final FindOneAndUpdateOptions options) {
return new FindAndUpdateOperation<>(
- assertNotNull(namespace), writeConcern, retryWrites, getCodec(), assertNotNull(toBsonDocument(update)))
+ assertNotNull(namespace), writeConcern, retryWrites, maxAdaptiveRetriesSetting, getCodec(), assertNotNull(toBsonDocument(update)))
.filter(toBsonDocument(filter))
.projection(toBsonDocument(options.getProjection()))
.sort(toBsonDocument(options.getSort()))
@@ -440,7 +446,7 @@ public WriteOperation findOneAndUpdate(final Bson filter, final Bson update,
public WriteOperation findOneAndUpdate(final Bson filter, final List extends Bson> update,
final FindOneAndUpdateOptions options) {
return new FindAndUpdateOperation<>(
- assertNotNull(namespace), writeConcern, retryWrites, getCodec(), assertNotNull(toBsonDocumentList(update)))
+ assertNotNull(namespace), writeConcern, retryWrites, maxAdaptiveRetriesSetting, getCodec(), assertNotNull(toBsonDocumentList(update)))
.filter(toBsonDocument(filter))
.projection(toBsonDocument(options.getProjection()))
.sort(toBsonDocument(options.getSort()))
@@ -516,7 +522,7 @@ public WriteOperation insertMany(final List extends T> docume
}
return new MixedBulkWriteOperation(assertNotNull(namespace),
- requests, options.isOrdered(), writeConcern, retryWrites)
+ requests, options.isOrdered(), writeConcern, retryWrites, maxAdaptiveRetriesSetting)
.bypassDocumentValidation(options.getBypassDocumentValidation())
.comment(options.getComment());
}
@@ -586,7 +592,7 @@ public WriteOperation bulkWrite(final List extends WriteModel
}
return new MixedBulkWriteOperation(assertNotNull(namespace), writeRequests,
- options.isOrdered(), writeConcern, retryWrites)
+ options.isOrdered(), writeConcern, retryWrites, maxAdaptiveRetriesSetting)
.bypassDocumentValidation(options.getBypassDocumentValidation())
.comment(options.getComment())
.let(toBsonDocument(options.getLet()));
@@ -738,7 +744,7 @@ public ReadOperationExplainable listSearchIndexes(final Class resultCl
@Nullable final String indexName, @Nullable final Integer batchSize, @Nullable final Collation collation,
@Nullable final BsonValue comment, @Nullable final Boolean allowDiskUse) {
return new ListSearchIndexesOperation<>(assertNotNull(namespace),
- codecRegistry.get(resultClass), indexName, batchSize, collation, comment, allowDiskUse, retryReads);
+ codecRegistry.get(resultClass), indexName, batchSize, collation, comment, allowDiskUse, retryReads, maxAdaptiveRetriesSetting);
}
public WriteOperation dropIndex(final String indexName, final DropIndexOptions ignoredOptions) {
@@ -752,7 +758,7 @@ public WriteOperation dropIndex(final Bson keys, final DropIndexOptions ig
public ReadOperationCursor listCollections(final String databaseName, final Class resultClass,
final Bson filter, final boolean collectionNamesOnly, final boolean authorizedCollections, @Nullable final Integer batchSize,
final BsonValue comment, @Nullable final TimeoutMode timeoutMode) {
- return new ListCollectionsOperation<>(databaseName, codecRegistry.get(resultClass))
+ return new ListCollectionsOperation<>(databaseName, codecRegistry.get(resultClass), maxAdaptiveRetriesSetting)
.retryReads(retryReads)
.filter(toBsonDocument(filter))
.nameOnly(collectionNamesOnly)
@@ -765,7 +771,7 @@ public ReadOperationCursor listCollections(final String databaseName, fin
public ReadOperationCursor listDatabases(final Class resultClass, final Bson filter,
final Boolean nameOnly,
final Boolean authorizedDatabasesOnly, final BsonValue comment) {
- return new ListDatabasesOperation<>(codecRegistry.get(resultClass))
+ return new ListDatabasesOperation<>(codecRegistry.get(resultClass), maxAdaptiveRetriesSetting)
.retryReads(retryReads)
.filter(toBsonDocument(filter))
.nameOnly(nameOnly)
@@ -776,7 +782,8 @@ public ReadOperationCursor listDatabases(final Class resultClass, fina
public ReadOperationCursor listIndexes(final Class resultClass, @Nullable final Integer batchSize,
final BsonValue comment, @Nullable final TimeoutMode timeoutMode) {
return new ListIndexesOperation<>(assertNotNull(namespace),
- codecRegistry.get(resultClass))
+ codecRegistry.get(resultClass),
+ maxAdaptiveRetriesSetting)
.retryReads(retryReads)
.batchSize(batchSize == null ? 0 : batchSize)
.comment(comment)
@@ -792,7 +799,8 @@ public ReadOperationCursor changeStream(final FullDocument fullDocument,
assertNotNull(namespace),
fullDocument,
fullDocumentBeforeChange,
- assertNotNull(toBsonDocumentList(pipeline)), decoder, changeStreamLevel)
+ assertNotNull(toBsonDocumentList(pipeline)), decoder, changeStreamLevel,
+ maxAdaptiveRetriesSetting)
.batchSize(batchSize)
.collation(collation)
.comment(comment)
@@ -806,7 +814,7 @@ public ReadOperationCursor changeStream(final FullDocument fullDocument,
public WriteOperation clientBulkWriteOperation(
final List extends ClientNamespacedWriteModel> clientWriteModels,
@Nullable final ClientBulkWriteOptions options) {
- return new ClientBulkWriteOperation(clientWriteModels, options, writeConcern, retryWrites, codecRegistry);
+ return new ClientBulkWriteOperation(clientWriteModels, options, writeConcern, retryWrites, maxAdaptiveRetriesSetting, codecRegistry);
}
private Codec getCodec() {
diff --git a/driver-core/src/main/com/mongodb/internal/operation/SpecRetryPolicy.java b/driver-core/src/main/com/mongodb/internal/operation/SpecRetryPolicy.java
index b339c3b39e7..49ec3a595e3 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/SpecRetryPolicy.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/SpecRetryPolicy.java
@@ -29,6 +29,7 @@
import com.mongodb.internal.async.function.RetryPolicy.Decision.RetryAttemptInfo;
import com.mongodb.internal.connection.OperationContext;
import com.mongodb.internal.connection.OperationContext.ServerDeprioritization;
+import com.mongodb.internal.time.ExponentialBackoff;
import com.mongodb.lang.Nullable;
import java.time.Duration;
@@ -38,12 +39,15 @@
import java.util.Set;
import java.util.function.Supplier;
+import static com.mongodb.MongoException.RETRYABLE_ERROR_LABEL;
+import static com.mongodb.MongoException.SYSTEM_OVERLOADED_ERROR_LABEL;
import static com.mongodb.assertions.Assertions.assertFalse;
import static com.mongodb.assertions.Assertions.assertNotNull;
import static com.mongodb.assertions.Assertions.assertNull;
import static com.mongodb.assertions.Assertions.assertTrue;
import static com.mongodb.assertions.Assertions.fail;
import static com.mongodb.internal.TimeoutContext.createMongoTimeoutException;
+import static com.mongodb.internal.operation.CommandOperationHelper.DEFAULT_MAX_ADAPTIVE_RETRIES;
import static com.mongodb.internal.operation.CommandOperationHelper.NO_WRITES_PERFORMED_ERROR_LABEL;
import static com.mongodb.internal.operation.CommandOperationHelper.RETRYABLE_WRITE_ERROR_LABEL;
import static com.mongodb.internal.operation.CommandOperationHelper.addRetryableWriteErrorLabelIfNeeded;
@@ -61,7 +65,7 @@ final class SpecRetryPolicy implements RetryPolicy {
private static final int INFINITE_ATTEMPTS = Integer.MAX_VALUE;
private final IndividualPolicies policies;
- private final int maxAttempts;
+ private int maxAttempts;
private final ServerDeprioritization serverDeprioritization;
private Supplier commandDescriptionSupplier;
@@ -132,6 +136,8 @@ public Decision onAttemptFailure(final RetryContext retryContext, final Throwabl
decideRetryableAndAddRetryableWriteErrorLabelIfNeeded(state, attemptFailedResult)).orElse(false);
retryableError |= policies.read().map(state ->
decideRetryable(state, attemptFailedResult)).orElse(false);
+ retryableError |= policies.overload().map(state ->
+ decideRetryableAndUpdateMaxAttempts(state, attemptFailedResult)).orElse(false);
boolean maxAttemptsReached = attempt >= maxAttempts - 1;
if (policies.isRetrySettingEffectivelyTrue() && !retryableError) {
logUnableToRetryError(attemptFailedResult);
@@ -139,7 +145,7 @@ public Decision onAttemptFailure(final RetryContext retryContext, final Throwabl
boolean retry = retryableError && !maxAttemptsReached;
Decision decision = new Decision(
decideProspectiveFailedResult(retryContext.getProspectiveFailedResult().orElse(null), maybeInternalAttemptFailedResult),
- retry ? new RetryAttemptInfo(Duration.ZERO) : null);
+ retry ? new RetryAttemptInfo(calculateOverloadBackoff(attemptFailedResult, attempt + 1)) : null);
policies.write().ifPresent(IndividualPolicies.State.Write::resetRequirementsInfo);
return decision;
}
@@ -195,6 +201,26 @@ private static boolean decideRetryable(final IndividualPolicies.State.Read state
return isRetryableMongoSecurityException(mongoException) || isRetryableException(mongoException);
}
+ /**
+ * Returns {@code true} iff another attempt must be executed provided that {@link #maxAttempts} has not been reached.
+ */
+ private boolean decideRetryableAndUpdateMaxAttempts(final IndividualPolicies.State.Overload state, final Throwable exception) {
+ boolean hasRetryableErrorLabel;
+ boolean hasSystemOverloadErrorLabel;
+ if (exception instanceof MongoException) {
+ MongoException mongoException = (MongoException) exception;
+ hasRetryableErrorLabel = mongoException.hasErrorLabel(RETRYABLE_ERROR_LABEL);
+ hasSystemOverloadErrorLabel = mongoException.hasErrorLabel(SYSTEM_OVERLOADED_ERROR_LABEL);
+ } else {
+ hasRetryableErrorLabel = false;
+ hasSystemOverloadErrorLabel = false;
+ }
+ if (hasSystemOverloadErrorLabel && policies.isRetrySettingEffectivelyTrue()) {
+ maxAttempts = state.getMaxAttempts();
+ }
+ return state.isRequirementsMet() && hasRetryableErrorLabel && hasSystemOverloadErrorLabel;
+ }
+
private void logUnableToRetryError(final Throwable exception) {
assertFalse(exception instanceof OperationHelper.ResourceSupplierInternalException);
if (LOGGER.isDebugEnabled()) {
@@ -246,6 +272,15 @@ private static Throwable decideReadProspectiveFailedResult(
}
}
+ private static Duration calculateOverloadBackoff(final Throwable attemptFailedResult, final int immediateNextAttempt) {
+ assertFalse(attemptFailedResult instanceof OperationHelper.ResourceSupplierInternalException);
+ if (attemptFailedResult instanceof MongoException
+ && ((MongoException) attemptFailedResult).hasErrorLabel(SYSTEM_OVERLOADED_ERROR_LABEL)) {
+ return ExponentialBackoff.calculateOverloadBackoff(immediateNextAttempt);
+ }
+ return Duration.ZERO;
+ }
+
private static int maxAttempts(final int maxRetries) {
assertTrue(maxRetries < INFINITE_ATTEMPTS - 1);
return maxRetries + 1;
@@ -271,6 +306,7 @@ static final class IndividualPolicies {
CONFLICTS = new EnumMap<>(Descriptor.class);
CONFLICTS.put(Descriptor.WRITE, EnumSet.of(Descriptor.READ));
CONFLICTS.put(Descriptor.READ, EnumSet.of(Descriptor.WRITE));
+ CONFLICTS.put(Descriptor.OVERLOAD, EnumSet.noneOf(Descriptor.class));
assertTrue(CONFLICTS.keySet().containsAll(asList(Descriptor.values())));
}
@@ -322,6 +358,16 @@ IndividualPolicies includeRead(final OperationContext operationContext) {
return this;
}
+ /**
+ * See
+ * Client Backpressure,
+ * which specifies the overload retry policy.
+ */
+ IndividualPolicies includeOverload(@Nullable final Integer maxAdaptiveRetriesSetting) {
+ include(Descriptor.OVERLOAD, new State.Overload(effectiveRetrySetting, maxAdaptiveRetriesSetting));
+ return this;
+ }
+
private void include(final Descriptor descriptor, final State state) {
State previous = policies.put(descriptor, state);
assertNull(previous);
@@ -333,6 +379,8 @@ private int getMaxAttempts() {
maxRetries = Descriptor.WRITE.maxRetries;
} else if (policies.containsKey(Descriptor.READ)) {
maxRetries = Descriptor.READ.maxRetries;
+ } else if (policies.containsKey(Descriptor.OVERLOAD)) {
+ maxRetries = overload().map(State.Overload::getMaxAttempts).orElse(Descriptor.OVERLOAD.maxRetries);
} else {
throw fail(toString());
}
@@ -347,6 +395,10 @@ private Optional read() {
return Optional.ofNullable((State.Read) policies.get(Descriptor.READ));
}
+ private Optional overload() {
+ return Optional.ofNullable((State.Overload) policies.get(Descriptor.OVERLOAD));
+ }
+
@Override
public String toString() {
return "IndividualPolicies{"
@@ -363,7 +415,11 @@ private enum Descriptor {
/**
* See {@link #includeRead(OperationContext)}.
*/
- READ(1);
+ READ(1),
+ /**
+ * See {@link #includeOverload(Integer)}.
+ */
+ OVERLOAD(DEFAULT_MAX_ADAPTIVE_RETRIES);
private final int maxRetries;
@@ -462,6 +518,41 @@ public String toString() {
+ '}';
}
}
+
+ static final class Overload extends State {
+ private final boolean requirementsMet;
+ @Nullable
+ private final Integer maxAdaptiveRetriesSetting;
+
+ Overload(
+ final boolean effectiveRetrySetting,
+ @Nullable
+ final Integer maxAdaptiveRetriesSetting) {
+ requirementsMet = effectiveRetrySetting;
+ this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting;
+ }
+
+ int getMaxAttempts() {
+ return maxAttempts(maxAdaptiveRetriesSetting == null
+ ? IndividualPolicies.Descriptor.OVERLOAD.maxRetries
+ : maxAdaptiveRetriesSetting);
+ }
+
+ /**
+ * See {@link Read#isRequirementsMet()}.
+ */
+ boolean isRequirementsMet() {
+ return requirementsMet;
+ }
+
+ @Override
+ public String toString() {
+ return "Overload{"
+ + "requirementsMet=" + requirementsMet
+ + ", maxAdaptiveRetriesSetting=" + maxAdaptiveRetriesSetting
+ + '}';
+ }
+ }
}
}
diff --git a/driver-core/src/main/com/mongodb/internal/operation/SyncOperationHelper.java b/driver-core/src/main/com/mongodb/internal/operation/SyncOperationHelper.java
index 8fa8b0db1de..32476ede201 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/SyncOperationHelper.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/SyncOperationHelper.java
@@ -187,9 +187,10 @@ static T executeRetryableRead(
final CommandCreator commandCreator,
final Decoder decoder,
final CommandReadTransformer transformer,
- final boolean retryReadsSetting) {
+ final boolean retryReadsSetting,
+ @Nullable final Integer maxAdaptiveRetriesSetting) {
return executeRetryableRead(operationContext, binding::getReadConnectionSource, database, commandCreator,
- decoder, transformer, retryReadsSetting);
+ decoder, transformer, retryReadsSetting, maxAdaptiveRetriesSetting);
}
static T executeRetryableRead(
@@ -199,9 +200,11 @@ static T executeRetryableRead(
final CommandCreator commandCreator,
final Decoder decoder,
final CommandReadTransformer transformer,
- final boolean retryReadsSetting) {
+ final boolean retryReadsSetting,
+ @Nullable
+ final Integer maxAdaptiveRetriesSetting) {
RetryControl retryControl = createSpecRetryControl(
- new SpecRetryPolicy.IndividualPolicies(retryReadsSetting).includeRead(operationContext),
+ new SpecRetryPolicy.IndividualPolicies(retryReadsSetting).includeRead(operationContext).includeOverload(maxAdaptiveRetriesSetting),
operationContext);
Supplier read = decorateWithRetries(retryControl, operationContext, () ->
@@ -262,10 +265,11 @@ static R executeRetryableWrite(
final CommandCreator commandCreator,
final CommandWriteTransformer transformer,
final com.mongodb.Function retryCommandModifier,
- final boolean effectiveRetryWritesSetting) {
+ final boolean effectiveRetryWritesSetting,
+ @Nullable final Integer maxAdaptiveRetriesSetting) {
MutableValue command = new MutableValue<>();
RetryControl retryControl = createSpecRetryControl(
- new SpecRetryPolicy.IndividualPolicies(effectiveRetryWritesSetting).includeWrite(),
+ new SpecRetryPolicy.IndividualPolicies(effectiveRetryWritesSetting).includeWrite().includeOverload(maxAdaptiveRetriesSetting),
operationContext);
Supplier retryingWrite = decorateWithRetries(retryControl, operationContext, () -> {
boolean firstAttempt = retryControl.isFirstAttempt();
diff --git a/driver-core/src/main/com/mongodb/internal/operation/TransactionOperation.java b/driver-core/src/main/com/mongodb/internal/operation/TransactionOperation.java
index 6ee79ab6a17..6d8a254cb28 100644
--- a/driver-core/src/main/com/mongodb/internal/operation/TransactionOperation.java
+++ b/driver-core/src/main/com/mongodb/internal/operation/TransactionOperation.java
@@ -24,6 +24,7 @@
import com.mongodb.internal.binding.WriteBinding;
import com.mongodb.internal.connection.OperationContext;
import com.mongodb.internal.validator.NoOpFieldNameValidator;
+import com.mongodb.lang.Nullable;
import org.bson.BsonDocument;
import org.bson.BsonInt32;
import org.bson.codecs.BsonDocumentCodec;
@@ -45,9 +46,12 @@
*/
public abstract class TransactionOperation implements WriteOperation {
private final WriteConcern writeConcern;
+ @Nullable
+ private final Integer maxAdaptiveRetriesSetting;
- TransactionOperation(final WriteConcern writeConcern) {
+ TransactionOperation(final WriteConcern writeConcern, @Nullable final Integer maxAdaptiveRetriesSetting) {
this.writeConcern = notNull("writeConcern", writeConcern);
+ this.maxAdaptiveRetriesSetting = maxAdaptiveRetriesSetting;
}
public WriteConcern getWriteConcern() {
@@ -60,7 +64,7 @@ public Void execute(final WriteBinding binding, final OperationContext operation
TimeoutContext timeoutContext = operationContext.getTimeoutContext();
return executeRetryableWrite(binding, operationContext, "admin", null, NoOpFieldNameValidator.INSTANCE,
new BsonDocumentCodec(), getCommandCreator(),
- writeConcernErrorTransformer(timeoutContext), getRetryCommandModifier(timeoutContext), true);
+ writeConcernErrorTransformer(timeoutContext), getRetryCommandModifier(timeoutContext), true, maxAdaptiveRetriesSetting);
}
@Override
@@ -70,7 +74,7 @@ public void executeAsync(final AsyncWriteBinding binding, final OperationContext
executeRetryableWriteAsync(binding, operationContext, "admin", null, NoOpFieldNameValidator.INSTANCE,
new BsonDocumentCodec(), getCommandCreator(),
writeConcernErrorTransformerAsync(timeoutContext), getRetryCommandModifier(timeoutContext),
- true,
+ true, maxAdaptiveRetriesSetting,
errorHandlingCallback(callback, LOGGER));
}
diff --git a/driver-core/src/main/com/mongodb/internal/time/ExponentialBackoff.java b/driver-core/src/main/com/mongodb/internal/time/ExponentialBackoff.java
index db5d2efa996..11ea3c52e27 100644
--- a/driver-core/src/main/com/mongodb/internal/time/ExponentialBackoff.java
+++ b/driver-core/src/main/com/mongodb/internal/time/ExponentialBackoff.java
@@ -17,7 +17,9 @@
package com.mongodb.internal.time;
import com.mongodb.internal.VisibleForTesting;
+import com.mongodb.internal.async.function.RetryControl;
+import java.time.Duration;
import java.util.concurrent.ThreadLocalRandom;
import java.util.function.DoubleSupplier;
@@ -29,10 +31,8 @@
*/
public final class ExponentialBackoff {
- private static final double TRANSACTION_BASE_MS = 5.0;
@VisibleForTesting(otherwise = PRIVATE)
static final double TRANSACTION_MAX_MS = 500.0;
- private static final double TRANSACTION_GROWTH = 1.5;
// TODO-JAVA-6079
private static DoubleSupplier testJitterSupplier = null;
@@ -43,17 +43,38 @@ private ExponentialBackoff() {
/**
* Calculate the backoff in milliseconds for transaction retries.
*
- * @param attemptNumber attempt number > 0
+ * @param attemptNumber A positive 0-based attempt number. That is, the attempt must be a retry attempt.
* @return The calculated backoff in milliseconds.
+ *
+ * @see RetryControl#attempt()
*/
public static long calculateTransactionBackoffMs(final int attemptNumber) {
- assertTrue(attemptNumber > 0, "Attempt number must be at least 1 (1-based) in the context of transaction backoff calculation");
+ return calculateBackoffMs(5.0, TRANSACTION_MAX_MS, 1.5, attemptNumber);
+ }
+
+ /**
+ * Calculate the backoff for command retries caused by
+ * {@linkplain com.mongodb.MongoException#SYSTEM_OVERLOADED_ERROR_LABEL overload}.
+ * See {@link #calculateTransactionBackoffMs(int)} for more details.
+ */
+ public static Duration calculateOverloadBackoff(final int attemptNumber) {
+ return Duration.ofMillis(calculateBackoffMs(100, 10000, 2, attemptNumber));
+ }
+
+ /**
+ * Calculate the backoff in milliseconds for transaction retries.
+ *
+ * @param attemptNumber attempt number > 0
+ * @return The calculated backoff in milliseconds.
+ */
+ public static long calculateBackoffMs(final double baseMs, final double maxMs, final double growth, final int attemptNumber) {
+ assertTrue(attemptNumber > 0, String.valueOf(attemptNumber));
double jitter = testJitterSupplier != null
? testJitterSupplier.getAsDouble()
: ThreadLocalRandom.current().nextDouble();
return Math.round(jitter * Math.min(
- TRANSACTION_BASE_MS * Math.pow(TRANSACTION_GROWTH, attemptNumber - 1),
- TRANSACTION_MAX_MS));
+ baseMs * Math.pow(growth, attemptNumber - 1),
+ maxMs));
}
/**
diff --git a/driver-core/src/test/functional/com/mongodb/OperationFunctionalSpecification.groovy b/driver-core/src/test/functional/com/mongodb/OperationFunctionalSpecification.groovy
index 906f2e3ff68..5a8d65d51ad 100644
--- a/driver-core/src/test/functional/com/mongodb/OperationFunctionalSpecification.groovy
+++ b/driver-core/src/test/functional/com/mongodb/OperationFunctionalSpecification.groovy
@@ -108,13 +108,13 @@ class OperationFunctionalSpecification extends Specification {
void acknowledgeWrite(final SingleConnectionBinding binding) {
new MixedBulkWriteOperation(getNamespace(), [new InsertRequest(new BsonDocument())], true,
- ACKNOWLEDGED, false).execute(binding, createOperationContext())
+ ACKNOWLEDGED, false, null).execute(binding, createOperationContext())
binding.release()
}
void acknowledgeWrite(final AsyncSingleConnectionBinding binding) {
executeAsync(new MixedBulkWriteOperation(getNamespace(), [new InsertRequest(new BsonDocument())],
- true, ACKNOWLEDGED, false), binding)
+ true, ACKNOWLEDGED, false, null), binding)
binding.release()
}
diff --git a/driver-core/src/test/functional/com/mongodb/client/test/CollectionHelper.java b/driver-core/src/test/functional/com/mongodb/client/test/CollectionHelper.java
index e6c28d9d5bc..a54ed99b5db 100644
--- a/driver-core/src/test/functional/com/mongodb/client/test/CollectionHelper.java
+++ b/driver-core/src/test/functional/com/mongodb/client/test/CollectionHelper.java
@@ -286,7 +286,8 @@ public void insertDocuments(final List documents, final WriteConce
for (BsonDocument document : documents) {
insertRequests.add(new InsertRequest(document));
}
- new MixedBulkWriteOperation(namespace, insertRequests, true, writeConcern, false).execute(binding, ClusterFixture.createOperationContext());
+ new MixedBulkWriteOperation(namespace, insertRequests, true, writeConcern, false, null).execute(
+ binding, ClusterFixture.createOperationContext());
}
public void insertDocuments(final Document... documents) {
@@ -328,7 +329,7 @@ public List find() {
public Optional listSearchIndex(final String indexName) {
ListSearchIndexesOperation listSearchIndexesOperation =
- new ListSearchIndexesOperation<>(namespace, codec, indexName, null, null, null, null, true);
+ new ListSearchIndexesOperation<>(namespace, codec, indexName, null, null, null, null, true, null);
BatchCursor cursor = listSearchIndexesOperation.execute(getBinding(), ClusterFixture.createOperationContext());
List results = new ArrayList<>();
@@ -346,7 +347,8 @@ public void createSearchIndex(final SearchIndexRequest searchIndexModel) {
}
public List find(final Codec codec) {
- BatchCursor cursor = new FindOperation<>(namespace, codec)
+ BatchCursor cursor = new FindOperation<>(namespace, codec,
+ null)
.sort(new BsonDocument("_id", new BsonInt32(1)))
.execute(getBinding(), ClusterFixture.createOperationContext());
List results = new ArrayList<>();
@@ -366,7 +368,7 @@ public void updateOne(final Bson filter, final Bson update, final boolean isUpse
update.toBsonDocument(Document.class, registry),
WriteRequest.Type.UPDATE)
.upsert(isUpsert)),
- true, WriteConcern.ACKNOWLEDGED, false)
+ true, WriteConcern.ACKNOWLEDGED, false, null)
.execute(getBinding(), ClusterFixture.createOperationContext());
}
@@ -376,7 +378,7 @@ public void replaceOne(final Bson filter, final Bson update, final boolean isUps
update.toBsonDocument(Document.class, registry),
WriteRequest.Type.REPLACE)
.upsert(isUpsert)),
- true, WriteConcern.ACKNOWLEDGED, false)
+ true, WriteConcern.ACKNOWLEDGED, false, null)
.execute(getBinding(), ClusterFixture.createOperationContext());
}
@@ -391,7 +393,7 @@ public void deleteMany(final Bson filter) {
private void delete(final Bson filter, final boolean multi) {
new MixedBulkWriteOperation(namespace,
singletonList(new DeleteRequest(filter.toBsonDocument(Document.class, registry)).multi(multi)),
- true, WriteConcern.ACKNOWLEDGED, false)
+ true, WriteConcern.ACKNOWLEDGED, false, null)
.execute(getBinding(), ClusterFixture.createOperationContext());
}
@@ -416,7 +418,7 @@ private List aggregate(final List pipeline, final Decoder decode
for (Bson cur : pipeline) {
bsonDocumentPipeline.add(cur.toBsonDocument(Document.class, registry));
}
- BatchCursor cursor = new AggregateOperation<>(namespace, bsonDocumentPipeline, decoder, level)
+ BatchCursor cursor = new AggregateOperation<>(namespace, bsonDocumentPipeline, decoder, level, null)
.execute(getBinding(), ClusterFixture.createOperationContext());
List results = new ArrayList<>();
while (cursor.hasNext()) {
@@ -451,7 +453,7 @@ public List find(final BsonDocument filter, final BsonDocument sort, fina
}
public List find(final BsonDocument filter, final BsonDocument sort, final BsonDocument projection, final Decoder