如何在列中首次出现null后传播null?含分组id批量处理
解决方案:Polars实现匹配后首次不匹配的Null传播(含分组场景)
核心需求
- 按预期数据集的顺序逐行匹配观测数据的
name字段 - 首次出现匹配不一致(观测无对应行或
name不相等)时,该行及后续所有行的description置为null - 支持按
id分组,每个id独立执行上述逻辑
基础场景(无分组)
步骤说明
- 给两个数据集添加全局行索引,确保按原始顺序匹配
- 左连接保留预期数据集的所有行,匹配对应行的观测数据
- 标记匹配不一致的行,将
description设为null - 使用累积窗口函数将第一个
null传播至后续所有行
代码实现
import polars as pl # 初始化数据 expected = pl.DataFrame({ "name": ["start", "stop", "start", "stop", "start", "stop", "start", "stop"], "description": ["a", "b", "c", "d", "e", "f", "g", "h"], }) observed = pl.DataFrame({ "name": ["start", "stop", "start", "stop", "stop", "stop", "start"], "time": [0.1, 0.2, 0.3, 0.4, 0.5, 0.6, 0.7], }) # 1. 添加全局行索引 expected = expected.with_row_index("idx") observed = observed.with_row_index("idx") # 2. 左连接匹配对应行 joined = expected.join(observed, on="idx", how="left", suffix="_observed") # 3. 标记匹配不一致的行 joined = joined.with_columns( pl.when( pl.col("name_observed").is_null() | (pl.col("name") != pl.col("name_observed")) ).then(None).otherwise(pl.col("description")).alias("description_matched") ) # 4. 传播Null:首次出现Null后,后续所有行均为Null result = joined.with_columns( pl.when( pl.col("description_matched").is_null().cumany() ).then(None).otherwise(pl.col("description_matched")).alias("description_final") ).select("idx", "name", "description_final", "time") print(result)
分组场景(按id独立处理)
步骤说明
- 给两个数据集按
id分组,添加组内行索引(每个id内部按原始顺序编号) - 按
id+组内行索引左连接,确保每个id内部的行顺序匹配 - 标记每个id内匹配不一致的行
- 按
id分组使用累积窗口函数,传播Null至该id内后续所有行
代码实现
import polars as pl # 初始化分组数据 observed = pl.DataFrame({ "id": [1, 2, 1, 2, 2], "name": ["start", "start", "stop", "stop", "stop"], "time": [0.1, 0.2, 0.3, 0.4, 0.5], }) expected = pl.DataFrame({ "id": [1, 1, 2, 2], "name": ["start", "stop", "start", "stop"], "description": ["a", "b", "c", "d"], }) # 1. 添加组内行索引 expected = expected.with_columns( pl.int_range(0, pl.count()).over("id").alias("group_idx") ) observed = observed.with_columns( pl.int_range(0, pl.count()).over("id").alias("group_idx") ) # 2. 按id+group_idx左连接 joined = expected.join(observed, on=["id", "group_idx"], how="left", suffix="_observed") # 3. 标记每个id内匹配不一致的行 joined = joined.with_columns( pl.when( pl.col("name_observed").is_null() | (pl.col("name") != pl.col("name_observed")) ).then(None).otherwise(pl.col("description")).alias("description_matched") ) # 4. 按id分组传播Null result = joined.with_columns( pl.when( pl.col("description_matched").is_null().cumany().over("id") ).then(None).otherwise(pl.col("description_matched")).alias("description_final") ).select("id", "name", "description_final", "time") print(result)
关键函数说明
cumany():累积布尔函数,一旦出现True(即首次Null),后续所有行均为True,以此判断是否需要置为Nullover("id"):将窗口计算限定在每个id分组内,确保分组间的Null传播互不干扰
内容的提问来源于stack exchange,提问作者DJDuque
相关产品推荐
相关产品推荐

