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
相关产品推荐
相关产品推荐

