Kafka 作为长时间运行任务的消息队列

Ger*_*phy 8 apache-kafka kafka-consumer-api

我想知道我的设置是否有遗漏以促进长期运行的工作。

出于我的目的,可以进行At most once消息传递,这意味着不需要考虑提交偏移量(或者至少可以在接收到消息时提交每个消息偏移量)。

为了实现竞争消费者模式,我有以下几点:

  • 一个话题
  • 同一组中的 X 个消费者
  • 主题中的 P 个分区(其中 P >= X 始终)

我的问题是我的消息可能需要大约 15 分钟(但这可能会波动高达 50%)才能处理。为了避免消费者的分区分配被撤销,我增加了 的值max.poll.interval.ms以反映这一点。然而,这会带来一些负面影响:

  • 如果某些消息超过此时间长度,那么在最坏的情况下,处理此消息的使用者将不得不等待max.poll.interval.ms重新平衡的值
  • 如果我需要根据负载扩展和增加消费者的数量,那么任何新消费者也可能必须等待max.poll.interval.ms重新平衡的值才能处理任何新消息

就目前而言,我认为我可以按以下方式进行:

  • 设置max.poll.interval.ms为一个小值并接受每个消费者处理每条消息都会超时并经历撤销分配并等待少量时间进行重新平衡的过程

但是我不喜欢这个,并且正在考虑为我的消息队列寻找替代技术,因为我没有看到任何明显的解决方法。诚然,我是 Kafka 的新手,并且只是一种直觉,以上是不可取的。我过去曾在这些场景中使用过 RabbitMQ,但是目前我们的架构中需要 Kafka 用于其他目的,如果 Kafka 可以实现这一点,那么不必引入另一种技术就好了。

我很感激任何人都可以就此主题提供的任何建议。

sen*_*iwu 5

使用 Kafka 作为作业队列来调度长时间运行的进程并不是一个好主意,因为 Kafka 不是最严格意义上的队列,并且故障处理和重试的语义是有限的。尽管您可以通过使用某些配置来重新平衡或超时来实现折衷,但它很可能仍然是脆弱的设计。简单的答案是 Kafka 不是为这些用例设计的。

的想法max.poll.interval.ms是为了防止活锁情况(请参阅参考资料),但在您的情况下,消费者将向 Kafka 代理发送误报并触发重新平衡,因为无法区分活锁和合法的长进程。

你应该考虑一下你提到的负面后果与生活之间的权衡。引入一项新技术,可帮助您以更好的方式对作业队列进行建模。对于更复杂的用例,请查看slack 是如何做到的


Ger*_*phy 3

我们解决所遇到问题的方法是按照评论中的建议进行的。我们决定将消息处理与消费者轮询分离。

每个工作线程/消费者上有 2 个线程,一个用于执行实际处理,另一个用于定期向 Kafka 打电话。

我们还做了一些工作来尝试减少消息的处理时间。然而,有些消息仍然需要花费几分钟的时间。这对我们来说已经工作了一段时间,没有任何问题。

感谢@Donal 在评论中提出的建议