You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用Polars高效对大型数据集执行Upsert(更新+插入)操作

如何使用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内存完全不够,直接触发内存溢出导致内核崩溃。

额外优化小技巧

  1. 全程用LazyFrame:避免在内存中保留大的eager DataFrame,处理完就写入磁盘,用scan_*系列函数读取。
  2. 选对压缩算法:Parquet用snappy(速度快)或zstd(压缩比高),既省磁盘又加快读取速度。
  3. 调整streaming块大小:如果还是有内存压力,可以通过pl.Config.set_streaming_chunk_size(1024*1024*1024)设置每个块的大小(比如1GB),让Polars每次处理更小的批次。

备注:内容来源于stack exchange,提问作者Olibarer

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.14 14:18:00