配置双输入主题的函数路由时KafkaBinderMetrics异常问题
多Kafka主题绑定函数路由时出现ConcurrentModificationException异常
我参照Stack Overflow的方案配置了一个接收两个Kafka主题输入的函数路由,简化后的配置如下:
spring: cloud: function: routing-expression: "headers['eventType']" stream: bindings: functionRouter-in-0: destination: order.requests,order.events binder: kafka content-type: application/json group: ${spring.application.name}-app consumer: concurrency: 12 autoStartup: true
应用启动时抛出如下异常:
2023-12-19 11:45:25.505 DEBUG [order-fulfilment-svc,,,] 21186 --- [pool-3-thread-1] o.s.c.s.binder.kafka.KafkaBinderMetrics : Cannot generate metric for topic: order.requests java.util.ConcurrentModificationException: KafkaConsumer is not safe for multi-threaded access. currentThread(name: pool-3-thread-1, id: 65) otherThread(id: 66) at org.apache.kafka.clients.consumer.KafkaConsumer.acquire(KafkaConsumer.java:2551) at org.apache.kafka.clients.consumer.KafkaConsumer.acquireAndEnsureOpen(KafkaConsumer.java:2532) at org.apache.kafka.clients.consumer.KafkaConsumer.partitionsFor(KafkaConsumer.java:1993) at org.apache.kafka.clients.consumer.KafkaConsumer.partitionsFor(KafkaConsumer.java:1969) at org.springframework.cloud.stream.binder.kafka.KafkaBinderMetrics.findTotalTopicGroupLag(KafkaBinderMetrics.java:204) at org.springframework.cloud.stream.binder.kafka.KafkaBinderMetrics.computeUnconsumedMessages(KafkaBinderMetrics.java:189) at org.springframework.cloud.stream.binder.kafka.KafkaBinderMetrics.lambda$bindTo$0(KafkaBinderMetrics.java:154) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) at java.base/java.util.concurrent.FutureTask.runAndReset$$$capture(FutureTask.java:305) at java.base/java.util.concurrent.FutureTask.runAndReset(FutureTask.java) at java.base/java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:305) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829)
我尝试将并发数改为1,但问题仍未解决;移除第二个主题后,异常消失。当前使用版本为SCS 3.2.9和spring-kafka 2.9.13。
请问这是否是一个Bug?
内容的提问来源于stack exchange,提问作者Pete
相关产品推荐
相关产品推荐

