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

如何用Python的asyncio分组任务,实现串行与并发结合执行?

解决方案:组内串行、组间并发控制数据库任务

方案一:手动分组实现组内串行、组间并发

直接修改原函数,通过二次分组控制并发数:

from asyncio import gather
from typing import List, Any, datetime

async def pre_processing_list_sic_codes(
        self,
        start_date: datetime,
        final_date: datetime,
        sic_codes: List[str]
        ) -> Any:
        batch_length = 80
        batches = [sic_codes[i : i + batch_length] for i in range(0, len(sic_codes), batch_length)]
        
        # 控制同时运行的组数量(即实际并发的数据库任务数)
        concurrent_group_size = 3
        # 将批次划分为若干并发组
        task_groups = [batches[i:i+concurrent_group_size] for i in range(0, len(batches), concurrent_group_size)]
        
        # 定义组内串行执行的逻辑
        async def execute_group(group: List[List[str]]):
            group_results = []
            for batch in group:
                # 组内逐个await,确保串行执行
                result = await self.processing(
                    sic_codes=batch,
                    start_date=start_date,
                    final_date=final_date,
                )
                group_results.append(result)
            return group_results
        
        # 并发执行所有组
        group_tasks = [execute_group(group) for group in task_groups]
        all_group_results = await gather(*group_tasks)
        
        # 可选:将二维结果展平为一维
        flattened_results = [item for sublist in all_group_results for item in sublist]
        return flattened_results

关键说明

  • concurrent_group_size:设置为3或4,对应同时运行的数据库任务数,根据你的数据库连接限制调整
  • 每组内的任务通过await逐个执行,避免同一组内同时占用多个连接
  • 不同组的执行函数通过gather并发运行,保证效率的同时控制总并发数

方案二:用信号量(Semaphore)自动控制并发数

如果不想手动分组,使用asyncio.Semaphore可以更灵活地限制并发任务数:

from asyncio import gather, Semaphore
from typing import List, Any, datetime

async def pre_processing_list_sic_codes(
        self,
        start_date: datetime,
        final_date: datetime,
        sic_codes: List[str]
        ) -> Any:
        batch_length = 80
        batches = [sic_codes[i : i + batch_length] for i in range(0, len(sic_codes), batch_length)]
        
        # 限制同时运行的processing任务数为3
        semaphore = Semaphore(3)
        
        # 包装processing函数,添加信号量限制
        async def limited_processing(batch):
            async with semaphore:
                return await self.processing(
                    sic_codes=batch,
                    start_date=start_date,
                    final_date=final_date,
                )
        
        # 创建所有受限任务并并发执行
        tasks = [limited_processing(batch) for batch in batches]
        all_results = await gather(*tasks)
        return all_results

关键说明

  • Semaphore(3):最多允许3个processing任务同时执行,自动控制并发连接数
  • async with semaphore:确保每次只有指定数量的任务进入数据库操作,无需手动分组
  • 这种方式代码更简洁,适合需要动态调整并发数的场景

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 07:23:35