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

使用Spring Apache Kafka将非Confluent版Kafka Topic数据导入Cassandra

结论先行

原生Apache Kafka 完全支持该需求,无需依赖任何Confluent专属组件,仅通过Spring官方提供的Spring for Apache Kafka和Spring Data for Apache Cassandra两个模块即可实现全流程开发。


实现步骤

1. 引入依赖

不需要添加任何Confluent相关组件,仅在pom.xml(Maven)中新增Spring Data Cassandra依赖即可,你已引入的Spring Kafka依赖无需修改:

<!-- Spring Data Cassandra 核心依赖 -->
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-data-cassandra</artifactId>
</dependency>

2. 基础配置

在application.yml中添加Cassandra配置,你现有原生Kafka的配置无需调整,不需要新增任何Confluent相关配置项:

spring:
  # Cassandra 连接配置
  cassandra:
    keyspace-name: 你的KeySpace名称
    contact-points: 127.0.0.1
    port: 9042
    username: cassandra
    password: cassandra
    local-datacenter: datacenter1
  # 你现有原生Kafka配置保持不变即可
  kafka:
    bootstrap-servers: 127.0.0.1:9092
    consumer:
      group-id: kafka-to-cassandra-group
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      auto-offset-reset: earliest

3. 核心编码实现

定义Cassandra实体类

对应你要存储的业务数据结构:

import org.springframework.data.cassandra.core.mapping.PrimaryKey;
import org.springframework.data.cassandra.core.mapping.Table;

@Table("user_input") // 对应Cassandra中的表名
public class UserInput {
    @PrimaryKey
    private String id;
    private String content; // 用户输入内容
    private Long createTime;
    // 构造方法、getter、setter自行补充
}

定义Cassandra操作Repository

import org.springframework.data.cassandra.repository.CassandraRepository;

public interface UserInputRepository extends CassandraRepository<UserInput, String> {
}

Kafka消费者写入Cassandra

直接在你已有的Kafka消费者逻辑中加入写入逻辑即可:

import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
import com.fasterxml.jackson.databind.ObjectMapper;

@Component
public class KafkaSyncConsumer {

    private final ObjectMapper objectMapper;
    private final UserInputRepository userInputRepository;

    // 构造注入依赖
    public KafkaSyncConsumer(ObjectMapper objectMapper, UserInputRepository userInputRepository) {
        this.objectMapper = objectMapper;
        this.userInputRepository = userInputRepository;
    }

    @KafkaListener(topics = "你的Kafka Topic名称", groupId = "kafka-to-cassandra-group")
    public void handleMessage(String message) {
        // 反序列化Kafka消息为实体对象
        UserInput userInput = objectMapper.readValue(message, UserInput.class);
        // 写入Cassandra,可自行补充异常处理、重试、死信队列逻辑
        userInputRepository.save(userInput);
    }
}

可选优化方案
  • 数据一致性保障:开启Kafka消费者手动ACK,写入Cassandra成功后再提交offset,避免数据丢失,仅需在配置中新增spring.kafka.consumer.enable-auto-commit: false,再调整消费者代码手动提交offset即可
  • 高吞吐场景优化:开启Kafka批量消费,批量写入Cassandra,大幅提升同步性能

内容的提问来源于stack exchange,提问作者vigneshwar reddy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 23:06:01