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

Reactive API的Cassandra自定义配置启动失败求助

Cassandra自定义配置启动失败,默认配置正常运行问题

问题描述

尝试多种Cassandra自定义配置,应用均无法启动;但使用Spring默认自动配置时可正常运行。由于需要从KeyVault获取用户名和密码,必须采用自定义配置,当前配置代码如下:

@Configuration
@EnableReactiveCassandraRepositories(basePackages = {"com.ey.nexus.map.viewer.repository"})
@EntityScan({"com.ey.nexus.map.viewer.model"})
@RequiredArgsConstructor
public class CassandraConfiguration extends AbstractReactiveCassandraConfiguration {

    private final KeyVaultConfigManager configManager;
    @Value("${spring.data.cassandra.keyspace-name}")
    private String keyspace;
    @Value("${spring.data.cassandra.contact-points}")
    private String contactPoints;
    @Value("${spring.data.cassandra.port}")
    private Integer port;
    @Value("${spring.data.cassandra.ssl}")
    private boolean sslEnabled;
    @Value("${spring.data.cassandra.username}")
    private String usernameKey;
    @Value("${spring.data.cassandra.password}")
    private String passKey;
    @Value("${spring.data.cassandra.local-datacenter}")
    private String localDatacenter;
    @Value("${spring.data.cassandra.socket-connection-timeout-millis}")
    private int socketConnectionTimeoutMillis;
    @Value("${spring.data.cassandra.socket-read-timeout-mills}")
    private int socketReadTimeoutMillis;
    @Value("${spring.data.cassandra.pooling-heartbeat-interval-seconds}")
    private int poolingHeartbeatIntervalSeconds;
    @Value("${spring.data.cassandra.pooling-idle-timeout-seconds}")
    private int poolingIdleTimeoutSeconds;
    @Value("${spring.data.cassandra.pooling-timeout-millis}")
    private int poolingTimeoutMillis;
    @Value("${spring.data.cassandra.pooling-options.local.core-connections-per-host}")
    private int poolingLocalCoreConnectionsPerHost;
    @Value("${spring.data.cassandra.pooling-options.local.max-connections-per-host}")
    private int poolingLocalMaxConnectionsPerHost;
    @Value("${spring.data.cassandra.pooling-options.remote.core-connections-per-host}")
    private int poolingRemoteCoreConnectionsPerHost;
    @Value("${spring.data.cassandra.pooling-options.remote.max-connections-per-host}")
    private int poolingRemoteMaxConnectionsPerHost;
    @Value("${spring.data.cassandra.pooling-options.local.max-requests-per-connection}")
    private int poolingLocalMaxRequestsPerConnection;
    @Value("${spring.data.cassandra.pooling-options.remote.max-requests-per-connection}")
    private int poolingRemoteMaxRequestsPerConnection;
    @Value("${spring.data.cassandra.reconnect-delay-millis}")
    private long reconnectDelayMillis;
    @Value("${spring.data.cassandra.reconnect-base-delay-millis}")
    private long reconnectBaseDelayMillis;
    @Value("${spring.data.cassandra.reconnect-max-delay-millis}")
    private long reconnectMaxDelayMillis;
    @Value("${spring.data.cassandra.retry-policy.class}")
    private String retryPolicyClass;
    @Value("${spring.data.cassandra.retry-policy.max-retry-count}")
    private int retryPolicyMaxRetryCount;
    @Value("${spring.data.cassandra.retry-policy.growing-back-off-time-millis}")
    private int retryPolicyGrowingBackOffTimeMillis;
    @Value("${spring.data.cassandra.retry-policy.fixed-back-off-time-millis}")
    private int retryPolicyFixedBackOffTimeMillis;
    @Value("${spring.data.cassandra.load-balancing-policy.read-datacenter}")
    private String readDatacenter;
    @Value("${spring.data.cassandra.load-balancing-policy.write-datacenter}")
    private String writeDatacenter;
    @Value("${spring.data.cassandra.load-balancing-policy.global-endpoint}")
    private String globalEndpoint;
    @Value("${spring.data.cassandra.load-balancing-policy.dnsExpirationInSeconds}")
    private int dnsExpirationInSeconds;
    private String username;
    private String password;

