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

配置双输入主题的函数路由时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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 18:52:41