【问题标题】:Kafka as a message queue for long running tasksKafka 作为长时间运行任务的消息队列
【发布时间】:2020-05-22 16:57:05
【问题描述】:

我想知道我的设置是否缺少一些东西来促进长期运行的工作。

就我的目的而言,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 可以实现这一点,那么不必引入其他技术会很好。

感谢任何人就这个问题提供的任何建议。

【问题讨论】:

  • 我看到的问题是 15 分钟的消息处理时间。不能以某种方式对该消息进行切片,并行化,以便 kafka 在您的处理器/消费者之间分配负载吗? - 不是需要 15 分钟的大消息,而是 1000 条小消息(也可能需要 15 分钟,但避免了您的投票和重新分区问题)
  • 我们遇到的问题是,我们的一些处理通过大量增加并行性,端到端过程实际上变慢了,因为每条消息都需要相同的数据才能处理每条消息,而不管数量多少处理每条消息所包含的工作量。我们也有调用外部系统同步数据的场景,根据同步数据的大小,消息处理时间可能会有所不同。当然,在可能的情况下,我们的目标是让每条消息的工作量非常小,并且并行度很高,但这并不总是可行的。
  • 消费和处理如何解耦?就像一堆提供给处理器池的消费者线程之类的东西。这些消费者负责控制 Kafka 民意调查的时间,而其他消费者则负责。例如,如果您需要同步数据,则将负责该进程的线程设为与 kafka 消费者不同的线程,将其解耦;这样,即使在另一个线程上上传了 30 分钟,您的消费者也会继续缓冲消息和/或提供其他处理器。
  • @GerardMurphy 你能解决这个问题吗?如果有,你能说说你是怎么处理的吗?
  • 解决这个问题的方法是将消费者轮询与处理消息的线程分离,并手动提交偏移量。当线程完成对消息的处理后,就可以提交偏移量。完成此操作后,消费者必须在最后提交的偏移量之后重新寻找下一条消息。几年前我为一个客户实现了这个,但我无权访问该代码。我目前正在寻找一个支持此功能的库,希望不必再次编写它。如果你找到了,请分享。

标签: apache-kafka kafka-consumer-api


【解决方案1】:

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

max.poll.interval.ms 的想法是防止活锁情况 (see),但在您的情况下,消费者将向 Kafka 代理发送误报并触发重新平衡,因为无法区分活锁和一个合法的漫长过程。

您应该考虑在承受您提到的负面后果与承受负面后果之间进行权衡。引入了一项新技术,可以帮助您以更好的方式对作业队列进行建模。对于更复杂的用例,请查看how slack is doing it

【讨论】:

  • 我不明白为什么使用 Kafka 是一个坏主意,因为它支持具有可配置保留期的不可变消息主题,使其成为作业队列的强大选项。消费者组提供了一种很好的并行消息处理方式,易于扩展,如果配置正确,一条消息一次只能由同一消费者组中的一个消费者处理。消息处理时间较长的问题在于,它增加了消费者轮询管理的复杂性,必须将其与消息处理分离。
  • 如果 Kafka 已经在您的技术堆栈中,那没关系。否则,这只是 金锤反模式 的另一种情况,其中引入了一种不是设计为队列并且需要 3 个运行时实例(用于生产就绪设置)的技术,只是为了实现作业调度程序跨度>
【解决方案2】:

我们解决问题的方法与 cmets 中的建议一致。 我们决定将消息处理与消费者轮询分离。

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

我们还尝试减少消息的处理时间。 然而,有些消息仍然需要时间,可以以分钟为单位来衡量。 这已经为我们工作了一段时间,没有任何问题。

感谢 cmets @Donal 的建议

【讨论】:

  • 所以基本上在收到消息后,您将其发送到线程池进行处理,然后进行手动提交。我说的对吗?
  • 好吧,在我们的例子中,我们没有线程池,我们也只在消息被处理后提交。我们在 Kubernetes 中运行我们的工作负载,并根据需要扩大/缩小消费者。每个消费者容器有 2 个线程。我们将处理卸载到不同的线程,然后在主线程上以特定的时间间隔轮询 kafka。这实际上告诉 kafka 消费者仍在响应,因此消费者不会被踢出消费者组并将其分区分配给不同的消费者。
  • 我刚刚使用 reactor-kafka 完成了一个大数据解决方案摄取。 (Azure 事件中心的 Kafka)以前,我按照您在此处所做的那样实施了解决方案。这次我必须找到已经开发的解决方案。 reactor-kafka 在单独的线程上管理轮询,并实现背压,仅在有需求时才将消息推送给消费者。这是美丽的观察工作。这确实需要使用 Project Reactor,但我所看到的好处使其值得。总的来说,代码少了很多。没有未来,没有执行者,但通过反应堆调度程序进行大规模扩展。
猜你喜欢
  • 2015-11-29
  • 1970-01-01
  • 2015-06-20
  • 1970-01-01
  • 2019-01-09
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多