    public int getSocketConnectionTimeoutMillis() {
        return socketConnectionTimeoutMillis;
    }

    public int getSocketReadTimeoutMillis() {
        return socketReadTimeoutMillis;
    }

    public int getPoolingHeartbeatIntervalSeconds() {
        return poolingHeartbeatIntervalSeconds;
    }

    public int getPoolingIdleTimeoutSeconds() {
        return poolingIdleTimeoutSeconds;
    }

    public int getPoolingTimeoutMillis() {
        return poolingTimeoutMillis;
    }

    public int getPoolingLocalCoreConnectionsPerHost() {
        return poolingLocalCoreConnectionsPerHost;
    }

    public int getPoolingLocalMaxConnectionsPerHost() {
        return poolingLocalMaxConnectionsPerHost;
    }

    public int getPoolingRemoteCoreConnectionsPerHost() {
        return poolingRemoteCoreConnectionsPerHost;
    }

    public int getPoolingRemoteMaxConnectionsPerHost() {
        return poolingRemoteMaxConnectionsPerHost;
    }

    public int getPoolingLocalMaxRequestsPerConnection() {
        return poolingLocalMaxRequestsPerConnection;
    }

    public int getPoolingRemoteMaxRequestsPerConnection() {
        return poolingRemoteMaxRequestsPerConnection;
    }

    public long getReconnectDelayMillis() {
        return reconnectDelayMillis;
    }

    public long getReconnectBaseDelayMillis() {
        return reconnectBaseDelayMillis;
    }

    public long getReconnectMaxDelayMillis() {
        return reconnectMaxDelayMillis;
    }

    public String getRetryPolicyClass() {
        return retryPolicyClass;
    }

    public int getRetryPolicyMaxRetryCount() {
        return retryPolicyMaxRetryCount;
    }

    public int getRetryPolicyGrowingBackOffTimeMillis() {
        return retryPolicyGrowingBackOffTimeMillis;
    }

    public int getRetryPolicyFixedBackOffTimeMillis() {
        return retryPolicyFixedBackOffTimeMillis;
    }

    public String getReadDatacenter() {
        return readDatacenter;
    }

    public String getWriteDatacenter() {
        return writeDatacenter;
    }

    public String getGlobalEndpoint() {
        return globalEndpoint;
    }

    public int getDnsExpirationInSeconds() {
        return dnsExpirationInSeconds;
    }

    protected String getLocalDataCenter() {
        return localDatacenter;
    }

    protected String getKeyspaceName() {
        return keyspace;
    }
    protected String getContactPoints() {
        return contactPoints;
    }
    protected int getPort() {
        return port;
    }
    public SchemaAction getSchemaAction() {
        return SchemaAction.NONE;
    }

    public boolean isSslEnabled() {
        return sslEnabled;
    }

    public String getUsername(){
        return username;
    }

    public String getPassword(){
        return password;
    }

    protected boolean getMetricsEnabled() { return false; }

    @Bean
    @NonNull
    public CqlSessionFactoryBean cassandraSession() {
        final CqlSessionFactoryBean cqlSessionFactoryBean = new CqlSessionFactoryBean();
        cqlSessionFactoryBean.setContactPoints(contactPoints);
        cqlSessionFactoryBean.setKeyspaceName(keyspace);
        cqlSessionFactoryBean.setLocalDatacenter(localDatacenter);
        cqlSessionFactoryBean.setPort(port);
        cqlSessionFactoryBean.setUsername(configManager.getProperty(usernameKey));
        cqlSessionFactoryBean.setPassword(configManager.getProperty(passKey));
        return cqlSessionFactoryBean;
    }
}

