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

Spark Scala:如何将列传入UDF并实现条件分组取最新行

Got it, let's work through this Spark Scala data processing task step by step. First, let's recap the core requirements to make sure we're aligned:

We need to dynamically group rows in a DataFrame and keep only the row with the latest Timestamp per group, following these rules:

  • When FFAction is I: group by the first 6 columns
  • When FFAction is O (and UpdateReason_updateReasonId is null) OR FFAction is D: group by the first 5 columns
  • When a group (by first 5 columns) contains both I rows and either qualifying O/D rows: use the first 5 columns as the grouping key for all rows in that set, then pick the latest timestamp

Let's implement this with Spark Scala

First, import the necessary functions and window expressions:

import org.apache.spark.sql.functions._
import org.apache.spark.sql.expressions.Window

Assuming your DataFrame is named rawDf with columns: col1, col2, col3, col4, col5, col6, FFAction, UpdateReason_updateReasonId, Timestamp (replace these with your actual column names).

Step 1: Mark initial grouping level for each row

We start by tagging each row with its intended grouping level based on the first two rules:

val dfWithInitialGroup = rawDf.withColumn(
  "initial_group_level",
  when(col("FFAction") === "I", 6)
    .when(
      (col("FFAction") === "O" && col("UpdateReason_updateReasonId").isNull) || 
      col("FFAction") === "D", 
      5
    )
    .otherwise(6) // Handle edge cases (e.g., unrecognized FFAction) as needed
)

Step 2: Detect mixed-type groups (I + qualifying O/D)

Next, we check if any group (by first 5 columns) has both I rows and qualifying O/D rows. We use a window partitioned by the first 5 columns to collect all distinct FFAction values in each group:

val group5Window = Window.partitionBy("col1", "col2", "col3", "col4", "col5")

val dfWithGroupCheck = dfWithInitialGroup.withColumn(
  "distinct_actions", collect_set("FFAction").over(group5Window)
).withColumn(
  "has_mixed_types",
  // Check if group contains I AND (qualifying O or D)
  array_contains(col("distinct_actions"), "I") && (
    array_contains(col("distinct_actions"), "O") || 
    array_contains(col("distinct_actions"), "D")
  )
)

Step 3: Define final grouping keys

Now we set the final grouping key: if the group has mixed types, we use the first 5 columns for all rows; otherwise, we use the initial grouping level's columns:

val dfWithFinalGroup = dfWithGroupCheck.withColumn(
  "final_group_key",
  when(col("has_mixed_types"), 
       struct("col1", "col2", "col3", "col4", "col5")
      )
      .when(col("initial_group_level") === 6, 
            struct("col1", "col2", "col3", "col4", "col5", "col6")
           )
      .otherwise(
        struct("col1", "col2", "col3", "col4", "col5")
      )
)

Step 4: Get the latest row per final group

Finally, we use a window partitioned by our final grouping key, sort rows by Timestamp descending, and pick the first row (latest timestamp) from each group:

val latestRowWindow = Window.partitionBy("final_group_key").orderBy(col("Timestamp").desc)

val finalResult = dfWithFinalGroup.withColumn(
  "row_rank", row_number().over(latestRowWindow)
).filter(col("row_rank") === 1)
// Drop helper columns we created
.drop("initial_group_level", "distinct_actions", "has_mixed_types", "final_group_key", "row_rank")

How this works with sample data

Let's say you have these rows:

col1col2col3col4col5col6FFActionUpdateReason_updateReasonIdTimestamp
ABCDEFInull2023-10-01 10:00:00
ABCDEFOnull2023-10-01 11:00:00
ABCDEGInull2023-10-01 09:00:00
ABCDEGDnull2023-10-01 12:00:00
XYZQWRO1232023-10-01 08:00:00
XYZQWRInull2023-10-01 07:00:00

The result will be:

  • For the A,B,C,D,E group: mixed I + O/D, so we keep the row with timestamp 2023-10-01 12:00:00 (D row)
  • For X,Y,Z,Q,W,R: the O row has a non-null UpdateReason_updateReasonId, so it doesn't qualify for rule 2. Thus, we group by all 6 columns, keeping both the O (08:00) and I (07:00) rows since they're in separate 6-column groups.

内容的提问来源于stack exchange,提问作者Shailendra

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 09:36:01