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

如何在列中首次出现null后传播null?含分组id批量处理

解决方案:Polars实现匹配后首次不匹配的Null传播(含分组场景)

核心需求

  • 按预期数据集的顺序逐行匹配观测数据的name字段
  • 首次出现匹配不一致(观测无对应行或name不相等)时,该行及后续所有行的description置为null
  • 支持按id分组,每个id独立执行上述逻辑

基础场景(无分组)

步骤说明

  1. 给两个数据集添加全局行索引,确保按原始顺序匹配
  2. 左连接保留预期数据集的所有行,匹配对应行的观测数据
  3. 标记匹配不一致的行,将description设为null
  4. 使用累积窗口函数将第一个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独立处理)

步骤说明

  1. 给两个数据集按id分组,添加组内行索引(每个id内部按原始顺序编号)
  2. 按id+组内行索引左连接,确保每个id内部的行顺序匹配
  3. 标记每个id内匹配不一致的行
  4. 按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,以此判断是否需要置为Null
  • over("id"):将窗口计算限定在每个id分组内,确保分组间的Null传播互不干扰

内容的提问来源于stack exchange,提问作者DJDuque

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 13:27:15