You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.30 01:36:18