如何使用Polars高效对大型数据集执行Upsert(更新+插入)操作
我完全理解你遇到的困境——10GB的数据集在16GB内存里直接用update操作肯定会爆内存,毕竟update是全量在内存中处理数据,加上临时计算的开销,内存根本扛不住。下面我给你几个针对大数据场景优化的内存友好解决方案:
方案一:LazyFrame + 磁盘分块处理(通用兜底方案)
核心思路是绝对避免把全量大数据锁在内存里,用Polars的Lazy延迟执行模式,配合磁盘存储实现分块处理。我们可以用join+concat手动实现upsert逻辑,再开启streaming让Polars分批加载数据:
步骤1:先把大数据集持久化到Parquet
如果你的df_old是刚生成的,先把它写入压缩的Parquet文件(Parquet是列式存储,压缩比高,完美适配Polars的Lazy读取):
# 用snappy压缩写入,兼顾速度和磁盘占用 df_old.write_parquet("large_dataset.parquet", compression="snappy")
步骤2:用Lazy模式实现Upsert逻辑
用scan_parquet读取大文件(不会全量加载到内存),通过anti join筛选出旧数据中未被更新的行,最后和新数据合并:
import polars as pl # Lazy模式读取大文件(按需加载,不占满内存) lf_old = pl.scan_parquet("large_dataset.parquet") # 新数据也转成LazyFrame,保持统一的延迟执行逻辑 lf_new = df_new.lazy() # 实现Upsert:旧数据未被更新的行 + 所有新数据行(更新+新增) df_upsert = ( pl.concat( [ # 筛选旧数据中不在新数据里的行(保留未被更新的原有数据) lf_old.join(lf_new, on=["group", "id"], how="anti"), # 所有新数据行(包括替换旧数据的更新行,以及完全新增的行) lf_new ], how="vertical" ) .sort(["group", "id"]) .collect(streaming=True) # 开启streaming,让Polars分块处理数据,内存占用可控 ) # 验证结果 print(f"Upsert后数据集大小:{round(df_upsert.estimated_size()/10**9, 3)} GB") print(df_upsert.head())
这个方案的内存占用会被Polars自动控制在3-5GB左右(远低于16GB),因为streaming会把大文件拆成小批次处理,处理完一批就释放对应内存,不会触发OOM。
方案二:分区存储+增量更新(进阶高效方案)
如果你的更新操作经常按group这类列筛选,那按更新键分区存储能进一步减少需要处理的数据量——比如按group把大文件拆成多个小Parquet文件,更新时只需要处理涉及到的分区,其他分区直接复用:
步骤1:写入分区Parquet文件
df_old.write_parquet( "partitioned_large_dataset", # 生成一个文件夹,每个group对应一个子文件夹 partition_by="group", # 按group分区,更新时只处理目标group的分区 compression="snappy" )
步骤2:只处理需要更新的分区
# 先拿到新数据中涉及的group列表 update_groups = df_new.select("group").unique().to_series().to_list() # 1. 读取需要更新的旧分区(比如A、D),筛选出未被更新的行 lf_updated_old = pl.scan_parquet( "partitioned_large_dataset", filters=[("group", "in", update_groups)] # 只读取目标group的分区 ) lf_updated_part = lf_updated_old.join(lf_new, on=["group", "id"], how="anti") # 2. 读取不需要更新的旧分区(除A、D、XYZ之外的所有group),直接复用 lf_unchanged = pl.scan_parquet( "partitioned_large_dataset", filters=[("group", "not in", update_groups)] ) # 3. 合并所有部分:未修改的旧分区 + 筛选后的待更新分区 + 所有新数据 df_upsert = ( pl.concat([lf_unchanged, lf_updated_part, lf_new.lazy()], how="vertical") .sort(["group", "id"]) .collect(streaming=True) )
这个方案的优势是只处理极小一部分数据(比如原10GB数据,只处理A、D两个group的几百MB数据),内存占用极低,速度也快很多。
为什么原方法会失败?
你用的df_old.update(...)是eager模式的全量内存操作:它需要把df_old和df_new同时放在内存里,还要创建临时结构来对齐、合并数据——10GB的df_old加上临时计算的开销,16GB内存完全不够,直接触发内存溢出导致内核崩溃。
额外优化小技巧
- 全程用LazyFrame:避免在内存中保留大的eager DataFrame,处理完就写入磁盘,用
scan_*系列函数读取。 - 选对压缩算法:Parquet用
snappy(速度快)或zstd(压缩比高),既省磁盘又加快读取速度。 - 调整streaming块大小:如果还是有内存压力,可以通过
pl.Config.set_streaming_chunk_size(1024*1024*1024)设置每个块的大小(比如1GB),让Polars每次处理更小的批次。
备注:内容来源于stack exchange,提问作者Olibarer

