【发布时间】:2018-09-26 23:59:20
【问题描述】:
我想使用 Kafka 2.0.0 创建一个 KafkaConsumer,它会一次性使用所有可用消息,然后立即退出。这与标准控制台使用者实用程序略有不同,因为该实用程序 waits for a specified timeout 用于新消息,并且仅在超时到期后退出。
使用KafkaConsumer,这个看似简单的任务似乎异常困难。我的直觉反应是以下伪代码:
consumer.assign(all partitions)
consumer.seekToBeginning(all partitions)
do
result = consumer.poll(Duration.ofMillis(0))
// onResult(result)
while result is not empty
但这不起作用,因为poll 总是返回一个空集合,即使该主题有很多消息。
对此进行研究,看起来一个原因可能是分配/订阅是considered lazy,并且在poll 循环完成之前不会分配分区(尽管我在文档中找不到对这个断言的任何支持) .但是,以下伪代码 也 在每次调用 poll 时返回一个空集合:
consumer.assign(all partitions)
consumer.seekToBeginning(all partitions)
// returns nothing
result = consumer.poll(Duration.ofMillis(0))
// returns nothing
result = consumer.poll(Duration.ofMillis(0))
// returns nothing
result = consumer.poll(Duration.ofMillis(0))
// deprecated poll also returns nothing
result = consumer.poll(0)
// returns nothing
result = consumer.poll(0)
// returns nothing
result = consumer.poll(0)
...
显然“懒惰”不是问题。
javadoc 声明:
如果有可用记录,此方法立即返回。
这似乎暗示上面的第一个伪代码应该可以工作。但是,它没有。
似乎唯一可行的是在poll 上指定一个非零超时,而不仅仅是任何非零值,例如1 不起作用。这表明poll 内部发生了一些不确定的行为,它假定poll 将始终在无限循环中执行,尽管有消息,它偶尔会返回一个空集合并不重要。代码seems to confirm this 与各种调用来检查超时是否已过期遍布poll 实现及其被调用者。
因此,对于幼稚的方法,显然需要更长的超时时间(理想情况下Long.MAX_VALUE 以避免更短轮询间隔的不确定行为),但不幸的是,这将导致消费者在最后一次轮询时阻塞,这在这种情况下是不需要的。使用幼稚的方法,我们现在可以在我们希望行为的确定性与我们在最后一次民意调查中无缘无故地等待多长时间之间进行权衡。我们如何避免这种情况?
【问题讨论】:
标签: apache-kafka