请求基于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
相关产品推荐
相关产品推荐

