异步Kafka实现中produce()方法仅被调用一次是什么原因?
问题原因
confluent_kafka的SerializingProducer.produce()仅会将消息写入本地发送缓冲区,不会主动触发网络发送逻辑,必须调用producer.poll(0)或producer.flush()处理内部网络IO、错误回调等事件,否则消息会持续堆积在本地,甚至内部报错后静默终止生产逻辑,导致你只能看到第一次put!打印。- 你在同步的主代码块直接使用
await asyncio.gather(...)不符合Python语法要求,同步上下文不能直接调用await,会导致asyncio事件循环调度异常,生产协程执行一次后就不会再被调度。 - 你没有为
produce方法配置投递回调,无法感知序列化失败、Broker连接异常等错误,问题被静默掩盖无法定位。
修复方案
修改生产函数,添加poll逻辑和投递回调
async def produce(topic_name, serializer): p = SerializingProducer({ "bootstrap.servers": "PLAINTEXT://localhost:9092", "value.serializer": serializer }) def delivery_report(err, msg): if err is not None: print(f"消息发送失败: {err}") else: print(f"消息投递成功,topic: {msg.topic()}, offset: {msg.offset()}") while True: try: p.produce( topic=topic_name, value=CIS(), on_delivery=delivery_report ) # 处理内部事件,触发消息发送 p.poll(0) print("put!") except Exception as e: print(f"生产异常: {e}") await asyncio.sleep(1)
修复主函数的事件循环启动逻辑
将主代码中直接await的部分替换为loop.run_until_complete()启动:
if __name__ == "__main__": # 原有代码保持不变 topic_name = "my_topic" schema_str = json.dumps( { "type": "record", "name": "cis", "namespace": "interaction", "fields": [ {"name": "user_id", "type": "string"}, {"name": "question_id", "type": "int"}, {"name": "is_correct", "type": "boolean"} ] } ) def to_dict(obj, ctx): return asdict(obj) def to_obj(obj, ctx): return CIS( user_id=obj["user_id"], question_id=obj["question_id"], is_correct=obj["is_correct"], ) avro_serializer = AvroSerializer(registry_client, schema_str, to_dict) avro_deserializer = AvroDeserializer(registry_client, schema_str, to_obj) loop = asyncio.get_event_loop() t1 = loop.create_task(produce(topic_name, avro_serializer)) t2 = loop.create_task(consume(topic_name, avro_deserializer)) # 替换原有await语句 loop.run_until_complete(asyncio.gather(t1, t2))
可选优化消费逻辑
调整sleep位置,降低消费延迟,同时添加错误处理:
async def consume(topic_name, deserialzier): c = DeserializingConsumer( { 'bootstrap.servers': "PLAINTEXT://localhost:9092", 'value.deserializer': deserialzier, 'group.id': "123", 'auto.offset.reset': "latest" } ) c.subscribe([topic_name]) while True: message = c.poll(0.1) if message is None: print(message) await asyncio.sleep(1) continue elif message.error(): print(f"消费错误: {message.error()}") continue else: print(message.value())
修改完成后重新运行即可看到生产逻辑会每秒触发一次,消费逻辑也能正常拉取到消息。
内容的提问来源于stack exchange,提问作者user3595632
相关产品推荐
相关产品推荐

