如何在Spring WebFlux应用中连接两个不同的Cassandra集群?
Spring WebFlux 连接多Cassandra集群实现示例
1. 依赖配置
首先在Maven pom.xml 中引入必要依赖:
<dependencies> <!-- Spring WebFlux --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-webflux</artifactId> </dependency> <!-- Reactive Cassandra 数据访问 --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-cassandra-reactive</artifactId> </dependency> </dependencies>
2. 多集群配置文件
在 application.yml 中配置两个Cassandra集群的连接参数:
spring: cassandra: cluster1: contact-points: 192.168.1.100 port: 9042 keyspace-name: cluster1_keyspace local-datacenter: datacenter1 username: cassandra password: cassandra cluster2: contact-points: 192.168.1.101 port: 9042 keyspace-name: cluster2_keyspace local-datacenter: datacenter1 username: cassandra password: cassandra
3. 第一个Cassandra集群配置类
创建独立的配置类,绑定cluster1的配置并初始化对应的Reactive操作模板:
package com.example.config; import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.data.cassandra.ReactiveSessionFactory; import org.springframework.data.cassandra.config.AbstractReactiveCassandraConfiguration; import org.springframework.data.cassandra.config.CqlSessionFactoryBean; import org.springframework.data.cassandra.config.SchemaAction; import org.springframework.data.cassandra.core.ReactiveCassandraTemplate; import org.springframework.data.cassandra.repository.config.EnableReactiveCassandraRepositories; @Configuration @EnableReactiveCassandraRepositories( basePackages = "com.example.repository.cluster1", reactiveCassandraTemplateRef = "cluster1ReactiveTemplate" ) @ConfigurationProperties(prefix = "spring.cassandra.cluster1") public class Cluster1CassandraConfig extends AbstractReactiveCassandraConfiguration { private String contactPoints; private int port; private String keyspaceName; private String localDatacenter; private String username; private String password; // Getter和Setter方法(可通过Lombok简化) public String getContactPoints() { return contactPoints; } public void setContactPoints(String contactPoints) { this.contactPoints = contactPoints; } public int getPort() { return port; } public void setPort(int port) { this.port = port; } public String getKeyspaceName() { return keyspaceName; } public void setKeyspaceName(String keyspaceName) { this.keyspaceName = keyspaceName; } public String getLocalDatacenter() { return localDatacenter; } public void setLocalDatacenter(String localDatacenter) { this.localDatacenter = localDatacenter; } public String getUsername() { return username; } public void setUsername(String username) { this.username = username; } public String getPassword() { return password; } public void setPassword(String password) { this.password = password; } @Bean public CqlSessionFactoryBean cluster1Session() { CqlSessionFactoryBean session = new CqlSessionFactoryBean(); session.setContactPoints(contactPoints); session.setPort(port); session.setKeyspaceName(keyspaceName); session.setLocalDatacenter(localDatacenter); session.setUsername(username); session.setPassword(password); return session; } @Bean("cluster1ReactiveTemplate") public ReactiveCassandraTemplate cluster1ReactiveTemplate(ReactiveSessionFactory sessionFactory) { return new ReactiveCassandraTemplate(sessionFactory); } @Override protected String getKeyspaceName() { return keyspaceName; } @Override protected SchemaAction getSchemaAction() { return SchemaAction.NONE; // 根据需求调整,例如CREATE_IF_NOT_EXISTS } }
4. 第二个Cassandra集群配置类
同理,创建cluster2的配置类,注意包路径和模板引用名需区分:
package com.example.config; import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.data.cassandra.ReactiveSessionFactory; import org.springframework.data.cassandra.config.AbstractReactiveCassandraConfiguration; import org.springframework.data.cassandra.config.CqlSessionFactoryBean; import org.springframework.data.cassandra.config.SchemaAction; import org.springframework.data.cassandra.core.ReactiveCassandraTemplate; import org.springframework.data.cassandra.repository.config.EnableReactiveCassandraRepositories; @Configuration @EnableReactiveCassandraRepositories( basePackages = "com.example.repository.cluster2", reactiveCassandraTemplateRef = "cluster2ReactiveTemplate" ) @ConfigurationProperties(prefix = "spring.cassandra.cluster2") public class Cluster2CassandraConfig extends AbstractReactiveCassandraConfiguration { private String contactPoints; private int port; private String keyspaceName; private String localDatacenter; private String username; private String password; // Getter和Setter方法 public String getContactPoints() { return contactPoints; } public void setContactPoints(String contactPoints) { this.contactPoints = contactPoints; } public int getPort() { return port; } public void setPort(int port) { this.port = port; } public String getKeyspaceName() { return keyspaceName; } public void setKeyspaceName(String keyspaceName) { this.keyspaceName = keyspaceName; } public String getLocalDatacenter() { return localDatacenter; } public void setLocalDatacenter(String localDatacenter) { this.localDatacenter = localDatacenter; } public String getUsername() { return username; } public void setUsername(String username) { this.username = username; } public String getPassword() { return password; } public void setPassword(String password) { this.password = password; } @Bean public CqlSessionFactoryBean cluster2Session() { CqlSessionFactoryBean session = new CqlSessionFactoryBean(); session.setContactPoints(contactPoints); session.setPort(port); session.setKeyspaceName(keyspaceName); session.setLocalDatacenter(localDatacenter); session.setUsername(username); session.setPassword(password); return session; } @Bean("cluster2ReactiveTemplate") public ReactiveCassandraTemplate cluster2ReactiveTemplate(ReactiveSessionFactory sessionFactory) { return new ReactiveCassandraTemplate(sessionFactory); } @Override protected String getKeyspaceName() { return keyspaceName; } @Override protected SchemaAction getSchemaAction() { return SchemaAction.NONE; } }
5. Repository层定义
将两个集群的Repository分别放在指定包下:
Cluster1的Repository
package com.example.repository.cluster1; import com.example.entity.User; import org.springframework.data.cassandra.repository.ReactiveCassandraRepository; import reactor.core.publisher.Mono; public interface Cluster1UserRepository extends ReactiveCassandraRepository<User, String> { Mono<User> findByUsername(String username); }
Cluster2的Repository
package com.example.repository.cluster2; import com.example.entity.Order; import org.springframework.data.cassandra.repository.ReactiveCassandraRepository; import reactor.core.publisher.Flux; public interface Cluster2OrderRepository extends ReactiveCassandraRepository<Order, String> { Flux<Order> findByUserId(String userId); }
对应的实体类示例(User):
package com.example.entity; import org.springframework.data.cassandra.core.mapping.PrimaryKey; import org.springframework.data.cassandra.core.mapping.Table; @Table("users") public class User { @PrimaryKey private String id; private String username; private String email; // 构造方法、Getter、Setter public User() {} public User(String id, String username, String email) { this.id = id; this.username = username; this.email = email; } public String getId() { return id; } public void setId(String id) { this.id = id; } public String getUsername() { return username; } public void setUsername(String username) { this.username = username; } public String getEmail() { return email; } public void setEmail(String email) { this.email = email; } }
6. 业务层与控制器示例
在Service中注入两个Repository,实现跨集群业务逻辑:
package com.example.service; import com.example.entity.User; import com.example.entity.Order; import com.example.repository.cluster1.Cluster1UserRepository; import com.example.repository.cluster2.Cluster2OrderRepository; import org.springframework.stereotype.Service; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @Service public class CrossClusterService { private final Cluster1UserRepository userRepository; private final Cluster2OrderRepository orderRepository; // 构造方法注入 public CrossClusterService(Cluster1UserRepository userRepository, Cluster2OrderRepository orderRepository) { this.userRepository = userRepository; this.orderRepository = orderRepository; } public Mono<User> getUserById(String userId) { return userRepository.findById(userId); } public Flux<Order> getUserOrders(String userId) { return orderRepository.findByUserId(userId); } // 跨集群联合查询示例 public Flux<Order> getUserWithOrders(String userId) { return userRepository.findById(userId) .flatMapMany(user -> orderRepository.findByUserId(userId)); } }
对应的WebFlux控制器:
package com.example.controller; import com.example.entity.User; import com.example.entity.Order; import com.example.service.CrossClusterService; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.PathVariable; import org.springframework.web.bind.annotation.RestController; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @RestController public class CrossClusterController { private final CrossClusterService crossClusterService; public CrossClusterController(CrossClusterService crossClusterService) { this.crossClusterService = crossClusterService; } @GetMapping("/users/{userId}") public Mono<User> getUser(@PathVariable String userId) { return crossClusterService.getUserById(userId); } @GetMapping("/users/{userId}/orders") public Flux<Order> getUserOrders(@PathVariable String userId) { return crossClusterService.getUserWithOrders(userId); } }
注意事项
- 确保两个集群的
local-datacenter配置正确,否则会出现连接超时问题。 - Repository的包路径必须与配置类中
basePackages指定的一致,避免Spring扫描混淆。 - 如果需要自定义CQL操作,可以通过
@Qualifier("cluster1ReactiveTemplate")注入对应模板。 - 生产环境建议开启连接池配置,可通过
spring.cassandra.cluster1.pool.*相关属性调整。
内容的提问来源于stack exchange,提问作者Erick Ramirez
相关产品推荐
相关产品推荐

