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

PySpark中按ID筛选列匹配行 无匹配则保留分组首行

实现方案

你原来的写法只过滤了全局满足numb == real的行,没有处理分组内无匹配行的场景,可通过窗口函数加标记的方式实现需求:

PySpark 实现代码

from pyspark.sql.functions import col, when, max, row_number
from pyspark.sql.window import Window

# 定义分组窗口,按ID分区,按numb升序排序(如需按原始数据顺序,可提前新增monotonically_increasing_id作为排序键)
window_spec = Window.partitionBy("ID").orderBy("numb")

# 新增两个标记列:标记分组是否有匹配行、分组内行号
df = df.withColumn("has_match", max(when(col("numb") == col("real"), 1).otherwise(0)).over(window_spec)) \
       .withColumn("row_num", row_number().over(window_spec))

# 按规则筛选
df_result = df.filter(
    # 有匹配的分组,保留所有匹配行
    (col("has_match") == 1) & (col("numb") == col("real")) |
    # 无匹配的分组,保留首行
    (col("has_match") == 0) & (col("row_num") == 1)
).drop("has_match", "row_num")

# 输出结果验证
df_result.show()

运行后输出结果和预期一致:

+------+----+----+
|    ID|numb|real|
+------+----+----+
|213412|2008|2008|
|213410|2009|null|
|393859|2017|2017|
|393859|2021|2021|
+------+----+----+

Pandas 实现代码

如果你用的是Pandas,可以用groupby apply的方式实现:

import pandas as pd
import numpy as np

def filter_group(g):
    # 先找当前分组匹配的行
    match_rows = g[g["numb"] == g["real"]]
    if not match_rows.empty:
        return match_rows
    # 没有匹配返回首行
    return g.head(1)

df_result = df.groupby("ID", group_keys=False).apply(filter_group).reset_index(drop=True)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 09:48:01