【发布时间】: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