【问题标题】:KafkaConsumer: consume all available messages once, and exitKafkaConsumer:消费所有可用消息一次,然后退出
【发布时间】: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


    【解决方案1】:

    实现这一点的唯一方法似乎是使用一些额外的逻辑来自我管理偏移量。伪代码如下:

    consumer.assign(all partitions)
    consumer.seekToBeginning(all partitions)
    // record the current ending offsets and poll until we get there
    endOffsets = consumer.endOffsets(all partitions)
    
    do
      result = consumer.poll(NONTRIVIAL_TIMEOUT)
      // onResult(result)
    while given any partition p, consumer.position(p) < endOffsets[p]
    

    以及在 Kotlin 中的实现:

    val topicPartitions = consumer.partitionsFor(topic).map { TopicPartition(it.topic(), it.partition()) }
    consumer.assign(topicPartitions)
    consumer.seekToBeginning(consumer.assignment())
    
    val endOffsets = consumer.endOffsets(consumer.assignment())
    fun pendingMessages() = endOffsets.any { consumer.position(it.key) < it.value }
    
    do {
      records = consumer.poll(Duration.ofMillis(1000))
      onResult(records)
    } while(pendingMessages())
    

    现在可以将轮询持续时间设置为合理的值(例如 1 秒),而不必担心丢失消息,因为循环会一直持续到消费者到达循环开始时确定的结束偏移量。

    还有另一种正确处理的极端情况:如果结束偏移量已更改,但当前偏移量和结束偏移量之间实际上没有消息,则轮询将阻塞并超时。因此,重要的是不要将超时设置得太低(否则消费者将在检索可用的消息之前超时)并且它也不能设置得太高(否则消费者将花费太长时间检索可用的消息时超时)。如果这些消息被删除,或者如果主题被删除并重新创建,则可能会发生后一种情况。

    【讨论】:

      【解决方案2】:

      如果没有人同时生产,你也可以使用endOffsets获取最后一条消息的位置,直到那个时候才消费。

      所以,在伪代码中:

      long currentOffset = -1
      long endOffset = consumer.endOffset(partition)
      while (currentOffset < endOffset) {
        records = consumer.poll(NONTRIVIAL_TIMEOUT) // discussed in your answer
        currentOffset = records.offsets().max()
      }
      

      这样我们可以避免最终的非零挂断,因为我们总是可以确定有东西可以接收。

      如果您的消费者的位置等于结束偏移量,您可能需要添加保护措施(因为在那里您不会收到任何消息)。

      此外,您可能希望将 max.poll.records 设置为 1,这样如果有人并行生产,您就不会使用位于 结束偏移之后的消息。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2018-03-06
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2017-12-21
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多