如何用Polars高效实现List[Struct]的累计聚合转换?
问题:Polars高效实现带动态聚合与过滤的累计行转换
给定如下Polars DataFrame:
import polars as pl original_dataframe = pl.DataFrame({ 'index': ['A', 'B', 'C', 'D', 'E', 'F', 'G'], 'content': [ {'key': 3, 'val': 20}, {'key': 4, 'val': 50}, {'key': 3, 'val': 8}, {'key': 5, 'val': 70}, {'key': 4, 'val': -60}, {'key': 2, 'val': 30}, {'key': 4, 'val': 5} ] })
原始结构:
┌───────┬───────────┐ │ index ┆ content │ │ --- ┆ --- │ │ str ┆ struct[2] │ ╞═══════╪═══════════╡ │ A ┆ {3,20} │ │ B ┆ {4,50} │ │ C ┆ {3,8} │ │ D ┆ {5,70} │ │ E ┆ {4,-60} │ │ F ┆ {2,30} │ │ G ┆ {4,5} │ └───────┴───────────┘
需按以下规则转换为目标结构:
- 逐行将
content中的struct累计添加到列表中 - 列表中存在相同
key的struct时,对val字段求和聚合 - 聚合后
val<=0时,立即从当前行的列表中删除该struct;后续行若该key再次出现且val>0,重新累计聚合 - 每个列表按struct的
key字段排序 - 若struct的
key或val为null,直接删除该struct
转换后目标结果:
┌───────┬──────────────────────────┐ │ index ┆ content │ │ --- ┆ --- │ │ str ┆ list[struct[2]] │ ╞═══════╪══════════════════════════╡ │ A ┆ [{3,20}] │ │ B ┆ [{3,20}, {4,50}] │ │ C ┆ [{3,28}, {4,50}] │ │ D ┆ [{3,28}, {4,50}, {5,70}] │ │ E ┆ [{3,28}, {5,70}] │ │ F ┆ [{2,30}, {3,28}, {5,70}] │ │ G ┆ [{2,30}, {3,28}, {4,5}, {5,70}] │ └───────┴──────────────────────────┘
补充示例验证:
输入:
pl.DataFrame({ 'index': ['A', 'B', 'C', 'D', 'E', 'F'], 'content': [ {'key': 3, 'val': 20}, {'key': 4, 'val': 50}, {'key': 3, 'val': 8}, {'key': 2, 'val': 30}, {'key': 4, 'val': -60}, {'key': 4, 'val': 5} ] })
输出:
┌───────┬──────────────────────────┐ │ index ┆ content │ │ --- ┆ --- │ │ str ┆ list[struct[2]] │ ╞═══════╪══════════════════════════╡ │ A ┆ [{3,20}] │ │ B ┆ [{3,20}, {4,50}] │ │ C ┆ [{3,28}, {4,50}] │ │ D ┆ [{2,30}, {3,28}, {4,50}] │ │ E ┆ [{2,30}, {3,28}] │ │ F ┆ [{2,30}, {3,28}, {4,5}] │ └───────┴──────────────────────────┘
当前使用iter_rows()结合Python原生列表、字典迭代实现,速度较慢,求纯Polars函数的高效实现方案。
解决方案:纯Polars函数实现高效转换
核心思路是利用Polars的窗口函数、分组聚合和列表操作,完全基于向量化计算实现,避免Python循环的性能损耗。
完整代码实现:
import polars as pl def transform_df(df: pl.DataFrame) -> pl.DataFrame: return ( df # 提取struct字段并过滤null值 .with_columns( pl.col("content").struct.field("key").alias("key"), pl.col("content").struct.field("val").alias("val") ) .filter(pl.col("key").is_not_null() & pl.col("val").is_not_null()) # 添加行号用于窗口计算 .with_row_index("row_num", offset=1) # 计算每个key的累计和,累计和<=0时重置(仅保留当前val>0的情况) .with_columns( pl.col("val") .cum_sum() .over("key") .alias("cum_val") ) .with_columns( pl.when(pl.col("cum_val") <= 0) .then(pl.when(pl.col("val") > 0).then(pl.col("val")).otherwise(0)) .otherwise(pl.col("cum_val")) .over("key") .alias("effective_cum_val") ) # 逐行收集有效key-val对,排序后打包为struct列表 .with_columns( pl.struct(pl.col("key"), pl.col("effective_cum_val").alias("val")) .filter(pl.col("effective_cum_val") > 0) .sort_by("key") .over(pl.int_range(1, pl.count() + 1)) .alias("content") ) # 保留目标列并恢复原始顺序 .select("index", "content") ) # 测试原始数据 original_dataframe = pl.DataFrame({ 'index': ['A', 'B', 'C', 'D', 'E', 'F', 'G'], 'content': [ {'key': 3, 'val': 20}, {'key': 4, 'val': 50}, {'key': 3, 'val': 8}, {'key': 5, 'val': 70}, {'key': 4, 'val': -60}, {'key': 2, 'val': 30}, {'key': 4, 'val': 5} ] }) result = transform_df(original_dataframe) print(result) # 测试补充示例 supplement_df = pl.DataFrame({ 'index': ['A', 'B', 'C', 'D', 'E', 'F'], 'content': [ {'key': 3, 'val': 20}, {'key': 4, 'val': 50}, {'key': 3, 'val': 8}, {'key': 2, 'val': 30}, {'key': 4, 'val': -60}, {'key': 4, 'val': 5} ] }) supplement_result = transform_df(supplement_df) print(supplement_result)
代码关键逻辑说明:
- 展开与过滤:将struct拆分为独立的
key和val列,直接过滤含null的行,满足规则5。 - 动态累计和计算:通过
cum_sum().over("key")计算每个key的累计值,再用条件逻辑处理累计和<=0的情况——若当前val>0则重置为当前val,否则置0,确保后续该key再次出现时能重新累计。 - 逐行生成结果列表:使用
over(pl.int_range(1, pl.count()+1))实现逐行的窗口聚合,收集所有有效(累计和>0)的key-val对,排序后打包为struct列表,满足规则1-4。
性能优势:
完全依赖Polars的向量化引擎,避免了Python循环的额外开销,处理大数据集时速度远优于iter_rows()的实现。
内容的提问来源于stack exchange,提问作者sejabs
相关产品推荐
相关产品推荐

