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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 11:05:04