基于A列垂直条件仅保留更新记录的PySpark实现问询
问题描述
针对每个ID,依据以下条件筛选出A列的更新记录:
- 追踪A列的所有变更,仅保留该值的首次出现
- 保留NULL值
- A列的值可大于、小于或等于前一个值
- 同一日内A列的不同值因dt_run不可排序导致:若当日值与前后日期的有效值相同则删除;若当日存在与前后值不同的则保留该值
- 输出中每个ID-dt_run仅对应一行
输入数据
| ID | A | dt_run |
|---|---|---|
| 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 | Null | 2022-02-17 |
| 2 | Null | 2022-02-18 |
| 2 | 92 | 2022-02-19 |
期望输出
| ID | A | dt_run |
|---|---|---|
| 1 | 45 | 2022-02-11 |
| 1 | 72 | 2022-02-15 |
| 1 | 45 | 2022-02-16 |
| 2 | 88 | 2022-02-16 |
| 2 | 77 | 2022-02-17 |
| 2 | Null | 2022-02-18 |
| 2 | 92 | 2022-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
相关产品推荐
相关产品推荐

