num.stream.threads>1时偶发ConcurrentModificationException问题求助
Kafka Streams多线程场景下ConcurrentModificationException排查指南
当Kafka Streams的num.stream.threads参数设置大于1时,流处理流程偶尔会抛出如下ConcurrentModificationException异常:
[wrwks-bef-projekt-aggregat-wirtschaftseinheitAndMietobjektUpdater-2577dae3-7c43-4782-bdd8-e51669a18469-StreamThread-6] ERROR o.a.k.s.p.internals.TaskManager - stream-thread [wrwks-bef-projekt-aggregat-wirtschaftseinheitAndMietobjektUpdater-2577dae3-7c43-4782-bdd8-e51669a18469-StreamThread-6] Failed to process stream task 1_2 due to the following error: org.apache.kafka.streams.errors.StreamsException: Exception caught in process. taskId=1_2, processor=KSTREAM-SOURCE-0000000011, topic=wrw-technischerplatz-mietobjekt-aggregat-oeffentlich-1, partition=2, offset=7263, stacktrace=java.util.ConcurrentModificationException at java.base/java.util.ArrayList$Itr.checkForComodification(Unknown Source) at java.base/java.util.ArrayList$Itr.next(Unknown Source) at org.apache.kafka.common.header.internals.RecordHeaders$FilterByKeyIterator.makeNext(RecordHeaders.java:184) at org.apache.kafka.common.header.internals.RecordHeaders$FilterByKeyIterator.makeNext(RecordHeaders.java:171) at org.apache.kafka.common.utils.AbstractIterator.maybeComputeNext(AbstractIterator.java:79) at org.apache.kafka.common.utils.AbstractIterator.hasNext(AbstractIterator.java:45) at org.springframework.cloud.stream.binder.kafka.streams.AbstractKafkaStreamsBinderProcessor$3.process(AbstractKafkaStreamsBinderProcessor.java:640) at org.apache.kafka.streams.processor.internals.ProcessorNode.process(ProcessorNode.java:146) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forwardInternal(ProcessorContextImpl.java:253) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:232) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:191) at org.apache.kafka.streams.kstream.internals.KStreamMapValues$KStreamMapProcessor.process(KStreamMapValues.java:42) at org.apache.kafka.streams.processor.internals.ProcessorNode.process(ProcessorNode.java:146) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forwardInternal(ProcessorContextImpl.java:253) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:232) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:191) at org.apache.kafka.streams.kstream.internals.KTableSource$KTableSourceProcessor.process(KTableSource.java:152) at org.apache.kafka.streams.processor.internals.ProcessorNode.process(ProcessorNode.java:146) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forwardInternal(ProcessorContextImpl.java:253) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:232) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:191) at org.apache.kafka.streams.processor.internals.SourceNode.process(SourceNode.java:84) at org.apache.kafka.streams.processor.internals.StreamTask.lambda$process$1(StreamTask.java:731) at org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.maybeMeasureLatency(StreamsMetricsImpl.java:809) at org.apache.kafka.streams.processor.internals.StreamTask.process(StreamTask.java:731) at org.apache.kafka.streams.processor.internals.TaskManager.process(TaskManager.java:1296) at org.apache.kafka.streams.processor.internals.StreamThread.runOnce(StreamThread.java:784) at org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:604) at org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:576) at org.apache.kafka.streams.processor.internals.StreamTask.process(StreamTask.java:758) at org.apache.kafka.streams.processor.internals.TaskManager.process(TaskManager.java:1296) at org.apache.kafka.streams.processor.internals.StreamThread.runOnce(StreamThread.java:784) at org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:604) at org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:576) Caused by: java.util.ConcurrentModificationException: null at java.base/java.util.ArrayList$Itr.checkForComodification(Unknown Source) at java.base/java.util.ArrayList$Itr.next(Unknown Source) at org.apache.kafka.common.header.internals.RecordHeaders$FilterByKeyIterator.makeNext(RecordHeaders.java:184) at org.apache.kafka.common.header.internals.RecordHeaders$FilterByKeyIterator.makeNext(RecordHeaders.java:171) at org.apache.kafka.common.utils.AbstractIterator.maybeComputeNext(AbstractIterator.java:79) at org.apache.kafka.common.utils.AbstractIterator.hasNext(AbstractIterator.java:45) at org.springframework.cloud.stream.binder.kafka.streams.AbstractKafkaStreamsBinderProcessor$3.process(AbstractKafkaStreamsBinderProcessor.java:640) at org.apache.kafka.streams.processor.internals.ProcessorNode.process(ProcessorNode.java:146) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forwardInternal(ProcessorContextImpl.java:253) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:232) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:191) at org.apache.kafka.streams.kstream.internals.KStreamMapValues$KStreamMapProcessor.process(KStreamMapValues.java:42) at org.apache.kafka.streams.processor.internals.ProcessorNode.process(ProcessorNode.java:146) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forwardInternal(ProcessorContextImpl.java:253) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:232) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:191) at org.apache.kafka.streams.kstream.internals.KTableSource$KTableSourceProcessor.process(KTableSource.java:152) at org.apache.kafka.streams.processor.internals.ProcessorNode.process(ProcessorNode.java:146) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forwardInternal(ProcessorContextImpl.java:253) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:232) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:191) at org.apache.kafka.streams.processor.internals.SourceNode.process(SourceNode.java:84) at org.apache.kafka.streams.processor.internals.StreamTask.lambda$process$1(StreamTask.java:731) at org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.maybeMeasureLatency(StreamsMetricsImpl.java:809) at org.apache.kafka.streams.processor.internals.StreamTask.process(StreamTask.java:731) ... 4 common frames omitted
目前无法确定该异常是由Kafka Streams客户端Bug还是业务代码导致,相关业务代码如下:
@Configuration class WirtschaftseinheitAndMietobjektUpdaterStreamConfiguration { @Bean fun wirtschaftseinheitAndMietobjektUpdater() = Function { projekte: KTable<String, ProjektAggregat> -> Function { mietobjekte: KTable<String, MietobjektAggregat> -> Function { wirtschaftseinheiten: KTable<String, WirtschaftseinheitAggregat> -> projekte .filter { _, projektAggregat -> projektAggregat.action == AGGREGATE } .leftJoin( mietobjekte, { it.projekt?.projekt?.technischerPlatz }, { projektAggregat, mietobjekt -> if (mietobjekt != null) if (mietobjekt.tplnr.length > WIRTSCHAFTSEINHEIT_LENGTH) (projektAggregat + mietobjekt).copy(action = MO_WE_UPDATE) else projektAggregat else projektAggregat }, ) .leftJoin( wirtschaftseinheiten, { it.projekt?.projekt?.technischerPlatz?.take(WIRTSCHAFTSEINHEIT_LENGTH) }, { projektAggregat, wirtschaftseinheit -> if (wirtschaftseinheit != null) (projektAggregat + wirtschaftseinheit).copy(action = MO_WE_UPDATE) else projektAggregat }, ) .toStream() .filter { _, projektAggregat -> projektAggregat?.action == MO_WE_UPDATE } .transform({ EventTypeHeaderTransformer() }) } } } }
排查提示
- 检查
EventTypeHeaderTransformer的线程安全性:异常栈显示问题出在RecordHeaders的迭代过程中,代码最后调用了.transform({ EventTypeHeaderTransformer() })。如果EventTypeHeaderTransformer内部存在对RecordHeaders的并发修改(比如迭代headers时添加/删除header),或Transformer实例被多线程共享(未实现线程安全),就会触发异常。确认Transformer是否为线程安全实现,每次transform调用是否使用独立实例,或对headers操作时先复制再修改。 - 验证Spring Cloud Stream Kafka Streams绑定器版本:异常栈涉及
AbstractKafkaStreamsBinderProcessor,检查当前使用的Spring Cloud Stream和Kafka Streams绑定器版本是否存在已知并发Bug。部分旧版本绑定器在多线程场景下可能存在RecordHeaders处理的线程不安全问题,尝试升级到最新稳定版。 - 确认聚合对象的不可变性:业务代码中使用了
projektAggregat + mietobjekt和.copy(action = MO_WE_UPDATE),确保ProjektAggregat、MietobjektAggregat等数据类是完全不可变的(Kotlin data class默认不可变,但内部若有可变集合属性需额外注意)。若聚合对象内部存在可变状态且被多线程访问,可能间接引发header处理的并发问题。 - 启用调试日志:开启
org.apache.kafka.streams和org.springframework.cloud.stream.binder.kafka.streams的DEBUG级别日志,观察异常发生前后的线程交互、任务分配情况,确认是否存在多线程任务共享非线程安全资源的情况。 - 模拟复现问题:在测试环境模拟多线程高并发场景,通过增加消息量、调整线程数触发问题。结合栈信息中提到的topic(
wrw-technischerplatz-mietobjekt-aggregat-oeffentlich-1)和partition(2),针对性发送测试消息,稳定复现后定位具体触发条件。
内容的提问来源于stack exchange,提问作者Andras Hatvani
相关产品推荐
相关产品推荐