问题排查与解决方案

核心问题点

  1. 配置逻辑冲突:继承AbstractReactiveCassandraConfiguration时,手动创建CqlSessionFactoryBean会和父类的自动配置逻辑冲突,导致连接初始化异常。
  2. 配置项未生效:代码中注入了大量Cassandra连接参数,但未在会话构建时应用这些参数,导致自定义配置缺失关键连接属性,和默认配置不一致。
  3. 拼写错误导致参数失效:原代码中socket-read-timeout-mills的mills应为millis,导致超时参数无法正确注入。
  4. SSL配置缺失:未根据sslEnabled参数启用SSL连接,若集群要求SSL则会连接失败。

修复后的配置代码

@Configuration
@EnableReactiveCassandraRepositories(basePackages = {"com.ey.nexus.map.viewer.repository"})
@EntityScan({"com.ey.nexus.map.viewer.model"})
@RequiredArgsConstructor
public class CassandraConfiguration extends AbstractReactiveCassandraConfiguration {

    private final KeyVaultConfigManager configManager;
    @Value("${spring.data.cassandra.keyspace-name}")
    private String keyspace;
    @Value("${spring.data.cassandra.contact-points}")
    private String contactPoints;
    @Value("${spring.data.cassandra.port}")
    private Integer port;
    @Value("${spring.data.cassandra.ssl}")
    private boolean sslEnabled;
    @Value("${spring.data.cassandra.username}")
    private String usernameKey;
    @Value("${spring.data.cassandra.password}")
    private String passKey;
    @Value("${spring.data.cassandra.local-datacenter}")
    private String localDatacenter;
    @Value("${spring.data.cassandra.socket-connection-timeout-millis}")
    private int socketConnectionTimeoutMillis;
    @Value("${spring.data.cassandra.socket-read-timeout-millis}")
    private int socketReadTimeoutMillis;
    @Value("${spring.data.cassandra.pooling-heartbeat-interval-seconds}")
    private int poolingHeartbeatIntervalSeconds;
    @Value("${spring.data.cassandra.pooling-idle-timeout-seconds}")
    private int poolingIdleTimeoutSeconds;
    @Value("${spring.data.cassandra.pooling-timeout-millis}")
    private int poolingTimeoutMillis;
    @Value("${spring.data.cassandra.pooling-options.local.core-connections-per-host}")
    private int poolingLocalCoreConnectionsPerHost;
    @Value("${spring.data.cassandra.pooling-options.local.max-connections-per-host}")
    private int poolingLocalMaxConnectionsPerHost;
    @Value("${spring.data.cassandra.pooling-options.remote.core-connections-per-host}")
    private int poolingRemoteCoreConnectionsPerHost;
    @Value("${spring.data.cassandra.pooling-options.remote.max-connections-per-host}")
    private int poolingRemoteMaxConnectionsPerHost;
    @Value("${spring.data.cassandra.pooling-options.local.max-requests-per-connection}")
    private int poolingLocalMaxRequestsPerConnection;
    @Value("${spring.data.cassandra.pooling-options.remote.max-requests-per-connection}")
    private int poolingRemoteMaxRequestsPerConnection;
    @Value("${spring.data.cassandra.reconnect-base-delay-millis}")
    private long reconnectBaseDelayMillis;
    @Value("${spring.data.cassandra.reconnect-max-delay-millis}")
    private long reconnectMaxDelayMillis;
    @Value("${spring.data.cassandra.retry-policy.class}")
    private String retryPolicyClass;

    @Override
    protected String getLocalDataCenter() {
        return localDatacenter;
    }

    @Override
    protected String getKeyspaceName() {
        return keyspace;
    }

    @Override
    protected String getContactPoints() {
        return contactPoints;
    }

    @Override
    protected int getPort() {
        return port;
    }

