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; } }
问题排查与解决方案
核心问题点
- 配置逻辑冲突:继承
AbstractReactiveCassandraConfiguration时,手动创建CqlSessionFactoryBean会和父类的自动配置逻辑冲突,导致连接初始化异常。 - 配置项未生效:代码中注入了大量Cassandra连接参数,但未在会话构建时应用这些参数,导致自定义配置缺失关键连接属性,和默认配置不一致。
- 拼写错误导致参数失效:原代码中
socket-read-timeout-mills的mills应为millis,导致超时参数无法正确注入。 - 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; }; } }
关键修复说明
- 移除手动创建的
CqlSessionFactoryBean:通过重写getSessionBuilderConfigurer()方法统一配置会话参数,避免和父类配置逻辑冲突。 - 补全所有配置项:将注入的连接池、超时、重连、重试等参数全部应用到
SessionBuilder,确保自定义配置和默认配置参数一致。 - 修正参数拼写错误:将
socket-read-timeout-mills改为socket-read-timeout-millis,确保超时参数正常注入。 - 动态启用SSL:根据配置的
sslEnabled参数开启SSL连接,匹配集群要求。 - 直接加载KeyVault凭证:通过
sessionBuilder.withAuthCredentials()传入从KeyVault获取的用户名密码,确保认证生效。
额外排查建议
- 查看应用启动日志,定位具体错误信息(如连接超时、认证失败、配置缺失等)。
- 验证
KeyVaultConfigManager能正确获取用户名和密码,可在启动前打印日志确认。 - 对比默认自动配置的参数,确保自定义配置的所有属性和默认值一致(除凭证加载方式外)。
内容的提问来源于stack exchange,提问作者shivakumar pattanashetti
相关产品推荐
相关产品推荐

