PySpark中仅保留修改行的数据集清洗实现方案
数据集清洗:保留指定字段的状态变更记录
需求说明:
- 按
id分组,针对指定字段(示例为city和sport)过滤数据 - 仅保留相对于同组前一行发生状态变更的记录,每个状态首次出现时保留
- 若某行回到之前的非紧邻状态(比如之前出现过但不是上一行的状态),仍需保留该行
输入数据集(df1)
| id | city | sport | date |
|---|---|---|---|
| abc | london | football | 2022-02-11 |
| abc | paris | football | 2022-02-12 |
| abc | paris | football | 2022-02-13 |
| abc | paris | football | 2022-02-14 |
| abc | paris | football | 2022-02-15 |
| abc | london | football | 2022-02-16 |
| abc | paris | football | 2022-02-17 |
| def | paris | volley | 2022-02-10 |
| def | paris | volley | 2022-02-11 |
| ghi | manchester | basketball | 2022-02-09 |
期望输出数据集
| id | city | sport | date |
|---|---|---|---|
| abc | london | football | 2022-02-11 |
| abc | paris | football | 2022-02-12 |
| abc | london | football | 2022-02-16 |
| abc | paris | football | 2022-02-17 |
| def | paris | volley | 2022-02-10 |
| ghi | manchester | basketball | 2022-02-09 |
实现方案(Python Pandas)
通过分组对比相邻行的指定字段值即可实现需求:
import pandas as pd # 构造输入数据集 df = pd.DataFrame([ ["abc", "london", "football", "2022-02-11"], ["abc", "paris", "football", "2022-02-12"], ["abc", "paris", "football", "2022-02-13"], ["abc", "paris", "football", "2022-02-14"], ["abc", "paris", "football", "2022-02-15"], ["abc", "london", "football", "2022-02-16"], ["abc", "paris", "football", "2022-02-17"], ["def", "paris", "volley", "2022-02-10"], ["def", "paris", "volley", "2022-02-11"], ["ghi", "manchester", "basketball", "2022-02-09"] ], columns=["id", "city", "sport", "date"]) # 生成过滤掩码:同组内当前行与前一行的指定字段是否存在差异 mask = df.groupby("id")[["city", "sport"]].apply( lambda x: x.ne(x.shift()).any(axis=1) ).reset_index(drop=True) # 筛选符合条件的行 result_df = df[mask] print(result_df)
核心逻辑:
- 按
id分组,保证仅在同一用户的记录间做状态对比 - 用
shift()获取前一行的字段值,通过ne()判断是否存在差异,只要任一指定字段变更就保留该行 - 每组的第一行默认保留(因
shift()后前一行为空,ne()返回True)
内容的提问来源于stack exchange,提问作者Jresearcher
相关产品推荐
相关产品推荐

