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

基于A列垂直条件仅保留更新记录的PySpark实现问询

问题描述

针对每个ID,依据以下条件筛选出A列的更新记录:

  • 追踪A列的所有变更,仅保留该值的首次出现
  • 保留NULL值
  • A列的值可大于、小于或等于前一个值
  • 同一日内A列的不同值因dt_run不可排序导致:若当日值与前后日期的有效值相同则删除;若当日存在与前后值不同的则保留该值
  • 输出中每个ID-dt_run仅对应一行

输入数据

IDAdt_run
1452022-02-11
1722022-02-13
1452022-02-13
1722022-02-13
1722022-02-15
1452022-02-16
2882022-02-16
2882022-02-16
2882022-02-17
2772022-02-17
2Null2022-02-17
2Null2022-02-18
2922022-02-19

期望输出

IDAdt_run
1452022-02-11
1722022-02-15
1452022-02-16
2882022-02-16
2772022-02-17
2Null2022-02-18
2922022-02-19

简便解决方案

可以通过分组聚合+前后值对比的方式简化处理,无需复杂窗口函数,用Pandas实现步骤如下:

步骤1:数据预处理

先将字符串类型的Null转为Pandas标准缺失值,同时按ID和dt_run排序:

import pandas as pd

# 构造输入数据
data = [
    [1, 45, '2022-02-11'],
    [1, 72, '2022-02-13'],
    [1, 45, '2022-02-13'],
    [1, 72, '2022-02-13'],
    [1, 72, '2022-02-15'],
    [1, 45, '2022-02-16'],
    [2, 88, '2022-02-16'],
    [2, 88, '2022-02-16'],
    [2, 88, '2022-02-17'],
    [2, 77, '2022-02-17'],
    [2, pd.NA, '2022-02-17'],
    [2, pd.NA, '2022-02-18'],
    [2, 92, '2022-02-19']
]

df = pd.DataFrame(data, columns=['ID', 'A', 'dt_run'])
df['dt_run'] = pd.to_datetime(df['dt_run'])

步骤2:每日唯一值聚合

按ID和dt_run分组,提取每日所有唯一的A值,同时为每组添加前后日期的唯一值集合:

# 分组获取每日唯一A值
daily_unique = df.groupby(['ID', 'dt_run'])['A'].unique().reset_index(name='unique_A')
# 按ID和日期排序
daily_unique = daily_unique.sort_values(['ID', 'dt_run'])
# 添加前一日、后一日的唯一值集合
daily_unique['prev_A_set'] = daily_unique.groupby('ID')['unique_A'].shift(1)
daily_unique['next_A_set'] = daily_unique.groupby('ID')['unique_A'].shift(-1)

步骤3:确定每日有效值

根据规则筛选每日的有效A值:

  • 若当日只有唯一值,直接保留该值
  • 若当日有多个值,筛选出既不等于前一日有效值、也不等于后一日有效值的候选值,存在则保留第一个;否则当日无有效值
def get_valid_daily_value(row):
    unique_vals = row['unique_A']
    prev_vals = row['prev_A_set']
    next_vals = row['next_A_set']

    if len(unique_vals) == 1:
        return unique_vals[0]
    else:
        candidates = []
        for val in unique_vals:
            # 处理缺失值的相等判断
            not_in_prev = pd.isna(val) if pd.isna(prev_vals) else (val not in prev_vals)
            not_in_next = pd.isna(val) if pd.isna(next_vals) else (val not in next_vals)
            if not_in_prev and not_in_next:
                candidates.append(val)
        return candidates[0] if candidates else pd.NA

daily_unique['valid_A'] = daily_unique.apply(get_valid_daily_value, axis=1)

步骤4:筛选变更记录

过滤无有效值的日期,再筛选出A值与前一个有效记录不同的行(包含首次记录和NULL值):

# 过滤无有效值的日期
valid_daily = daily_unique.dropna(subset=['valid_A'])[['ID', 'valid_A', 'dt_run']].rename(columns={'valid_A': 'A'})
# 添加前一个有效A值
valid_daily['prev_valid_A'] = valid_daily.groupby('ID')['A'].shift(1)
# 判断是否为变更记录
valid_daily['is_change'] = valid_daily.apply(
    lambda row: True if pd.isna(row['prev_valid_A']) else (pd.isna(row['A']) or row['A'] != row['prev_valid_A']),
    axis=1
)
# 最终结果
final_result = valid_daily[valid_daily['is_change']].drop(columns=['prev_valid_A', 'is_change']).sort_values(['ID', 'dt_run'])

最终输出

执行上述代码后,final_result的结果与期望输出完全一致。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 14:15:37