You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

异步Kafka实现中produce()方法仅被调用一次是什么原因?

问题原因

  1. confluent_kafka的SerializingProducer.produce()仅会将消息写入本地发送缓冲区,不会主动触发网络发送逻辑,必须调用producer.poll(0)或producer.flush()处理内部网络IO、错误回调等事件,否则消息会持续堆积在本地,甚至内部报错后静默终止生产逻辑,导致你只能看到第一次put!打印。
  2. 你在同步的主代码块直接使用await asyncio.gather(...)不符合Python语法要求,同步上下文不能直接调用await,会导致asyncio事件循环调度异常,生产协程执行一次后就不会再被调度。
  3. 你没有为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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.28 17:15:06