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

如何修改Python代码实现列表的异步迭代处理?

问题描述

我正在写第一段异步Python代码,把场景简化成了两个脚本:main.py和s3_accessor.py。main脚本生成一个包含数百个值的列表,需要把每个值传入build_inventory函数处理,但同步处理速度太慢了。

原代码

main.py

import s3_accessor # 另一个.py脚本
import asyncio

async def main():
     some_list = [a, b, c, d, e, f, g]
     
     for val in some_list:
          inventory_df = s3_accessor.build_inventory(val)
          s3_accessor.write_df_to_s3('some_key', inventory_df)

if __name__ == "__main__":
     loop = asyncio.get_event_loop()
     loop.run_until_complete(main()) 

s3_accessor.py

async def build_inventory(val):
     ...

疑问:
怎么改代码才能实现列表的异步迭代处理?我试过用await asyncio.wait()把整个for循环包起来,但好像不对。


解决方案

要实现异步并发处理列表里的每个值,核心是把每个值的处理逻辑包装成独立的异步任务,让它们同时执行,而非逐个等待前一个完成。同时要注意几个关键细节:

  1. 异步函数必须用await调用,否则只会返回一个协程对象,不会实际执行。
  2. 如果write_df_to_s3是同步IO操作,会阻塞整个事件循环,要么改成异步实现,要么用asyncio.to_thread把它放到线程池执行,避免阻塞异步逻辑。
  3. 批量处理异步任务推荐用asyncio.gather,它能一次性启动所有任务并等待全部完成,是处理并发任务最常用的方式。

修改后的代码

先调整s3_accessor.py

如果原来的write_df_to_s3是同步函数,先把它改成异步兼容的版本:

import asyncio

# 假设原同步写入逻辑如下
def _write_df_to_s3_sync(key, df):
    # 原有的同步S3写入代码
    ...

async def build_inventory(val):
     # 原有的异步逻辑,比如异步拉取数据生成DataFrame
     ...

async def write_df_to_s3(key, df):
    # 用to_thread把同步操作放到线程池,不阻塞事件循环
    await asyncio.to_thread(_write_df_to_s3_sync, key, df)

修改后的main.py

import s3_accessor
import asyncio

async def process_single_val(val):
    # 把单个值的处理逻辑封装成独立函数
    inventory_df = await s3_accessor.build_inventory(val)
    # 注意这里也要await异步的写入方法,避免异步任务未完成就结束
    await s3_accessor.write_df_to_s3(f'some_key_{val}', inventory_df) # 建议给每个值用不同key,避免数据覆盖

async def main():
     some_list = [a, b, c, d, e, f, g]
     
     # 为每个值创建一个异步任务
     tasks = [process_single_val(val) for val in some_list]
     # 并发执行所有任务,等待全部完成
     await asyncio.gather(*tasks)

if __name__ == "__main__":
     # Python 3.7+可用更简洁的asyncio.run替代原有loop写法
     asyncio.run(main())

额外优化:控制并发量

如果并发数过高导致S3限流报错,可以用asyncio.Semaphore限制同时运行的任务数:

async def process_single_val(val, semaphore):
    async with semaphore:
        inventory_df = await s3_accessor.build_inventory(val)
        await s3_accessor.write_df_to_s3(f'some_key_{val}', inventory_df)

async def main():
     some_list = [a, b, c, d, e, f, g]
     # 限制最多10个任务同时执行,可根据实际情况调整
     semaphore = asyncio.Semaphore(10)
     tasks = [process_single_val(val, semaphore) for val in some_list]
     await asyncio.gather(*tasks)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 00:45:24