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

无法消费异步消息:Kafka程序卡在迭代1挂起求助

Kafka异步消息消费挂起问题解决方案

问题场景

Kafka、Zookeeper运行正常,以下是producer和consumer代码,运行后程序卡在迭代1并挂起:

producer.py代码

async def publish():
    producer = AIOKafkaProducer(bootstrap_servers='localhost:9092',
    enable_idempotence=True)  
    await producer.start()

    consumer = AIOKafkaConsumer(
    topicAKG,
    bootstrap_servers='localhost:9092',group_id='test',
    max_poll_interval_ms=60000,
    max_poll_records=50)
    await consumer.start()

    try:
        for i in range(1, 6):
            await producer.send_and_wait(topic, value='from producer'.encode())
            print(f"Iteration: {i}")
            async for message in consumer:
                print("Received ========== ", message.value.decode())
                await consumer.commit()
    finally:
        await producer.stop()
        await consumer.stop()

consumer.py代码

import asyncio
from aiokafka import AIOKafkaConsumer, AIOKafkaProducer

topic = 'app'
topicAKG = 'back'

async def consume():
    consumer = AIOKafkaConsumer(topic, bootstrap_servers='localhost:9092',
    group_id="test",
    max_poll_interval_ms=60000,
    max_poll_records=50)
    await consumer.start()

    producer = AIOKafkaProducer(bootstrap_servers='localhost:9092',
                                enable_idempotence=True)
    await producer.start()

    try:
        async for message in consumer:
            print("Received",message.value.decode())
            await asyncio.sleep(2)  # Delay for 3 seconds
            await consumer.commit()  # Commit the offset to avoid re-consuming the same message
            await producer.send_and_wait(topicAKG, value='from consumer'.encode())
    finally:
        await consumer.stop()
        await producer.stop()

loop = asyncio.get_event_loop()
loop.run_until_complete(consume())

运行输出

  • 生产者输出:
    Iteration: 1
    Received ==========  from consumer
    
  • 消费者输出:
    Received from producer
    

问题原因

producer代码中,async for message in consumer是无限异步循环,会持续监听topicAKG的消息且不会主动退出。第一次迭代时,producer发送消息后收到consumer返回的一条消息,随后循环会一直等待下一条消息,但此时consumer已完成第一次处理,没有新消息发送,导致整个程序阻塞在该循环中,无法进入下一次for i in range(1,6)迭代。

解决方案

方案一:单次获取单条消息

修改producer的消费逻辑,用consumer.getone()替代async for循环,只获取一条消息处理后继续</think_never_used_51bce0c785ca2f68081bfa7d91973934>主循环:

# 补充缺失的topic定义
topic = 'app'
topicAKG = 'back'

async def publish():
    producer = AIOKafkaProducer(bootstrap_servers='localhost:9092',
    enable_idempotence=True)  
    await producer.start()

    consumer = AIOKafkaConsumer(
    topicAKG,
    bootstrap_servers='localhost:9092',group_id='test',
    max_poll_interval_ms=60000,
    max_poll_records=50)
    await consumer.start()

    try:
        for i in range(1, 6):
            await producer.send_and_wait(topic, value='from producer'.encode())
            print(f"Iteration: {i}")
            # 获取单条消息后处理,退出当前消费逻辑
            message = await consumer.getone()
            print("Received ========== ", message.value.decode())
            await consumer.commit()
    finally:
        await producer.stop()
        await consumer.stop()

方案二:后台异步任务消费

如果需要持续监听topicAKG的消息,同时执行发送循环,可将消费逻辑放到后台异步任务,避免阻塞主循环:

import asyncio
from aiokafka import AIOKafkaProducer, AIOKafkaConsumer

topic = 'app'
topicAKG = 'back'

async def consume_background(consumer):
    # 后台持续消费消息
    async for message in consumer:
        print("Received ========== ", message.value.decode())
        await consumer.commit()

async def publish():
    producer = AIOKafkaProducer(bootstrap_servers='localhost:9092',
    enable_idempotence=True)  
    await producer.start()

    consumer = AIOKafkaConsumer(
    topicAKG,
    bootstrap_servers='localhost:9092',group_id='test',
    max_poll_interval_ms=60000,
    max_poll_records=50)
    await consumer.start()

    # 创建后台消费任务
    consume_task = asyncio.create_task(consume_background(consumer))

    try:
        for i in range(1, 6):
            await producer.send_and_wait(topic, value='from producer'.encode())
            print(f"Iteration: {i}")
            await asyncio.sleep(0.1)  # 给后台任务留处理时间
    finally:
        # 取消并等待后台任务结束
        consume_task.cancel()
        try:
            await consume_task
        except asyncio.CancelledError:
            pass
        await producer.stop()
        await consumer.stop()

注意事项

  • 确保producer代码中补充topic和topicAKG的定义(原代码缺失)
  • 方案二中的sleep是为了让后台消费任务有足够时间处理消息,避免主循环过快发送导致消息堆积
  • 根据业务需求选择方案:若每次发送只需接收一条响应,用方案一;若需持续监听响应同时发送,用方案二

内容的提问来源于stack exchange,提问作者Tanjin Alam

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 10:47:47