如何在PySpark DataFrame中实现指定分组窗口函数标记逻辑?
PySpark DataFrame实现分组标记逻辑
数据结构示例
+---------+---------+--------------------+-----------------+---------+-----------------+--------------+-----------------------------+-------------+---------------+ |CURR_COL1|CURR_COL2| HASH_VALUE| CURR_COL3|CURR_COL4| CURR_COL45|ENTER12333_NUM|RECALC_ENTER12333_CALCUALTION|ROW_OPERATION|WINNER_IN222222| +---------+---------+--------------------+-----------------+---------+-----------------+--------------+-----------------------------+-------------+---------------+ | 75757| hello|9bc7d98d527bb54b1...| 79.0000000000| pb| 55.0000000000|440.0000000000| null| I| null| | 46| hello|9bc7d98d527bb54b1...| 79.0000000000| pb| 55.0000000000|590.0000000000| null| I| null| | 4545| Senorita|d95ee5d8db9958f6e...| 79.0000000000| null| 79.0000000000|590.0000000000| null| U| null| | 189899| hello|a93d52dad9dcc3bd0...| 79.0000000000| Purnima| 79.0000000000|890.0000000000| null| N| null| | 234223| goodbye|4271325117076d7b5...|454646.0000000000| bhatia|454646.0000000000|890.0000000000| null| D| null| +---------+---------+--------------------+-----------------+---------+-----------------+--------------+-----------------------------+-------------+---------------+
需求说明
按ENTER12333_NUM字段分组,若分组内存在ROW_OPERATION值为I/U/D的行,则将RECALC_ENTER12333_CALCUALTION字段标记为Y;若分组内没有符合条件的行,则标记为N。
已实现的SQL逻辑
select a.*, case when ( sum(case when ROW_OPERATION in ('I','U','D') then 1 else 0 end ) over (partition by ENTER12333_NUM) ) > 0 then 'Y' else 'N' end RECALC_ENTER12333_CALCUALTION from delta_firmographic_data a
PySpark DataFrame实现方案
可以通过窗口函数实现和SQL完全对应的逻辑,步骤如下:
- 导入PySpark相关函数与窗口类:
from pyspark.sql import Window from pyspark.sql import functions as F
- 定义分组窗口:
# 按ENTER12333_NUM字段分组的窗口规则 window_spec = Window.partitionBy("ENTER12333_NUM")
- 计算并添加标记字段:
# 复制原DataFrame并添加目标字段 result_df = df.withColumn( "RECALC_ENTER12333_CALCUALTION", F.when( # 统计分组内符合ROW_OPERATION条件的行数总和 F.sum( F.when(F.col("ROW_OPERATION").isin("I", "U", "D"), 1).otherwise(0) ).over(window_spec) > 0, "Y" ).otherwise("N") )
简化实现(可选)
如果只需要判断分组内是否存在符合条件的行,也可以用max函数替代sum,逻辑更直观:
result_df = df.withColumn( "RECALC_ENTER12333_CALCUALTION", F.when( F.max( F.when(F.col("ROW_OPERATION").isin("I", "U", "D"), 1).otherwise(0) ).over(window_spec) == 1, "Y" ).otherwise("N") )
以上两种实现都和原SQL逻辑一致,最终会为每一行添加对应的标记值。
内容的提问来源于stack exchange,提问作者pbh
相关产品推荐
相关产品推荐

