Java Cassandra Driver v4如何通过DriverConfigLoader配置重试策略
解决方案
1. 自定义实现FallthroughRetryPolicy
4.x版本驱动移除了3.x内置的FallthroughRetryPolicy,该策略核心逻辑是所有异常场景都不触发重试,直接将异常抛出给上层业务,你可以自行实现RetryPolicy接口完成逻辑适配,无状态场景可使用单例模式:
import com.datastax.oss.driver.api.core.ConsistencyLevel; import com.datastax.oss.driver.api.core.context.DriverContext; import com.datastax.oss.driver.api.core.retry.RetryDecision; import com.datastax.oss.driver.api.core.retry.RetryPolicy; import com.datastax.oss.driver.api.core.servererrors.CoordinatorException; import com.datastax.oss.driver.api.core.servererrors.WriteType; import com.datastax.oss.driver.api.core.session.Request; import org.jspecify.annotations.NonNull; public class FallthroughRetryPolicy implements RetryPolicy { public static final FallthroughRetryPolicy INSTANCE = new FallthroughRetryPolicy(); // 保留无参构造,支持配置反射加载 public FallthroughRetryPolicy() {} // 保留驱动上下文构造,支持配置文件声明式加载 public FallthroughRetryPolicy(DriverContext context, String profileName) {} @Override public RetryDecision onReadTimeout(@NonNull Request request, @NonNull ConsistencyLevel cl, int blockFor, int received, boolean dataPresent, int retryCount) { return RetryDecision.rethrow(); } @Override public RetryDecision onWriteTimeout(@NonNull Request request, @NonNull ConsistencyLevel cl, @NonNull WriteType writeType, int blockFor, int received, int retryCount) { return RetryDecision.rethrow(); } @Override public RetryDecision onUnavailable(@NonNull Request request, @NonNull ConsistencyLevel cl, int required, int alive, int retryCount) { return RetryDecision.rethrow(); } @Override public RetryDecision onRequestAborted(@NonNull Request request, @NonNull Throwable error, int retryCount) { return RetryDecision.rethrow(); } @Override public RetryDecision onErrorResponse(@NonNull Request request, @NonNull CoordinatorException error, int retryCount) { return RetryDecision.rethrow(); } @Override public void close() {} }
2. 自定义实现LoggingRetryPolicy
4.x版本同时移除了装饰器模式的LoggingRetryPolicy,你可以自行实现包装类,持有实际生效的重试策略实例,在每次做重试决策前打印对应日志即可:
import com.datastax.oss.driver.api.core.ConsistencyLevel; import com.datastax.oss.driver.api.core.retry.RetryDecision; import com.datastax.oss.driver.api.core.retry.RetryPolicy; import com.datastax.oss.driver.api.core.servererrors.CoordinatorException; import com.datastax.oss.driver.api.core.servererrors.WriteType; import com.datastax.oss.driver.api.core.session.Request; import org.jspecify.annotations.NonNull; import org.slf4j.Logger; import org.slf4j.LoggerFactory; public class LoggingRetryPolicy implements RetryPolicy { private static final Logger log = LoggerFactory.getLogger(LoggingRetryPolicy.class); private final RetryPolicy delegate; public LoggingRetryPolicy(RetryPolicy delegate) { this.delegate = delegate; } @Override public RetryDecision onReadTimeout(@NonNull Request request, @NonNull ConsistencyLevel cl, int blockFor, int received, boolean dataPresent, int retryCount) { RetryDecision decision = delegate.onReadTimeout(request, cl, blockFor, received, dataPresent, retryCount); logEvent("read timeout", request, cl, retryCount, decision); return decision; } @Override public RetryDecision onWriteTimeout(@NonNull Request request, @NonNull ConsistencyLevel cl, @NonNull WriteType writeType, int blockFor, int received, int retryCount) { RetryDecision decision = delegate.onWriteTimeout(request, cl, writeType, blockFor, received, retryCount); logEvent("write timeout", request, cl, retryCount, decision); return decision; } @Override public RetryDecision onUnavailable(@NonNull Request request, @NonNull ConsistencyLevel cl, int required, int alive, int retryCount) { RetryDecision decision = delegate.onUnavailable(request, cl, required, alive, retryCount); logEvent("coordinator unavailable", request, cl, retryCount, decision); return decision; } @Override public RetryDecision onRequestAborted(@NonNull Request request, @NonNull Throwable error, int retryCount) { RetryDecision decision = delegate.onRequestAborted(request, error, retryCount); log.warn("Request aborted, retry count: {}, decision: {}", retryCount, decision, error); return decision; } @Override public RetryDecision onErrorResponse(@NonNull Request request, @NonNull CoordinatorException error, int retryCount) { RetryDecision decision = delegate.onErrorResponse(request, error, retryCount); log.warn("Received coordinator error, retry count: {}, decision: {}", retryCount, decision, error); return decision; } private void logEvent(String eventType, Request request, ConsistencyLevel cl, int retryCount, RetryDecision decision) { if (decision.isRetry()) { log.info("{} occurred on request {} with consistency level {}, retry count: {}, will retry on {}", eventType, request, cl, retryCount, decision.getRetryTarget()); } else { log.info("{} occurred on request {} with consistency level {}, retry count: {}, final decision: {}", eventType, request, cl, retryCount, decision); } } @Override public void close() throws Exception { delegate.close(); } }
3. 适配原有业务逻辑
注意:4.x版本已经废弃了3.x的
Cluster入口,统一使用CqlSession构建客户端实例。
你不需要强制通过DriverConfigLoader指定类名配置策略,直接在构建Session时手动传入重试策略实例即可,完全兼容原有动态选择策略、按需开启日志包装的逻辑:
// 适配4.x的策略转换逻辑,和原有3.x逻辑完全对齐 private static RetryPolicy retryPolicyDataConvert(String retryPolicyStr) { if (CassandraConstants.CASSANDRACONNECTION_RETRYPOLICY_DEFAULT.equals(retryPolicyStr)) { return DefaultRetryPolicy.INSTANCE; } else if (CassandraConstants.CASSANDRACONNECTION_RETRYPOLICY_DOWNGRADING.equals(retryPolicyStr)) { return DowngradingConsistencyRetryPolicy.INSTANCE; } else if (CassandraConstants.CASSANDRACONNECTION_RETRYPOLICY_FALLTHROUGH.equals(retryPolicyStr)) { return FallthroughRetryPolicy.INSTANCE; } // 默认兜底走默认重试策略 return DefaultRetryPolicy.INSTANCE; } public static RetryPolicy getRetryPolicy(String retryPolicyStr, boolean isLogingPolicy) { RetryPolicy basePolicy = retryPolicyDataConvert(retryPolicyStr); return isLogingPolicy ? new LoggingRetryPolicy(basePolicy) : basePolicy; } // 构建客户端时直接注入策略 CqlSession session = CqlSession.builder() // 此处补充你的接触点、本地数据中心、认证信息等其他配置 .withRetryPolicy(getRetryPolicy(connectionInfo.getRetryPolicy(), connectionInfo.isLoggingRetryPolicy())) .build();
如果你一定要使用DriverConfigLoader编程式配置:无参构造的基础策略(DefaultRetryPolicy、DowngradingConsistencyRetryPolicy、自定义的FallthroughRetryPolicy)可以直接通过withClass指定类名,但LoggingRetryPolicy因为需要传入被包装的策略实例,无法通过反射无参构造加载,更推荐使用上述手动传入实例的方式,灵活度更高。
内容的提问来源于stack exchange,提问作者prathyusha magam
相关产品推荐
相关产品推荐

