【问题标题】:How can I fix ConcurrentModificationException errors in Kafka? (0.9.0.1)如何修复 Kafka 中的 ConcurrentModificationException 错误? (0.9.0.1)
【发布时间】:2016-04-02 06:37:37
【问题描述】:

我知道 Kafka 声明 KafkaConsumer 不是线程安全的。

所以我这样做了:(Scala)

val m = Map(new TopicPartition(msg.topic(), msg.partition()) -> new OffsetAndMetadata(msg.offset()))
consumer.synchronized{ consumer.commitSync(m) }

我将我对消费者的访问权限放在一个同步块中,但我仍然在使用 consumer.commitSync(m) 时遇到 ConcurrentModificationException 错误。

为什么,我该怎么办?

我使用的是 Akka 流,所以肯定会有线程的奥秘,但同步块不应该解决这个问题吗?

【问题讨论】:

  • 我猜你有多个 consumer 实例,所以当你在不同的对象上调用 synchronized 时同步不起作用......你需要一个全局共享对象(即每个消费者线程用来获取锁的虚拟Object)。

标签: multithreading scala apache-kafka


【解决方案1】:

在文档中找到的一种方法是为 KafkaConsumer 创建一个单独的线程,并通过某种并发队列与外部工作进行通信。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2014-09-25
    • 2016-07-18
    • 1970-01-01
    • 2022-11-17
    • 2016-10-12
    • 2020-09-07
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多