如何用Python asyncio.Semaphore循环启动批量协程并分批执行?
优化asyncio批量协程分批执行方案
原代码存在的潜在问题
- 块范围计算错误:当
blocks_needed不是10的整数倍时,原代码会多处理块。比如blocks_needed=15,外层循环会生成latest_block_number和latest_block_number-10两个批次,总共处理20个块,超出了需要的15个。 - 任务列表复用风险:反复
clear()任务列表虽然可行,但不如每次批次创建新列表清晰,减少潜在的引用问题。 - 语义不够直观:两层嵌套循环的结构可读性较差,难以一眼看出批次划分逻辑。
优化后的实现方案
async def start(self): sem = asyncio.Semaphore(10) latest_block_number = await self.w3.eth.block_number # 生成所有需要处理的块号列表(从最新块往前取blocks_needed个) target_blocks = list(range( latest_block_number, latest_block_number - self.blocks_needed, -1 )) # 按每10个块为一批拆分 batch_size = 10 batches = [ target_blocks[i:i+batch_size] for i in range(0, len(target_blocks), batch_size) ] for batch in batches: # 为当前批次创建所有任务 tasks = [ asyncio.create_task( self.get_transactions_of_block(current_block_number=block_num, sem=sem) ) for block_num in batch ] # 等待当前批次所有任务完成 await asyncio.gather(*tasks) # 非最后一批才等待60秒 if batch != batches[-1]: await asyncio.sleep(60)
关键优化点说明
- 精确控制处理范围:先一次性生成所有需要处理的块号,再按批次拆分,确保不会多处理或漏处理块。
- 简化批次逻辑:用列表推导式拆分批次,结构更直观,可读性更强。
- 避免任务列表复用:每个批次创建独立的任务列表,消除潜在的状态残留问题。
- 优化等待逻辑:最后一批任务完成后无需再等待60秒,避免不必要的延迟。
- 保持信号量作用:Semaphore(10)确保同时运行的协程不超过10个,和原逻辑一致,但批次拆分更合理。
额外建议
如果get_transactions_of_block内部没有正确使用信号量,需要确保在函数内通过上下文管理器获取和释放:
async def get_transactions_of_block(self, current_block_number, sem): async with sem: # 实际的块交易获取逻辑 block = await self.w3.eth.get_block(current_block_number, full_transactions=True) # 后续交易处理逻辑...
内容的提问来源于stack exchange,提问作者svat_di
相关产品推荐
相关产品推荐

