【发布时间】: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