【问题标题】:How to add a failure callback for kafka-python kafka.KafkaProducer#send()?如何为 kafka-python kafka.KafkaProducer#send() 添加失败回调?
【发布时间】:2017-12-20 03:33:39
【问题描述】:

如果生成的记录失败,我想设置一个回调。最初,我只想记录失败的记录。

Confluent Kafka python 库提供了一种添加回调的机制:

produce(topic[, value][, key][, partition][, on_delivery][, timestamp])
...
    on_delivery(err,msg) (func) – Delivery report callback to call (from poll() or flush()) on successful or failed delivery

如何使用 kafka-python kafka.KafkaProducer#send() 实现类似的行为,而不必使用已弃用的 SimpleClient 使用 kafka.SimpleClient#send_produce_request()

【问题讨论】:

    标签: kafka-python


    【解决方案1】:

    虽然没有记录,但这相对简单。每当您发送消息时,您都会立即收到Future 回复。您可以将回调/errback 附加到 Future

    F = producer.send(topic=topic, value=message, key=key)
    F.add_callback(callback, message=message, **kwargs_to_pass_to_callback_method)
    F.add_errback(erback, message=message, **kwargs_to_pass_to_errback_method)
    

    相关源码在这里: https://github.com/dpkp/kafka-python/blob/1937ce59b4706b44091bb536a9b810ae657c3225/kafka/future.py#L48-L64

    我们真的应该记录这个,我提交了https://github.com/dpkp/kafka-python/issues/1256 来跟踪它。

    【讨论】:

      猜你喜欢
      • 2019-04-20
      • 1970-01-01
      • 2015-11-12
      • 2021-08-16
      • 1970-01-01
      • 2014-10-14
      • 1970-01-01
      • 1970-01-01
      • 2012-12-15
      相关资源
      最近更新 更多