为何Python-Polars数据集执行.unpivot()后体积大幅增大?
问题
我在Polars DataFrame中存储了一个约15万行×15万列的对称矩阵(含names索引列),占用约100GB内存。执行unpivot()操作后,内存占用飙升至600GB,最终数据集体积超过400GB,远大于原数据。示例代码如下:
# 示例数据 df = { "A": [1, 2, 3], "B": [2, 1, 2], "C": [3, 2, 1], "names": ["A", "B", "C"] } df = pl.DataFrame(df) # 执行unpivot df = df.unpivot(index = "names", variable_name = "names_2", value_name = "distance") # 输出结果 # shape: (9, 3) # ┌───────┬─────────┬──────────┐ # │ names ┆ names_2 ┆ distance │ # │ --- ┆ --- ┆ --- │ # │ str ┆ str ┆ i64 │ # ╞═══════╪═════════╪══════════╡ # │ A ┆ A ┆ 1 │ # │ B ┆ A ┆ 2 │ # │ C ┆ A ┆ 3 │ # │ A ┆ B ┆ 2 │ # │ B ┆ B ┆ 1 │ # │ C ┆ B ┆ 2 │ # │ A ┆ C ┆ 3 │ # │ B ┆ C ┆ 2 │ # │ C ┆ C ┆ 1 │ # └───────┴─────────┴──────────┘
原因分析
- 数据行数爆炸:原矩阵是15万行×15万列(含索引列),
unpivot()会将每一列的数值转换为一行,最终生成15万 × 15万 = 2.25e10行数据,是原数据行数的15万倍。 - 字符串重复存储:原数据中列名是元数据,仅存储一次;
unpivot()后列名被拆分为names_2列的每行值,同时names列也重复15万次。两个字符串列的重复存储是体积暴涨的核心原因——原100GB主要是数值数据,而新DataFrame中字符串存储的开销远大于数值。 - 对称数据冗余:原矩阵是对称的,
unpivot()后会同时保留(A,B)和(B,A)这类重复的对称数据,进一步加剧了数据量冗余。
优化方法
1. 仅保留非冗余的对称数据
利用矩阵的对称性,只保留上三角(含对角线)或下三角(含对角线)的数据,将行数从n²减少到n(n+1)/2(约1.125e10行),直接减半数据量:
# 获取names列的索引映射 name_to_idx = {name: idx for idx, name in enumerate(df["names"].to_list())} # 执行unpivot后过滤出上三角(含对角线)数据 df_unpivoted = df.unpivot(index="names", variable_name="names_2", value_name="distance") df_optimized = df_unpivoted.filter( pl.col("names").map_dict(name_to_idx) <= pl.col("names_2").map_dict(name_to_idx) )
2. 字符串编码为整数
将names和names_2的字符串替换为整数ID,大幅降低字符串存储开销:
# 提取唯一名称并生成映射 unique_names = df["names"].unique() name_id_map = pl.DataFrame({"names": unique_names, "id": pl.int_range(0, len(unique_names))}) # 替换原DataFrame的名称为ID df_id = df.join(name_id_map, on="names").drop("names").rename({"id": "names"}) # 将列名也替换为ID df_id = df_id.rename({col: str(name_id_map.filter(pl.col("names") == col)["id"].item()) for col in df_id.columns if col != "names"}) # 执行unpivot并处理names_2列 df_unpivoted_id = df_id.unpivot(index="names", variable_name="names_2", value_name="distance") df_unpivoted_id = df_unpivoted_id.with_columns(pl.col("names_2").cast(pl.Int64)) # (可选)如需恢复字符串,可保留映射表,后续按需join
3. 分块处理(缓解峰值内存)
如果操作过程中内存压力过大,可采用分块unpivot的方式,逐步处理并写入磁盘:
chunk_size = 1000 # 根据内存调整块大小 for i in range(0, len(df), chunk_size): chunk = df.slice(i, chunk_size) chunk_unpivoted = chunk.unpivot(index="names", variable_name="names_2", value_name="distance") # 过滤冗余数据 chunk_optimized = chunk_unpivoted.filter(pl.col("names").map_dict(name_to_idx) <= pl.col("names_2").map_dict(name_to_idx)) # 写入磁盘(如parquet) chunk_optimized.write_parquet(f"chunk_{i}.parquet") # 最后合并所有块 df_final = pl.read_parquet("chunk_*.parquet")
4. 选择更高效的存储类型
- 将
distance列的类型从i64改为更小的整数类型(如i32/i16),如果数值范围允许; - 最终结果存储为Parquet格式,利用其列存储和压缩特性进一步减小体积。
内容的提问来源于stack exchange,提问作者Nils R
相关产品推荐
相关产品推荐

