PySpark需求:标记delta>10的无效行并按person生成连续性递增列
Pandas DataFrame 连续性分组标记实现方案
需要对DataFrame完成以下操作:
- 当
delta列值大于10时,标记该行无效,对应continuity列值为NaN - 按
person字段分组后,每检测到一行无效行,后续所有有效行的continuity值自动递增1
输入DataFrame
| person | delta |
|---|---|
| X | 2 |
| X | 3 |
| X | 4 |
| X | 20 |
| X | 50 |
| X | 5 |
| Y | 1 |
| Y | 20 |
| Y | 2 |
| Y | 3 |
| Z | 9 |
| Z | 30 |
| Z | 2 |
| Z | 15 |
| Z | 3 |
期望输出DataFrame
| person | delta | continuity |
|---|---|---|
| X | 2 | 1 |
| X | 3 | 1 |
| X | 4 | 1 |
| X | 20 | NaN |
| X | 50 | NaN |
| X | 5 | 2 |
| Y | 1 | 1 |
| Y | 20 | NaN |
| Y | 2 | 2 |
| Y | 3 | 2 |
| Z | 9 | 1 |
| Z | 30 | NaN |
| Z | 2 | 2 |
| Z | 15 | NaN |
| Z | 3 | 3 |
实现代码
import pandas as pd # 构造输入DataFrame data = { 'person': ['X', 'X', 'X', 'X', 'X', 'X', 'Y', 'Y', 'Y', 'Y', 'Z', 'Z', 'Z', 'Z', 'Z'], 'delta': [2, 3, 4, 20, 50, 5, 1, 20, 2, 3, 9, 30, 2, 15, 3] } df = pd.DataFrame(data) # 标记无效行(delta>10) df['is_invalid'] = df['delta'] > 10 # 按person分组计算连续性分组值:累积无效行数量+1作为分组标识 df['continuity'] = df.groupby('person')['is_invalid'].cumsum() + 1 # 将无效行的continuity设为NaN df.loc[df['is_invalid'], 'continuity'] = pd.NA # 移除辅助列 df = df.drop(columns=['is_invalid']) print(df)
代码说明
- 新增
is_invalid辅助列快速标记无效行; - 按
person分组后对is_invalid做累积求和,每遇到一个无效行,累积值加1,再加上初始值1,得到有效行的连续性分组值; - 最后将无效行的
continuity设为NaN,并移除辅助列,得到目标结果。
内容的提问来源于stack exchange,提问作者user21017176
相关产品推荐
相关产品推荐

