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

请求基于Python asyncio实现数据获取与处理的异步并发示例脚本

使用Python asyncio实现流水线式迭代循环

要实现你描述的流水线式循环(处理当前数据时并行获取下一轮数据,且数据处理串行执行),可以通过提前启动下一轮的数据获取任务,同时保证数据处理步骤按顺序执行来实现。以下是符合需求的asyncio示例脚本:

import asyncio
import random

async def data_fetching(iteration):
    """模拟数据获取,比如从API/数据库拉取数据"""
    print(f"开始获取第 {iteration} 轮数据")
    # 模拟耗时操作,比如网络请求
    await asyncio.sleep(random.uniform(0.5, 1.5))
    fetched_data = f"第 {iteration} 轮获取的数据"
    print(f"完成获取第 {iteration} 轮数据")
    return fetched_data

async def data_processing(data, iteration):
    """模拟数据处理,比如解析、计算等"""
    print(f"开始处理第 {iteration} 轮数据: {data}")
    # 模拟耗时操作,比如数据计算
    await asyncio.sleep(random.uniform(1.0, 2.0))
    print(f"完成处理第 {iteration} 轮数据")
    return f"处理后的{data}"

async def main(total_iterations=5):
    # 预启动第一轮的数据获取任务
    next_fetch_task = asyncio.create_task(data_fetching(1))
    
    for iteration in range(1, total_iterations + 1):
        # 等待当前轮次的数据获取完成
        current_data = await next_fetch_task
        
        # 如果不是最后一轮,启动下一轮的数据获取任务(和当前处理并行)
        if iteration < total_iterations:
            next_fetch_task = asyncio.create_task(data_fetching(iteration + 1))
        
        # 处理当前轮次的数据(串行执行,不会和其他处理步骤并发)
        await data_processing(current_data, iteration)

if __name__ == "__main__":
    asyncio.run(main())

代码逻辑说明

  • 并行获取与串行处理:在处理第i轮数据前,已经启动了第i+1轮的获取任务,这样数据处理和下一轮的数据获取会并行执行,提升整体效率。
  • 内存控制:同一时间只有一个数据处理任务在执行,且仅保留当前待处理数据和一个正在执行的获取任务,不会堆积多个未处理的数据,符合内存限制要求。
  • 任务衔接:通过asyncio.create_task提前启动下一轮获取,用await保证当前处理完成后再进入下一轮,确保流程顺序正确。

内容的提问来源于stack exchange,提问作者Arash Vahabpour

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 18:24:58