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

如何在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完全对应的逻辑,步骤如下:

  1. 导入PySpark相关函数与窗口类:
from pyspark.sql import Window
from pyspark.sql import functions as F
  1. 定义分组窗口:
# 按ENTER12333_NUM字段分组的窗口规则
window_spec = Window.partitionBy("ENTER12333_NUM")
  1. 计算并添加标记字段:
# 复制原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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 12:15:35