无法消费异步消息: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
相关产品推荐
相关产品推荐

