如何基于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
相关产品推荐
相关产品推荐

