【问题标题】:Spring Cloud Stream: Proper handling for StreamListener when it takes a long time to process the messageSpring Cloud Stream:处理消息需要很长时间时对StreamListener的正确处理
【发布时间】:2019-02-12 11:49:50
【问题描述】:

当 StreamListener 需要很长时间(比max.poll.interval.ms)处理消息时,该特定消费者被占用,其他新消息将被分配到其他分区。时间大于max.poll.interval.ms后,再平衡发生,同样的情况也会发生在另一个消费者身上。因此,此消息将在所有分区中循环并继续占用资源。

但是,这种情况并不经常发生,只有几条消息不知何故需要这么长时间来处理,而且是无法控制的。

我们可以提交偏移量并在几次重新平衡后将其扔给 DLQ 吗?如果是,我们该怎么做?如果不是,这种情况的正确处理是什么?

【问题讨论】:

  • 为什么不直接增加max.poll.interval.ms
  • 因为此消息可能需要很长时间才能处理,但这不是常态。我不想为了覆盖这种信息而牺牲性能。我想忽略它,然后继续。

标签: spring-boot apache-kafka spring-cloud-stream spring-kafka


【解决方案1】:

增加max.poll.interval.ms 不会对性能产生影响(除非检测到真的死亡的消费者需要更长的时间)。

每次处理此“不良”记录时进行重新平衡对性能的损害更大。

但是,您可以使用自定义 SeekToCurrentErrorHandler together with a recoverer such as the DeadLetterPublishingRecoverer 做您想做的事情。您还需要一个重新平衡侦听器来计算重新平衡和一些机制来跨实例共享来自错误处理程序的状态(标准只将状态保存在内存中)。

我认为相当复杂。

【讨论】:

    猜你喜欢
    • 2019-06-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-04-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-04-15
    相关资源
    最近更新 更多