【发布时间】:2021-12-12 15:57:48
【问题描述】:
我知道,Kafka 旨在像处理无限的事件流一样处理,并且获取剩余消息计数不是内置功能。但我必须以某种方式监控我的消费者进程的运行情况以及我是否为他们提供了足够的资源。
总体场景是基本的 Kafka 使用,不同服务器上的几个生产者插入一个主题,group_a 中的消费者读取消息,执行一些 AI 工作并插入另一个主题以进行进一步处理。
传入消息的速率绝不是恒定的或可预测的,因此我需要检查我的消费者是否落后(假设group_a 在我的输入中有超过 1000 条未读消息主题)。
考虑到我可以完全控制 Kafka 设置以及消费者和生产者代码,我有哪些选择?
- 我认为在没有任何繁重处理和计数消息的情况下阅读整个主题既不干净也不高效 (?)
- 如果只有一对生产者/消费者,我可以计算生产和消费消息的数量并计算剩余的消息,但我正在运行多服务器设置,这不太可行。是否可以将 Kafka 本身用作共享数据存储并记录所有生产和消费的消息?
- 为了避免 XY 问题,计算未读消息数量的正确方法是了解我的消费者是否获得了“足够”的资源?
【问题讨论】:
标签: python apache-kafka