Spring响应式Cassandra启动失败:抛出AllNodesFailedException异常
问题背景
我构建了一个Spring响应式应用连接Cassandra,但响应式方案始终无法正常工作,非响应式方案却能正常连接。以下是两种方案的代码对比及问题详情:
方案1:响应式Cassandra(无法正常工作)
Maven依赖
<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-cassandra-reactive</artifactId> <version>3.0</version> </dependency>
application.yml配置
spring.data.cassandra.contact-points=<my connection points> spring.data.cassandra.username=abc spring.data.cassandra.password=xyz spring.data.cassandra.local-datacenter=datacenter1 spring.data.cassandra.keyspace-name=mykeyspace spring.data.cassandra.port=9042
配置类
@Configuration @EnableReactiveCassandraRepositories public class LocalBeanConfig extends AbstractReactiveCassandraConfiguration { @Override protected String getKeyspaceName() { return "mykeyspace"; } }
实体类
@AllArgsConstructor @Table("test_table") public class TableClass { @Getter @PrimaryKey @Column("id") private String id; @Column("value") @Getter private String value; }
仓库接口
@Repository public interface TestAppRepository extends ReactiveCassandraRepository<TableClass, String> { }
数据访问代码
public class MyService { private final TestAppRepository testRepository; public void get(CrawlTask crawlTask) { testRepository.findById("id_1").map(testAppConfig -> test1(testAppConfig)).switchIfEmpty(Mono.just(test2())); } }
异常信息
应用启动时抛出以下异常:
Caused by: com.datastax.oss.driver.api.core.AllNodesFailedException
方案2:非响应式Cassandra(正常工作)
Maven依赖
<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-cassandra</artifactId> </dependency>
application.yml配置
spring.cassandra.contact-points=<my connection points> spring.cassandra.username=abc spring.cassandra.password=xyz spring.cassandra.local-datacenter=datacenter1 spring.cassandra.keyspace-name=mykeyspace spring.cassandra.port=9042
实体类
@AllArgsConstructor @Table("test_table") public class TableClass { @Getter @PrimaryKey @Column("id") private String id; @Column("value") @Getter private String value; }
仓库接口
public interface TestAppRepository extends CassandraRepository<TableClass, String> { }
数据访问代码
public class MyService { private final TestAppRepository testRepository; public void get(CrawlTask crawlTask) { Mono.just(testRepository.findById("id_1").map(testAppConfig -> test1(testAppConfig)).orElse(test2())); } }
已尝试的操作与疑问
- 两种方案配置前缀不同:非响应式用
spring.cassandra,响应式用spring.data.cassandra - 尝试在响应式方案中使用
spring.cassandra前缀,服务能启动但无法查询到数据 - 不清楚为何非响应式正常,响应式却连接失败,遗漏了什么配置?
编辑补充:完整错误日志
s.d.r.c.RepositoryConfigurationDelegate : Bootstrapping Spring Data Reactive Cassandra repositories in DEFAULT mode. .s.d.r.c.RepositoryConfigurationDelegate : Finished Spring Data repository scanning in 65 ms. Found 1 Reactive Cassandra repository interfaces. .s.d.r.c.RepositoryConfigurationDelegate : Bootstrapping Spring Data Cassandra repositories in DEFAULT mode. .s.d.r.c.RepositoryConfigurationDelegate : Finished Spring Data repository scanning in 4 ms. Found 0 Cassandra repository interfaces. c.d.o.d.i.core.DefaultMavenCoordinates : DataStax Java driver for Apache Cassandra(R) (com.datastax.oss:java-driver-core) version 4.15.0 c.d.oss.driver.internal.core.time.Clock : Using native clock for microsecond precision c.d.o.d.i.c.control.ControlConnection : [s0] Error connecting to Node(endPoint=/127.0.0.1:9042, hostId=null, hashCode=3f2701f1), trying next node (ConnectionInitException: [s0|control|connecting...] Protocol initialization request, step 1 (OPTIONS): failed to send request
问题解决分析
核心问题原因
配置类与自动配置冲突
自定义的LocalBeanConfig继承了AbstractReactiveCassandraConfiguration,这个类会完全接管Cassandra的配置逻辑,导致application.yml中spring.data.cassandra前缀的配置无法被自动加载。而非响应式方案没有自定义配置类,Spring Boot自动配置能正确读取spring.cassandra的配置项。响应式流未订阅导致数据查询失败
当你尝试用spring.cassandra前缀启动响应式服务时,服务能启动但查不到数据,是因为响应式流必须被订阅才会执行,你的get方法中只构建了Mono流但没有触发订阅操作。
具体解决步骤
方案A:移除自定义配置类,依赖自动配置
直接删除LocalBeanConfig类,Spring Boot会自动根据spring.data.cassandra前缀的配置创建响应式Cassandra连接,这是最符合Spring Boot设计理念的方式。
方案B:保留配置类,手动注入配置属性
如果必须保留自定义配置类,需要通过@ConfigurationProperties绑定配置项:
@Configuration @EnableReactiveCassandraRepositories @ConfigurationProperties(prefix = "spring.data.cassandra") public class LocalBeanConfig extends AbstractReactiveCassandraConfiguration { private String keyspaceName; private String contactPoints; private int port; private String localDatacenter; private String username; private String password; @Override protected String getKeyspaceName() { return keyspaceName; } @Override protected String getContactPoints() { return contactPoints; } @Override protected int getPort() { return port; } @Override protected String getLocalDatacenter() { return localDatacenter; } @Override protected DriverConfigLoaderBuilderConfigurer getDriverConfigLoaderBuilderConfigurer() { return builder -> builder .withString(DefaultDriverOption.AUTH_PROVIDER_CLASS, PlainTextAuthProvider.class.getName()) .withString(DefaultDriverOption.AUTH_USERNAME, username) .withString(DefaultDriverOption.AUTH_PASSWORD, password); } // 生成对应的getter和setter方法 public void setKeyspaceName(String keyspaceName) { this.keyspaceName = keyspaceName; } public void setContactPoints(String contactPoints) { this.contactPoints = contactPoints; } public void setPort(int port) { this.port = port; } public void setLocalDatacenter(String localDatacenter) { this.localDatacenter = localDatacenter; } public void setUsername(String username) { this.username = username; } public void setPassword(String password) { this.password = password; } }
同时保持application.yml使用spring.data.cassandra前缀的配置。
方案C:使用spring.cassandra前缀配合自动配置
如果想用spring.cassandra前缀,无需自定义配置类,直接依赖自动配置即可。同时修改数据访问代码,确保响应式流被订阅:
// 方式1:直接在方法内订阅 public void get(CrawlTask crawlTask) { testRepository.findById("id_1") .map(testAppConfig -> test1(testAppConfig)) .switchIfEmpty(Mono.just(test2())) .subscribe(); // 必须订阅才会执行响应式流 } // 方式2:返回Mono让上层调用方订阅 public Mono<YourReturnType> get(CrawlTask crawlTask) { return testRepository.findById("id_1") .map(testAppConfig -> test1(testAppConfig)) .switchIfEmpty(Mono.just(test2())); }
额外注意点
- Spring Boot 3.0版本建议移除starter的
version标签,让parent依赖管理版本,避免版本冲突。 - 检查Cassandra节点的网络连通性,确保应用能访问到
contact-points的9042端口,响应式驱动的连接逻辑和非响应式一致,网络问题也可能导致连接失败。
内容的提问来源于stack exchange,提问作者mang4521

