Spring Batch用SynchronizedItemStreamReader实现多线程KafkaItemReader遇异常
核心问题本质
你遇到的ConcurrentModificationException根本原因是KafkaConsumer的线程绑定特性,而非同步锁的问题。SynchronizedItemStreamReader仅能保证单个ItemReader实例被多线程同步访问,但KafkaConsumer内部会严格校验调用线程的ID——只要不同线程先后操作同一个KafkaConsumer实例,哪怕是同步调用,也会触发线程安全异常。另外你提到只生成两个clientId的消费者,说明你在多线程步骤中复用了同一个KafkaItemReader实例,没有为每个线程创建独立的消费者。
你的操作错误点
- 复用单个KafkaItemReader实例:多线程环境下,多个线程共享同一个底层持有KafkaConsumer的KafkaItemReader,即使加了同步锁,跨线程调用KafkaConsumer依然会触发其内部的线程安全检查。
- 误解SynchronizedItemStreamReader的作用:它的设计目标是解决单ItemReader被多线程并发访问时的同步问题,但无法适配KafkaConsumer这种天生线程绑定的组件场景。
正确的多线程消费方案
方案1:为每个线程创建独立的KafkaItemReader
利用Spring Batch的@StepScope注解,让每个线程获取独立的KafkaItemReader实例,每个实例持有专属的KafkaConsumer:
@Bean @StepScope public KafkaItemReader<Integer> kafkaItemReader() { return new KafkaItemReaderBuilder<Integer>() .consumerFactory(consumerFactory()) .topics("your-numeric-topic") .partitionOffsets(new PartitionOffset("your-numeric-topic", 0, 0L)) .build(); } @Bean public TaskExecutor taskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(8); executor.setMaxPoolSize(8); executor.setThreadNamePrefix("Multi-No."); executor.initialize(); return executor; } @Bean public Step multiThreadedKafkaStep(JobRepository jobRepository, PlatformTransactionManager transactionManager) { return new StepBuilder("multiThreadedKafkaStep", jobRepository) .<Integer, Integer>chunk(10, transactionManager) .reader(kafkaItemReader()) // 每个线程会获取新的@StepScope实例 .writer(yourItemWriter()) .taskExecutor(taskExecutor()) .build(); }
这样每个线程都会拥有独立的KafkaConsumer,clientId也会自动区分,从根本上避免线程安全问题。
方案2:使用分区步骤消费多分区
如果你的Kafka主题有多个分区,更推荐用Spring Batch的分区步骤(Partitioned Step),每个分区分配一个独立的Step执行实例,每个实例消费指定分区的数据:
- 自定义
Partitioner来拆分Kafka的分区作为步骤的分区源 - 通过
PartitionHandler管理多个分区的并行执行 - 每个分区对应的Step实例持有独立的KafkaItemReader,天然规避线程安全问题
补充说明
KafkaConsumer的线程安全校验是强校验,它会记录首次调用的线程ID,后续任何不同线程的调用都会直接抛出异常,这和是否加同步锁无关。SynchronizedItemStreamReader无法绕过这个校验,因为它只是保证同一时间一个线程访问,无法改变KafkaConsumer的线程绑定特性。
内容的提问来源于stack exchange,提问作者Hyunjin Staats

