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

Python批量下载千级以上Zip文件的速度优化咨询

优化大规模Zip文件下载与解压的方案(基于Spark集群)

核心优化方向

针对你当前的场景(3000+文件已耗时1-2小时,十万级URL待处理),核心优化点集中在并行化利用Spark集群能力、消除不必要的磁盘IO和网络请求优化三个方面:

1. 分布式并行处理

当前代码用list(map)串行执行,完全浪费了Spark集群的分布式计算能力。通过将URL列表转为Spark RDD/DataFrame,可以把任务分散到集群所有节点同时执行,大幅缩短总耗时。

2. 合并下载与解压流程

当前先下载所有Zip到本地再解压,产生了40GB的磁盘写入+读取开销。合并流程后可以直接在内存中处理Zip文件,跳过中间存储,直接将CSV写入目标路径,这能显著减少IO耗时。

3. 网络请求优化

  • 启用HTTP会话保持,减少TCP握手次数
  • 添加超时与重试机制,避免因网络波动导致任务挂起或失败

基于Spark的优化代码

import requests
import zipfile
from io import BytesIO
from pyspark.sql import SparkSession
from requests.adapters import HTTPAdapter
from urllib3.util.retry import Retry

def download_and_extract(url, export_path):
    # 配置会话重试机制,处理网络波动
    session = requests.Session()
    retry = Retry(total=3, backoff_factor=1, status_forcelist=[429, 500, 502, 503, 504])
    adapter = HTTPAdapter(max_retries=retry)
    session.mount("http://", adapter)
    session.mount("https://", adapter)
    
    try:
        # 带超时的请求,避免任务挂起
        response = session.get(url, timeout=30)
        response.raise_for_status()  # 捕获HTTP错误状态码
        
        # 内存中直接解析Zip,跳过本地磁盘存储
        with zipfile.ZipFile(BytesIO(response.content)) as zip_ref:
            # 只处理CSV文件,过滤其他无关文件
            for file_name in zip_ref.namelist():
                if file_name.lower().endswith('.csv'):
                    # 写入共享存储路径(如HDFS/S3,确保集群所有节点可访问)
                    with open(f"{export_path}/{file_name}", "wb") as csv_file:
                        csv_file.write(zip_ref.read(file_name))
    except Exception as e:
        # 记录失败URL,后续可批量重试
        print(f"处理URL {url} 失败: {str(e)}")

if __name__ == "__main__":
    # 初始化Spark会话,根据集群配置调整参数
    spark = SparkSession.builder \
        .appName("LargeScaleZipProcessing") \
        .getOrCreate()
    
    # 替换为你的URL列表,可从文件/HDFS加载
    links = ["url1", "url2", ...]
    # 必须使用集群共享存储路径,不能用本地磁盘
    export_path = "hdfs://your-cluster-storage/csv-output"
    
    # 控制并行度:numSlices根据集群节点数和资源调整,比如100-500
    spark.sparkContext.parallelize(links, numSlices=200) \
        .foreach(lambda url: download_and_extract(url, export_path))
    
    spark.stop()

关键注意事项

  • 共享存储:Spark集群节点无法访问本地磁盘,必须使用HDFS、S3或NFS等共享存储作为输出路径,否则只有执行任务的单个节点能看到生成的CSV文件。
  • 并行度调整:numSlices参数要根据集群的CPU和网络带宽调整,避免并发过高导致目标服务器限流,或集群资源耗尽。
  • 错误处理:代码中添加了重试机制和异常捕获,可扩展为将失败URL写入日志文件,后续进行批量重试。
  • 函数合并:已将下载和解压合并为单一函数,彻底消除了中间Zip文件的磁盘IO开销,这是提升性能的核心举措之一。

额外优化建议

  • 若目标服务器支持,可使用wget或curl命令行工具替代requests,Spark中通过subprocess调用,能获得更底层的网络优化。
  • 对于十万级URL,建议分批次提交任务,比如每批处理5000个URL,避免一次性占用过多集群资源。
  • 监控集群的CPU、内存和网络使用率,动态调整并行度,平衡处理速度和集群负载。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 17:27:14