Python Asyncio:确保保存任务串行执行,且不阻塞数据采集
解决Asyncio中串行执行后台保存任务的问题
你的核心需求是让采集任务和保存任务并行,但同一时间只能有一个保存任务在运行。原代码的问题在于每次采集完成就直接创建新的保存任务,导致多个保存操作并发执行。以下是两种可行的解决方法:
方法一:维护保存任务链
通过跟踪当前正在执行的保存任务,让新的保存任务等待上一个完成后再启动,采集任务不受任何阻塞:
import asyncio import random async def save_data(): print("I'm saving a batch") await asyncio.sleep(2) print("I'm done saving") async def collect_data(): current_save_task = None while True: print("I'm collecting data") await asyncio.sleep(random.randint(1, 5)) async def wrapped_save(): nonlocal current_save_task # 等待上一个保存任务完成 if current_save_task is not None: await current_save_task await save_data() current_save_task = asyncio.create_task(wrapped_save()) asyncio.run(collect_data())
每次采集完成后,我们创建一个包装后的保存任务,它会先等待上一个保存任务结束,再执行自身。采集任务继续循环采集,完全不被保存操作拖慢。
方法二:使用异步队列
用asyncio.Queue实现生产者-消费者模型,采集任务负责往队列里发保存请求,单独的后台任务负责串行处理这些请求:
import asyncio import random async def save_data(): print("I'm saving a batch") await asyncio.sleep(2) print("I'm done saving") async def save_worker(queue): # 后台消费者,串行处理所有保存请求 while True: await queue.get() await save_data() queue.task_done() async def collect_data(): save_queue = asyncio.Queue() # 启动后台保存进程 asyncio.create_task(save_worker(save_queue)) while True: print("I'm collecting data") await asyncio.sleep(random.randint(1, 5)) # 发送保存请求到队列 await save_queue.put(None) asyncio.run(collect_data())
这种方式更适合后续扩展(比如需要传递采集到的数据给保存任务),队列会自动帮你管理保存请求的顺序,确保逐个执行。
内容的提问来源于stack exchange,提问作者BENJAMIN GILBERT
相关产品推荐
相关产品推荐

