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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 11:37:05