
本文详解如何使用 Dask + fsspec 正确读取并拼接多个 ZIP 文件中同名 CSV(如所有 a.csv)——因 fsspec 不支持跨 ZIP 的通配符 glob,需显式构造路径并手动 concat。
本文详解如何使用 dask + fsspec 正确读取并拼接多个 zip 文件中同名 csv(如所有 `a.csv`)——因 fsspec 不支持跨 zip 的通配符 glob,需显式构造路径并手动 concat。
在使用 Dask 处理超大规模数据时,一个常见需求是:从多个 ZIP 归档文件中提取同名 CSV 文件(如 a.csv)并纵向拼接(concat),最终形成单一大型 Dask DataFrame 以支持内存外计算。然而,直接使用通配符路径(如 'zip://a.csv::foo*.zip' 或 'foo*.zip::a.csv')往往失败——它仅返回首个匹配 ZIP 中的文件,而非全部。
根本原因在于 fsspec 的设计限制:其 glob 操作仅作用于「最内层文件系统」(即 ZIP 包内部),不支持对多个外部 ZIP 文件进行通配符枚举;同时,dd.read_csv() 单次调用始终绑定到一个确定的 filesystem 实例,无法自动遍历多个 ZIP 源。
✅ 正确做法是:显式生成每个 ZIP 中目标 CSV 的完整 zip:// 路径,分别读取为独立的 Dask DataFrame,再通过 dd.concat() 合并。以下是完整、可复用的实现方案:
import dask.dataframe as dd
import glob
import os
# 步骤 1:获取所有匹配的 ZIP 文件路径(如 foo1.zip, foo2.zip...)
zip_paths = sorted(glob.glob("foo*.zip")) # 确保顺序稳定,便于调试
# 步骤 2:为指定 CSV 文件名(如 "a.csv")构建所有 zip:// 路径
target_csv = "a.csv"
file_urls = [f"zip://{target_csv}::{os.path.abspath(zip_path)}" for zip_path in zip_paths]
# 步骤 3:逐个读取并构建 Dask DataFrame 列表
dfs = [
dd.read_csv(
url,
delimiter=";", # 根据实际分隔符调整
header=0,
index_col=False,
blocksize=None # 对于小 CSV 可设 None;大文件建议指定(如 "64MB")以优化分块
)
for url in file_urls
]
# 步骤 4:沿行方向拼接(ignore_index=True 避免重复索引)
ddf = dd.concat(dfs, axis=0, ignore_index=True, interleave_partitions=False)
# 步骤 5:触发计算(或后续保存/分析)
result = ddf.compute()
print(f"Concatenated {len(zip_paths)} '{target_csv}' files → shape: {result.shape}")
? 关键注意事项:
- ✅ 路径必须绝对化:os.path.abspath() 确保 zip:// URL 在分布式环境中(如 Dask Cluster)仍有效,避免相对路径解析失败;
- ✅ 显式指定 blocksize:若 ZIP 内 CSV 较大,建议设置 blocksize="64MB" 让 Dask 自动分块读取,提升并行效率;
- ⚠️ 避免 interleave_partitions=True:该参数在多源 concat 时可能引发列顺序错乱(尤其当各 CSV 列名/顺序不完全一致时),默认 False 更安全;
- ? 批量处理多个 CSV 名称? 将上述逻辑封装为函数,循环处理 ["a.csv", "b.csv", "c.csv"] 即可:
def concat_csv_from_zips(csv_name, zip_pattern="foo*.zip", **read_csv_kwargs):
zip_paths = sorted(glob.glob(zip_pattern))
urls = [f"zip://{csv_name}::{os.path.abspath(p)}" for p in zip_paths]
dfs = [dd.read_csv(url, **read_csv_kwargs) for url in urls]
return dd.concat(dfs, axis=0, ignore_index=True)
# 使用示例
dfa = concat_csv_from_zips("a.csv", delimiter=";")
dfb = concat_csv_from_zips("b.csv", delimiter=";")
? 进阶提示: 若 ZIP 数量极多(数百+),可考虑用 dask.delayed 替代列表推导式,实现更细粒度的任务调度与错误隔离;对于生产环境,建议添加 try/except 包裹单个 read_csv 调用,并记录失败 ZIP,提升鲁棒性。
综上,Dask 本身不提供“跨 ZIP 通配符读取”的语法糖,但通过显式路径管理与 dd.concat 组合,即可优雅、高效、可控地完成多源同名 CSV 的合并任务——这正是 Dask “显式优于隐式” 设计哲学的典型体现。










