【发布时间】:2017-06-03 11:55:22
【问题描述】:
我有一个 REST 服务,我们称之为 MDD,它有一个 kafka 消费者。当我第一次启动 rest 服务时,另一个服务告诉 MDD 的消费者订阅一个特定的主题,一切似乎都很好。
然后服务告诉 MDD 的消费者订阅另一个主题。我现在这样做的方式是通过 consumer.assign() 方法。基本上,如果引入了一个未分配消费者的新主题,我将这个新主题分配给消费者。因此,一个消费者现在被分配到 2 个不同的主题。
这个消费者轮询消息并将它们存入 HDFS。
现在我注意到的是,当订阅第二个主题时,有时我会收到关于未能附加到 HDFS 中的文件的错误,当我查看日志时,它试图附加一些应该附加的数据直到后来才被附加。 例如,kafka 的数据按 A、B、C 的顺序出现。当 MDD 完成将 A 附加到 HDFS 后,它会尝试附加 C(而不是 B),同时也尝试附加 B。另请注意,此时没有来自第一个主题的数据,只有来自第二个主题的数据正在流入。所以目前,在任何给定时间只有一个 kafka 主题有数据流入。
有人知道会发生什么吗?当我将一个消费者分配给多个主题时,是否会产生一些线程问题?因为当消费者被分配给一个主题时,一切似乎都很好,但是一旦它被分配给一个以上的主题,我就无法附加到 HDFS 中的文件,因为其他一些作者已经拥有了租约。这个错误不会经常发生,只是非常随机。
还有一个推荐的修复方法是每次创建新主题时,创建一个新的 kafka 消费者吗?
【问题讨论】:
标签: java multithreading hadoop apache-kafka