本文详解如何使用 Dask 与 fsspec 正确读取并拼接多个 ZIP 文件中同名 CSV(如所有 a.csv),解决通配符跨 ZIP 失效问题,并提供可复用的自动化方案。
本文详解如何使用 dask 与 fsspec 正确读取并拼接多个 zip 文件中同名 csv(如所有 `a.csv`),解决通配符跨 zip 失效问题,并提供可复用的自动化方案。
在使用 Dask 处理大规模结构化数据时,一个常见但易被忽视的挑战是:fsspec 的 glob 机制不支持跨多个 ZIP 文件进行通配匹配。例如,你期望 zip://a.csv::foo*.zip 能自动遍历 foo1.zip、foo2.zip、foo3.zip 并分别提取其中的 a.csv 后拼接——但实际仅返回首个 ZIP 中的文件。这是因为 fsspec 将每个 ZIP 视为独立的“子文件系统”,其通配逻辑(*)仅作用于 ZIP 内部路径,而非 ZIP 文件名本身;而 Dask 的 dd.read_csv() 在单次调用中也只绑定一个底层文件系统实例,无法动态切换 ZIP 源。
因此,必须采用显式构造多路径 + 手动拼接的策略。以下是推荐的生产级实现:
✅ 正确做法:显式生成 ZIP 内路径列表 + 并行读取 + 拼接
import dask.dataframe as dd
import glob
import os
# 步骤 1:获取所有匹配的 ZIP 文件路径(支持通配)
zip_files = sorted(glob.glob("foo*.zip")) # 如 ['foo1.zip', 'foo2.zip', 'foo3.zip']
# 步骤 2:为指定 CSV 文件名(如 'a.csv')构建完整的 fsspec URL 列表
csv_name = "a.csv"
file_urls = [f"zip://{csv_name}::{os.path.abspath(zf)}" for zf in zip_files]
# 步骤 3:并行读取每个 ZIP 中的 CSV(Dask 自动处理延迟计算)
dfs = [
dd.read_csv(
url,
delimiter=";",
header=0,
index_col=False,
assume_missing=True, # 建议启用,兼容列缺失场景
blocksize=None # 对小 CSV 可设 None;大文件建议指定(如 "64MB")
)
for url in file_urls
]
# 步骤 4:沿行方向拼接(保持分区逻辑,避免立即 compute)
ddf = dd.concat(dfs, interleave_partitions=False, ignore_index=True)
# ✅ 最终计算(仅在此刻触发实际 I/O 和计算)
result = ddf.compute()
print(f"成功合并 {len(zip_files)} 个 ZIP 中的 '{csv_name}',总计 {len(result)} 行")
⚠️ 关键注意事项
- 路径必须绝对化:os.path.abspath(zf) 避免相对路径导致 fsspec 解析失败(尤其在非工作目录运行时);
- 不要提前 .compute():对每个 dd.read_csv() 单独调用 .compute() 会丧失 Dask 的延迟执行优势,大幅降低性能;
- interleave_partitions=False 是关键参数:确保 dd.concat 不尝试重排分区,仅做简单垂直拼接,避免冗余 shuffle;
- 内存与性能权衡:若 ZIP 数量极大(如数百个),建议分批处理或使用 dask.delayed 控制并发粒度;
- 错误容错增强(可选):可在列表推导中加入 try/except 包裹 dd.read_csv,跳过损坏 ZIP 或缺失 CSV 的文件。
? 扩展:批量处理多个 CSV 名称(a.csv, b.csv, c.csv)
for csv_name in ["a.csv", "b.csv", "c.csv"]:
file_urls = [f"zip://{csv_name}::{os.path.abspath(zf)}" for zf in zip_files]
dfs = [dd.read_csv(url, delimiter=";", header=0, index_col=False) for url in file_urls]
ddf = dd.concat(dfs, ignore_index=True)
# 保存中间结果或直接参与后续 merge
ddf.to_parquet(f"merged_{csv_name.replace('.csv', '')}.parquet", write_index=False)
该方案完全规避了 fsspec 的 glob 限制,充分利用 Dask 的惰性计算与并行 I/O 能力,是处理跨 ZIP 同名文件合并任务的稳健范式。










