【问题标题】:Processing kafka messages taking long time处理kafka消息需要很长时间
【发布时间】:2019-11-25 01:58:04
【问题描述】:

我有一个 Python 进程(或者更确切地说,是一组在消费者组中并行运行的进程),它根据来自某个主题的 Kafka 消息的输入来处理数据。通常每条消息都会被快速处理,但有时,根据消息的内容,可能需要很长时间(几分钟)。在这种情况下,Kafka 代理会断开客户端与组的连接并启动重新平衡。我可以将session_timeout_ms 设置为一个非常大的值,但它会多出 10 分钟,这意味着如果客户端死亡,集群将在 10 分钟内无法正确重新平衡。这似乎是个坏主意。此外,大多数消息(大约 98% 的消息)速度很快,因此只为 1-2% 的消息支付这样的惩罚似乎很浪费。 OTOH,大消息频繁到足以导致大量重新平衡并消耗大量性能(因为当组重新平衡时,什么都没有完成,然后“死”客户端再次重新加入并导致另一个重新平衡)。

那么,我想知道,有没有其他方法可以处理需要很长时间才能处理的消息?有没有办法手动启动心跳来告诉代理“没关系,我还活着,我只是在处理消息”?我认为 Python 客户端(我使用 kafka-python 1.4.7)应该为我做到这一点,但它似乎没有发生。此外,API 似乎根本没有单独的“心跳”功能。据我了解,调用poll() 实际上会给我下一条消息——而我什至还没有处理完当前的消息,而且还会弄乱 Kafka 消费者的迭代器 API,这在 Python 中使用起来非常方便。

如果我没记错的话,Kafka 集群是 Confluent 2.3 版。

【问题讨论】:

    标签: apache-kafka kafka-consumer-api kafka-python


    【解决方案1】:

    在 Kafka 中,0.10.1+ 的 Kafka 轮询和会话心跳是相互解耦的。 可以得到解释here

    ma​​x.poll.interval.ms 在超时之前允许消费者实例完成处理的时间是指如果处理时间超过 max.poll.interval.ms 时间,消费者组将假定其从消费者组中删除并调用重新平衡。

    增加这将增加预期轮询之间的间隔,从而使消费者有更多时间处理从 poll(long) 返回的一批记录。 但同时,它也会延迟组重新平衡,因为消费者只会在轮询调用中加入重新平衡。

    session.timeout.ms 是用于识别消费者是否还活着并在定义的时间间隔 (heartbeat.interval.ms) 发送心跳的超时时间。一般来说,经验法则是 heartbeat.interval.ms 应该是会话超时的 1/3,因此在网络故障的情况下,消费者最多可以在会话超时之前错过 3 次心跳。

    1. session.timeout.ms:较低的值有利于更快地检测到故障。

    2. max.poll.interval.ms:较大的值会降低由于处理时间增加而导致失败的风险,但会增加重新平衡时间。

    注意:Consumer Group 消耗的大量分区和主题也会影响整体重新平衡时间

    如果您真的想摆脱重新平衡,则可以使用另一种方法手动分配每个消费者实例上的分区,使用分区分配。在这种情况下,每个消费者实例将使用自己分配的分区独立运行。但在这种情况下,您将无法利用重新平衡功能自动分配分区。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-08-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多