【问题标题】:Limit the number of messages to consume in KafkaConsumer of kafka-python library在 kafka-python 库的 KafkaConsumer 中限制要消费的消息数量
【发布时间】:2021-07-09 08:37:23
【问题描述】:

示例代码:

consumer = KafkaConsumer(config["kafka"]["input"],
                         bootstrap_servers=config["kafka"]["brokers"].split(','),
                         value_deserializer=lambda m: json.loads(m.decode('ascii')),
                         enable_auto_commit=config["kafka"]["auto_commit"],
                         auto_commit_interval_ms=config["kafka"]["commit_interval"],
                         group_id=config["kafka"]["group"],
                         consumer_timeout_ms=config["kafka"]["timeout"]
                        )

尝试了 max_poll_records、fetch_max_bytes 甚至 consumer.poll() 方法都没有奏效。

【问题讨论】:

    标签: kafka-consumer-api kafka-python


    【解决方案1】:

    Confluent-Kafka python 库允许限制消息的数量。

    用途:

    consumer.consume(num_messages=config["kafka"]["number_of_messages_to_consumer"], timeout=config["kafka"]["timeout"]) 限制消息的消费。

    Kafka-Python 库不支持这个。

    【讨论】:

      猜你喜欢
      • 2018-03-06
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2022-06-15
      • 1970-01-01
      • 2014-02-13
      • 1970-01-01
      相关资源
      最近更新 更多