如何利用asyncio实现produceMessages与processMessages异步并行运行?
问题分析与解决方案
你的核心问题在于:processMessages的初始化方法里直接调用asyncio.run(self.producer.run()),这个调用会永久阻塞(因为生产者是无限循环),导致后续的print_logs完全无法执行。要让生产者和消费者并发运行,必须将两者都放入同一个asyncio事件循环中,让任务异步调度执行。
修改后的代码
first_file.py(仅规范类名,功能不变)
import asyncio import random class ProduceMessages: def __init__(self, timeout=10): self.timeout = timeout self.message_logs = [] async def run(self): while True: self.message_logs.append(random.uniform(0, 1)) await asyncio.sleep(self.timeout)
second_file.py(核心修改)
import first_file import asyncio class ProcessMessages: def __init__(self): self.producer = first_file.ProduceMessages(timeout=2) # 调小超时方便测试 async def _process_logs(self): last_processed_idx = 0 while True: # 仅处理新增的消息(更符合实际消费逻辑) new_messages = self.producer.message_logs[last_processed_idx:] if new_messages: print(f"新增消息: {new_messages}") last_processed_idx = len(self.producer.message_logs) await asyncio.sleep(1) # 用asyncio.sleep替代time.sleep,不阻塞事件循环 async def run(self): # 将生产者和消费者包装为异步任务 producer_task = asyncio.create_task(self.producer.run()) consumer_task = asyncio.create_task(self._process_logs()) # 等待两个无限循环任务持续运行 await asyncio.gather(producer_task, consumer_task) if __name__ == "__main__": processor = ProcessMessages() asyncio.run(processor.run())
关键修改点说明
- 避免初始化时阻塞:移除
__init__中的asyncio.run调用,改为在专门的异步run方法中启动任务。 - 并发任务调度:使用
asyncio.create_task将生产者和消费者的异步函数包装为独立任务,让asyncio事件循环自动调度两者交替执行。 - 非阻塞休眠:用
asyncio.sleep替代time.sleep,确保休眠时事件循环可以切换到其他任务(比如生产者)执行。 - 优化消费逻辑:新增消息追踪索引,仅处理每次循环中的新消息,而非打印整个日志列表。
内容的提问来源于stack exchange,提问作者medihde
相关产品推荐
相关产品推荐

