Dask如何将数据处理流水线中间结果持久化到磁盘而非内存
Dask大体积中间结果磁盘持久化方案
Dask本身是惰性执行机制,只要中间结果没有被固化,每次触发下游action操作都会沿着计算图回溯重跑所有上游依赖步骤。.persist()是把结果常驻集群/本地内存的方案,数据量超过内存容量时确实不适用,你可以用下面几种磁盘持久化方案解决重复计算问题:
手动落盘列式存储文件(生产环境首选)
这是稳定性最高、跨场景兼容性最好的方案,优先选Parquet作为存储格式,压缩比高、支持谓词下推和列裁剪,读写性能远高于CSV、JSON这类文本格式。
操作逻辑很简单:把计算好的中间结果直接写入磁盘路径,后续所有依赖该结果的步骤直接从落盘路径读取数据,完全不会触发之前的转换逻辑重算。
参考代码:# 先调整中间结果的分区大小,单分区控制在100MB-256MB区间,避免小文件过多 intermediate_df = intermediate_df.repartition(partition_size="128MB") # 写入磁盘,选zstd压缩平衡压缩率和读写速度 intermediate_df.to_parquet( path="/data/pipeline/cached_sensor_intermediate/", engine="pyarrow", compression="zstd", overwrite=True, write_index=False ) # 后续所有步骤直接读取落盘的缓存数据即可 cached_df = dask.dataframe.read_parquet("/data/pipeline/cached_sensor_intermediate/") # 基于cached_df开展后续所有转换操作如果你的流水线跑在分布式集群上,把路径换成集群可访问的共享存储地址(HDFS、对象存储、分布式文件系统都可以),不要写单节点本地磁盘,避免其他计算节点访问不到。
基于Dask缓存接口自动做磁盘缓存
如果你不想手动管理缓存文件的读写、清理逻辑,可以用Dask的Cache模块配置磁盘后端,代码改动量极小,配置完成后正常调用.persist()就会自动把结果缓存到磁盘,不会全量占满内存。
参考代码:import os import tempfile from dask.cache import Cache # 选择SSD路径作为缓存目录,能大幅提升缓存读写速度 disk_cache_path = os.path.join(tempfile.gettempdir(), "dask_disk_cache") # 配置缓存可用磁盘空间上限,比如分配150GB空间给缓存 cache = Cache(cache_dir=disk_cache_path, available_bytes=150 * 1024**3) # 注册缓存到全局Dask调度器 cache.register() # 这里调用persist不会把数据全放内存,会自动按策略缓存到磁盘 cached_df = intermediate_df.persist()这个方案适合快速迭代的开发场景,缺点是缓存是Dask内部格式,没法直接被其他数据处理工具复用,缓存失效也需要手动清理目录。
注意:不要尝试用
.compute()把中间结果转成Pandas DataFrame来规避重算,数据量超过内存时会直接触发OOM,和.persist()的内存问题本质上没有区别。
内容的提问来源于stack exchange,提问作者Michael
相关产品推荐
相关产品推荐

