【发布时间】:2017-12-29 20:59:25
【问题描述】:
我有一个消费者项目,它使用来自 Kafka 主题的数据。该流中 90% 的数据可以实时处理,但对于特定记录 (~10%),我需要延迟处理。
我是否应该在同一个 JVM 中有两个独立的消费者并在一个消费者中消费 90% 的记录并忽略 10% 并让其他消费者处理它或将 10% 的消息推送到另一个主题并延迟处理其他话题?
如果我可以有一个消费者和两个检查点机制,一个用于 90%,另一个延迟 10%,但 Kafka 客户端似乎不支持这个用例,那就太好了。这将帮助我避免任何不必要的反序列化和网络 IO。
【问题讨论】:
-
“特定记录”是什么意思?他们有什么特殊的领域?另外,“延迟过程”是什么意思?正如您所说,如果它们不代表相同的事件,最合乎逻辑的做法是在另一个主题中生成这些数据并相应地使用它。另一种方法是让消费者简单地读取数据并将它们传输给某个工作人员。然后你会有两个工人:一个是实时的,一个是延迟消息的。早点阅读它们没有问题。我可能没有清楚了解您的需求,请不要犹豫,提供一些详细信息
-
@Treziac 我们正在处理事件数据并每分钟将它们以微批次的形式存储。如果事件时间戳来自前几天(约 10% 的记录),我们需要延迟处理并为历史数据提供更大的批次。我认为使用同一个消费者并将它们转移给另一个工作人员可能不是一个好主意,因为数据需要在内存中超过 15 分钟(延迟),我们不能在这个过程中提交检查点。此外,我无法控制数据的生成方式,因此我要么需要将消息推送到另一个主题,要么为每个用例设置两个消费者。
标签: java apache-kafka kafka-consumer-api