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