【问题标题】:Kafka Consumer - topic(s) with higher priorityKafka Consumer - 具有更高优先级的主题
【发布时间】:2017-07-17 08:43:44
【问题描述】:

我正在使用 Kafka Consumer 读取多个主题,我需要其中一个具有更高的优先级。处理需要很长时间,并且(低优先级)主题中总是有很多消息,但我需要尽快处理来自其他主题的消息。

Does Kafka support priority for topic or message? 的问题类似,但这个问题使用的是旧 API。

在新的 API (0.10.1.1) 中,有一些方法

KafkaConsumer::pause(Collection)
KafkaConsumer::resume(Collection)

但我不清楚,如何有效地检测高优先级主题中有新消息,需要暂停其他主题的消费。

有什么想法/例子吗?

【问题讨论】:

  • 您可以检查您正在监视的分区的 endOffsets 是否大于这些分区的最后提交的偏移量。这究竟是如何工作的将是特定于实现的,但这会让您在轮询之前知道是否有更多消息要使用
  • 请看这个,它可能是你要找的:stackoverflow.com/a/66013251/4602706

标签: apache-kafka kafka-consumer-api


【解决方案1】:

最后我解决了这个问题,正如 dawsaw 建议的那样 - 在处理循环中,我存储了我从中读取的所有主题/分区:

  • 开始偏移
  • endOffsets
  • committed - 我不能使用位置,因为我订阅主题,而不是分区。

每当(endOffset - commited) > 0 处理任何优先主题时,我都会调用consumer.pause() 处理非优先主题,并在(endOffset - commited) == 0 之后再次调用所有优先主题。

【讨论】:

  • 您能分享您解决问题的策略吗?假设我们有(总共 10 Gbs)低优先级消息和一些高优先级消息。我们有多个消费者和多个生产者。即使我们暂停消费者,我们也需要暂停所有其他主题的生产者,以使您的想法实现。对?请问您有这方面的经验吗,因为在 100 个服务和 10 个主题的生态系统中这似乎几乎是不可能的? - 是的,我已经阅读了你关于此事的其他相关问题。谢谢
  • 不 - 没有必要暂停任何生产者 - 想法是您有单个消费者订阅了多个主题(其中一些主题是高优先级的,其他主题是普通优先级的)。在轮询新消息之前,您需要检查高优先级主题的滞后。如果这些滞后中的任何一个不为零,则意味着您需要暂停对正常优先级主题的订阅,而不是“消耗”消费者的时间。处理完所有来自高优先级主题的消息后,您可以再次恢复正常优先级的消息。
  • 谢谢。我不能完全违抗。但对于较大的系统来说,它闻起来很糟糕。一旦大坝的大门为大量数据打开,我将不得不不时检查我是否在这个低优先级队列中浪费资源。我为什么要?对。反正。再次感谢
  • @miran 你知道 Kafka 的 Python 客户端是否有类似的实现?
  • 我相信 python 客户端提供与 java 相同的 API - 所以你绝对可以实现它...
【解决方案2】:

我猜你可以混合使用 position() 和 commit() 方法。 position() 方法获取将要获取的下一条记录的偏移量,并且 commit() 方法获取给定分区的最后提交的偏移量(如文档中所述)。 在轮询较低优先级之前,您可以检查 position() 和 commit() 以获得较高优先级。如果 position() 比 commit() 高,你可以 pause() 低优先级, poll() 高优先级(),然后恢复低优先级。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-08-19
    • 2021-12-09
    • 2021-02-25
    • 2017-02-20
    相关资源
    最近更新 更多