如何基于asyncio+aiohttp并发编排7z分卷文件的下载与解压流程?
如何基于asyncio+aiohttp并发编排7z分卷文件的下载与解压流程?
看起来你已经找对了核心方向!这个“下载当前文件分卷的同时解压上一个已完成下载的文件”的编排思路非常贴合你的需求,能充分利用网络IO和CPU资源,比同步串行的方式效率提升不少,先给你点个赞👍
先说说你的方案里做得好的地方:
- 用TaskGroup管理并发任务:这是Python 3.11+推荐的并发任务管理方式,比手动创建Task更简洁,还能自动处理任务异常和等待,非常适合这种批量并发场景。
- 用
asyncio.to_thread包装同步解压函数:subprocess.run是阻塞式操作,直接放在协程里会卡住整个事件循环,而to_thread能把这个阻塞操作放到线程池执行,不会影响其他下载任务的运行,这个细节处理得很到位。 - 按完整文件分组分卷:把同属一个7z文件的分卷归为一组,是实现“下载N+解压N-1”逻辑的基础,这个分组逻辑是对的。
可以优化的几个细节:
1. 分卷分组逻辑更严谨
你原来的分组逻辑用了正则替换,建议用更稳妥的方式提取完整文件名,避免匹配出错,同时保证分组顺序正确:
import re import os def group_parts(self, files): file_groups = {} for item in files: part_name = item["part"] # 从分卷名提取基础文件名(比如从file_1.7z.001得到file_1) match = re.match(r"(.*)\.7z\.\d+$", part_name) if match: base_name = match.group(1) if base_name not in file_groups: file_groups[base_name] = [] file_groups[base_name].append(part_name) # 按文件名序号排序,保证file_1、file_2的顺序 return sorted(file_groups.values(), key=lambda x: int(re.search(r"file_(\d+)", x[0]).group(1)))
2. 封装下载逻辑,提升代码复用性
把下载一组分卷的逻辑封装成单独协程,让主逻辑更简洁:
async def download_file_parts(self, parts): async with asyncio.TaskGroup() as tg: for part in parts: # 这里假设你的download_from_gitlab需要的url是拼接出来的,可根据实际调整 url = f"{self.base_url}/{part}" tg.create_task(self.download_from_gitlab(url, part))
3. 完善解压的错误处理
你的extract_synchronous可以增加返回码判断,避免静默失败:
def extract_synchronous(self, program: str, args: list[str]): try: cmd = [program, *args] proc = subprocess.run(cmd, stderr=subprocess.STDOUT, stdout=subprocess.PIPE, text=True) if proc.returncode != 0: raise RuntimeError(f"解压失败!命令:{' '.join(cmd)},返回码:{proc.returncode},输出:{proc.stdout}") print(f"✅ 解压成功:{' '.join(cmd)}") except Exception as e: print(f"❌ 解压出错:{str(e)}") raise # 抛出异常让TaskGroup捕获,方便上层处理
4. 合理配置下载信号量
在类初始化时设置信号量大小,控制同时下载的分卷数量,避免被服务器限流:
def __init__(self, base_url, save_dir): self.base_url = base_url self.save_dir = save_dir self._session = aiohttp.ClientSession() self._semaphore = asyncio.Semaphore(5) # 同时最多下载5个分卷,可根据服务器情况调整 os.makedirs(save_dir, exist_ok=True)
5. 处理下载失败的残留文件
如果分卷下载失败,自动删除不完整的文件,避免影响后续解压:
async def download_from_gitlab(self, url: str, output_document: str): try: async with self._semaphore: async with self._session.get(url=url) as r: r.raise_for_status() # HTTP状态码非200时抛出异常 with open(output_document, "wb") as f: chunk_size = 8192*1024 async for data in r.content.iter_chunked(chunk_size): f.write(data) print(f"✅ 下载完成:{output_document}") except Exception as e: # 下载失败,删除残留文件 if os.path.exists(output_document): os.remove(output_document) print(f"❌ 下载失败:{output_document},错误:{str(e)}") raise e
优化后的主逻辑示例
async def main(self, files: list): grouped_parts = self.group_parts(files) for idx, parts in enumerate(grouped_parts): current_file = parts[0].split('.')[0] if idx == 0: print(f"\n===== 开始下载第{idx+1}组文件:{current_file} =====") try: await self.download_file_parts(parts) print(f"===== 第{idx+1}组文件下载完成 =====") except Exception as e: print(f"===== 第{idx+1}组文件下载失败,跳过后续操作 =====") continue else: prev_file = grouped_parts[idx-1][0].split('.')[0] print(f"\n===== 开始下载第{idx+1}组文件:{current_file},同时解压第{idx}组文件:{prev_file} =====") try: async with asyncio.TaskGroup() as tg: # 下载当前组 tg.create_task(self.download_file_parts(parts)) # 解压上一组(7z只需指定第一个分卷) extract_args = ['x', grouped_parts[idx-1][0], f"-o{self.save_dir}"] tg.create_task(asyncio.to_thread(self.extract_synchronous, "7z", extract_args)) print(f"===== 第{idx+1}组下载+第{idx}组解压完成 =====") except Exception as e: print(f"===== 第{idx+1}组下载/第{idx}组解压失败 =====") continue # 最后处理最后一组的解压 if grouped_parts: last_parts = grouped_parts[-1] print(f"\n===== 开始解压最后一组文件:{last_parts[0].split('.')[0]} =====") extract_args = ['x', last_parts[0], f"-o{self.save_dir}"] await asyncio.to_thread(self.extract_synchronous, "7z", extract_args) print(f"===== 所有文件处理完成 =====") await self._session.close()
总结
你的核心思路完全正确,只要把上述细节优化一下,就能实现高效的并行下载与解压流程,让网络下载(IO密集)和解压(CPU/IO密集)操作互相重叠,比串行方式节省大量时间。
内容来源于stack exchange
相关产品推荐
相关产品推荐

