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

如何基于Spring Batch实现Kafka到Cassandra的高吞吐数据存储?

Spring Batch Kafka到Cassandra百万级数据存储优化方案

问题背景

需求:使用Spring Batch从Kafka消费包含String类型主键id与int类型状态码的对象,存储至Cassandra,需支撑每分钟百万级消息量的高效存储。

当前实现架构:

  • Kafka Listener批量消费消息后存入ConcurrentLinkedQueue
  • Spring Batch的ItemReader单条轮询队列获取数据,ItemWriter批量写入Cassandra
  • Topic配置:5个分区、2个副本
  • 硬件:Ryzen 7 4800H八核处理器、16GB RAM、SSD硬盘

遇到的问题:

  • 提升并发后步骤耗时反而增加
  • 当前存储100万条消息需2分钟,怀疑ConcurrentLinkedQueue是性能瓶颈
  • 曾尝试暂停Kafka Listener拆分批次入队,出现消息丢失问题

核心优化方案

1. 替换单条轮询的ItemReader,实现批量读取队列数据

当前ItemReader每次poll()单条数据,频繁的队列操作带来锁竞争和上下文切换开销,改成批量读取可大幅减少Reader调用次数:

@StepScope
public ItemReader<List<Message>> batchItemReader() {
    return new ItemReader<>() {
        @Override
        public List<Message> read() throws Exception {
            List<Message> batch = new ArrayList<>(1000); // 预分配容量,减少扩容开销
            Message item;
            // 批量读取,直到凑够目标大小或队列为空
            for (int i = 0; i < 1000; i++) {
                item = dataFromKafka.poll();
                if (item == null) {
                    break;
                }
                batch.add(item);
            }
            return batch.isEmpty() ? null : batch;
        }
    };
}

同时调整Step的chunk配置,确保与批量读取的大小匹配,避免资源浪费。

2. 优化Kafka消费并发,匹配Topic分区数

当前Kafka Listener的concurrency=1,无法充分利用Topic的5个分区并行性。将concurrency设置为与分区数一致,每个Listener线程对应一个分区,避免分区竞争:

@KafkaListener(
        id = "topic-kafka-listener",
        groupId = "topic-batch",
        containerFactory = "kafkaListenerContainerFactory",
        concurrency = "5", // 匹配Topic分区数
        topics = "topic"
)
public void receive(@NotNull @Payload List<Message> messages) {
    dataFromKafka.addAll(messages);
}

额外配置Kafka消费者的max.poll.records为1000-5000,减少网络请求次数,提升单次消费效率。

3. 替换ConcurrentLinkedQueue为更高效的批量缓存容器

ConcurrentLinkedQueue是单链表结构,批量操作效率低,推荐使用以下容器:

  • LinkedBlockingQueue:支持drainTo()方法一次性取出多个元素,大幅减少锁竞争:
// 初始化时设置合理容量,实现背压控制
LinkedBlockingQueue<Message> dataFromKafka = new LinkedBlockingQueue<>(100000);

@StepScope
public ItemReader<List<Message>> batchItemReader() {
    return () -> {
        List<Message> batch = new ArrayList<>(1000);
        dataFromKafka.drainTo(batch, 1000);
        return batch.isEmpty() ? null : batch;
    };
}
  • Disruptor:如果追求极致性能,LMAX Disruptor的环形队列结构锁开销远低于传统并发队列,适合高吞吐量场景。

4. 优化Cassandra写入性能

Cassandra的批量写入是性能核心,需从多维度优化:

  • 调整批量写入大小:当前chunk设为10000,可根据Cassandra配置调整为2000-5000,避免单批次过大导致超时。
  • 使用异步批量写入:改用异步API减少线程等待时间:
@StepScope
public ItemWriter<Message> asyncItemWriter() {
    return chunk -> {
        List<CompletableFuture<Void>> futures = chunk.getItems().stream()
                .map(message -> messageRepository.saveAsync(message))
                .collect(Collectors.toList());
        CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
    };
}
  • Cassandra集群配置调优:使用LOCAL_QUORUM一致性级别,开启unlogged batch(无事务场景),调整concurrent_writes等参数提升写入吞吐量。

5. 重构Spring Batch执行逻辑,避免定时任务触发Job

当前每5秒启动一次Job会带来频繁启停的额外开销,改为长运行Step持续处理数据:

public Step step1(JobRepository jobRepository, PlatformTransactionManager transactionManager) {
    return new StepBuilder("step1", jobRepository)
            .<List<Message>, Message>chunk(1000, transactionManager)
            .taskExecutor(threadPoolTaskExecutor())
            .reader(batchItemReader())
            .writer(asyncItemWriter())
            .allowStartIfComplete(true) // 允许Step重复执行
            .build();
}

// 启动一次Job后持续运行
JobParameters jobParameters = new JobParametersBuilder()
        .addLong("startTime", System.currentTimeMillis())
        .toJobParameters();
jobLauncher.run(job, jobParameters);

同时修改ItemReader,当队列为空时短暂休眠(如100ms),避免空轮询占用CPU。

6. 消除消息丢失风险

暂停Kafka Listener导致消息丢失的核心原因是offset未正确提交,正确做法是:

  • 将KafkaListener的ackMode设置为MANUAL_IMMEDIATE或BATCH,确保消费完成后再提交offset
  • 不暂停Listener,通过设置队列上限实现背压,当队列满时Listener自动阻塞,避免消息堆积

额外性能调优建议

  • 调整线程池参数:核心线程数设为8-12(匹配CPU核心数),队列容量设为1000,避免线程过多导致上下文切换
  • 内存与GC优化:调整JVM堆大小(如-Xms8G -Xmx12G),使用对象池复用Message对象减少GC开销
  • 监控瓶颈定位:用JProfiler或VisualVM分析CPU、内存、锁竞争情况,监控Cassandra写入延迟、吞吐量指标

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 22:02:29