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
FFActionisI: group by the first 6 columns - When
FFActionisO(andUpdateReason_updateReasonIdis null) ORFFActionisD: group by the first 5 columns - When a group (by first 5 columns) contains both
Irows and either qualifyingO/Drows: 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:
| col1 | col2 | col3 | col4 | col5 | col6 | FFAction | UpdateReason_updateReasonId | Timestamp |
|---|---|---|---|---|---|---|---|---|
| A | B | C | D | E | F | I | null | 2023-10-01 10:00:00 |
| A | B | C | D | E | F | O | null | 2023-10-01 11:00:00 |
| A | B | C | D | E | G | I | null | 2023-10-01 09:00:00 |
| A | B | C | D | E | G | D | null | 2023-10-01 12:00:00 |
| X | Y | Z | Q | W | R | O | 123 | 2023-10-01 08:00:00 |
| X | Y | Z | Q | W | R | I | null | 2023-10-01 07:00:00 |
The result will be:
- For the
A,B,C,D,Egroup: mixed I + O/D, so we keep the row with timestamp2023-10-01 12:00:00(D row) - For
X,Y,Z,Q,W,R: the O row has a non-nullUpdateReason_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