    @Override
    public SchemaAction getSchemaAction() {
        return SchemaAction.NONE;
    }

    @Override
    protected boolean getMetricsEnabled() {
        return false;
    }

    @Override
    protected SessionBuilderConfigurer getSessionBuilderConfigurer() {
        return sessionBuilder -> {
            // 加载KeyVault中的认证凭证
            sessionBuilder.withAuthCredentials(
                    configManager.getProperty(usernameKey),
                    configManager.getProperty(passKey)
            );

            // 启用SSL连接
            if (sslEnabled) {
                sessionBuilder.withSsl();
            }

            // 配置Socket超时参数
            sessionBuilder.withSocketOptions(SocketOptions.builder()
                    .setConnectTimeoutMillis(socketConnectionTimeoutMillis)
                    .setReadTimeoutMillis(socketReadTimeoutMillis)
                    .build());

            // 配置连接池参数
            PoolingOptions poolingOptions = PoolingOptions.builder()
                    .setHeartbeatIntervalSeconds(poolingHeartbeatIntervalSeconds)
                    .setIdleTimeoutSeconds(poolingIdleTimeoutSeconds)
                    .setTimeoutMillis(poolingTimeoutMillis)
                    .setCoreConnectionsPerHost(HostDistance.LOCAL, poolingLocalCoreConnectionsPerHost)
                    .setMaxConnectionsPerHost(HostDistance.LOCAL, poolingLocalMaxConnectionsPerHost)
                    .setMaxRequestsPerConnection(HostDistance.LOCAL, poolingLocalMaxRequestsPerConnection)
                    .setCoreConnectionsPerHost(HostDistance.REMOTE, poolingRemoteCoreConnectionsPerHost)
                    .setMaxConnectionsPerHost(HostDistance.REMOTE, poolingRemoteMaxConnectionsPerHost)
                    .setMaxRequestsPerConnection(HostDistance.REMOTE, poolingRemoteMaxRequestsPerConnection)
                    .build();
            sessionBuilder.withPoolingOptions(poolingOptions);

            // 配置重连策略
            sessionBuilder.withReconnectionPolicy(
                    new ExponentialReconnectionPolicy(reconnectBaseDelayMillis, reconnectMaxDelayMillis)
            );

            // 配置重试策略
            try {
                Class<? extends RetryPolicy> retryPolicyClassInstance = (Class<? extends RetryPolicy>) Class.forName(retryPolicyClass);
                RetryPolicy retryPolicy = retryPolicyClassInstance.getDeclaredConstructor().newInstance();
                sessionBuilder.withRetryPolicy(retryPolicy);
            } catch (Exception e) {
                throw new RuntimeException("初始化重试策略失败", e);
            }

            return sessionBuilder;
        };
    }
}

关键修复说明

  1. 移除手动创建的CqlSessionFactoryBean:通过重写getSessionBuilderConfigurer()方法统一配置会话参数,避免和父类配置逻辑冲突。
  2. 补全所有配置项:将注入的连接池、超时、重连、重试等参数全部应用到SessionBuilder,确保自定义配置和默认配置参数一致。
  3. 修正参数拼写错误:将socket-read-timeout-mills改为socket-read-timeout-millis,确保超时参数正常注入。
  4. 动态启用SSL:根据配置的sslEnabled参数开启SSL连接,匹配集群要求。
  5. 直接加载KeyVault凭证:通过sessionBuilder.withAuthCredentials()传入从KeyVault获取的用户名密码,确保认证生效。

额外排查建议

  • 查看应用启动日志,定位具体错误信息(如连接超时、认证失败、配置缺失等)。
  • 验证KeyVaultConfigManager能正确获取用户名和密码,可在启动前打印日志确认。
  • 对比默认自动配置的参数,确保自定义配置的所有属性和默认值一致(除凭证加载方式外)。

内容的提问来源于stack exchange,提问作者shivakumar pattanashetti

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 01:14:56