从Couchbase Java SDK2迁移至SDK3:实现等效重试策略
Couchbase Java SDK 2 → SDK3 重试机制迁移等效实现
核心差异说明
SDK3 摒弃了 SDK2 中基于 RxJava retryWhen 的重试方式,改用内置的 RetryStrategy 接口统一管理重试逻辑,支持单请求级配置或全局默认配置。以下是与原代码行为完全一致的迁移方案:
步骤1:异常映射
原 SDK2 中指定的异常在 SDK3 中大多兼容,仅部分拼写调整:
RequestCancelledException→RequestCanceledException(注意拼写:Cancelled → Canceled)- 其余异常(
TemporaryFailureException、TemporaryLockFailureException等)包路径调整为com.couchbase.client.core.error.*
步骤2:自定义 RetryStrategy 实现
对应原 getRetryFunction() 的逻辑,实现 SDK3 的 RetryStrategy 接口:
import com.couchbase.client.core.retry.RetryAction; import com.couchbase.client.core.retry.RetryStrategy; import com.couchbase.client.core.error.*; import reactor.core.publisher.Mono; import java.time.Duration; public class CustomCouchbaseRetryStrategy implements RetryStrategy { private final int maxAttempts; private final com.couchbase.client.core.retry.Delay delayStrategy; public CustomCouchbaseRetryStrategy(int maxAttempts, com.couchbase.client.core.retry.Delay delayStrategy) { this.maxAttempts = maxAttempts; this.delayStrategy = delayStrategy; } @Override public Mono<RetryAction> shouldRetry(RequestContext ctx) { Throwable cause = ctx.lastThrowable(); // 匹配原代码指定的异常类型 boolean isTargetException = cause instanceof TemporaryFailureException || cause instanceof TemporaryLockFailureException || cause instanceof BackpressureException || cause instanceof ReplicaNotAvailableException || cause instanceof RequestCanceledException || cause instanceof CouchbaseOutOfMemoryException; if (isTargetException && ctx.retryAttempts() < maxAttempts) { // 按原延迟策略计算重试间隔 Duration delay = delayStrategy.delay(ctx.retryAttempts()); return Mono.just(RetryAction.withDelay(delay)); } // 不满足重试条件则终止 return Mono.just(RetryAction.noRetry()); } }
注:若原
CouchbaseRetryStrategy.determineRetryStrategy()返回的是 SDK2 的Delay,需转换为 SDK3 的com.couchbase.client.core.retry.Delay实现(比如Delay.fixed(Duration.ofMillis(100))、Delay.exponential(Duration.ofMillis(100)))。
步骤3:Upsert 请求中应用重试策略
SDK3 中需通过 AsyncCollection(默认使用 _default Scope/Collection)执行异步操作,并通过 RequestOptions 绑定重试策略:
// 获取 SDK3 异步集合对象(对应原 SDK2 的 AsyncBucket) AsyncCollection asyncCollection = client.bucket("your-bucket") .defaultScope() .defaultCollection(); // 初始化自定义重试策略 CustomCouchbaseRetryStrategy retryStrategy = new CustomCouchbaseRetryStrategy( retryStrategy.getMaxAttempts(), // 转换原延迟策略为 SDK3 兼容的 Delay 实例 ((CouchbaseRetryStrategy) retryStrategy).determineRetryStrategy() ); // 配置请求选项并执行 upsert RequestOptions options = RequestOptions.requestOptions().retryStrategy(retryStrategy); asyncCollection.upsert(doc.id(), doc.content(), options);
可选:全局配置重试策略
若需所有请求默认应用该策略,可在 Cluster 初始化时配置:
Cluster cluster = Cluster.connect("cluster-connection-string", ClusterOptions.clusterOptions("username", "password") .environment(env -> env.retryStrategy(new CustomCouchbaseRetryStrategy(maxAttempts, delayStrategy))));
内容的提问来源于stack exchange,提问作者Keshav Singh
相关产品推荐
相关产品推荐

