【发布时间】:2018-02-23 00:31:17
【问题描述】:
看最新(v0.10)Kafka Consumer documentation:
“消费者的位置给出了下一条记录的偏移量。它将比消费者在该分区中看到的最高偏移量大一个。它会自动前进每个消费者接收数据调用poll(long)并接收消息的时间。"
有没有办法查询服务器端分区可用的最大偏移量,而不检索所有消息?
我试图实现的逻辑如下:
- 每秒查询主题中待处理消息的数量 (A)
- 如果 A > 阈值,则唤醒将继续检索所有消息并处理它们的处理器
- 否则什么都不做(睡眠 1)
动机是我需要做一些批处理,但我希望处理器只有在有足够的数据时才唤醒(并且我不想两次检索所有数据)。
【问题讨论】:
标签: apache-kafka kafka-consumer-api