Spark DataFrame中为run_key非零连续块生成唯一分组键
Alright, let's walk through how to add that unique group key column for your non-zero continuous run_key blocks. The core idea is to first identify each continuous non-zero segment, then generate a unique identifier for each segment—whether it's the first value, average, or just a sequential ID, as long as it's unique per block.
Step 1: Prep the DataFrame with a Non-Zero Marker
First, we'll add a helper column to flag rows where run_key is non-zero. This helps us track where the continuous blocks start and end.
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ // Assume your original DataFrame is named `df` val markedDF = df.withColumn("is_non_zero", when($"run_key" =!= 0, 1).otherwise(0))
Step 2: Identify Continuous Non-Zero Blocks
Next, we'll use a cumulative sum window function to assign a unique ID to each continuous non-zero segment. Important: You need an ordered column (like an ID, timestamp, or sequence number) to ensure the window function processes rows in the correct order—without this, Spark can't reliably identify continuous blocks.
// Define a window ordered by your sequential column (replace `id` with your actual order column) val orderWindow = Window.orderBy($"id") // Calculate a cumulative sum of the non-zero marker to group continuous segments val groupedDF = markedDF.withColumn("group_id", sum($"is_non_zero").over(orderWindow)) // Optional: Set group_id to null for zero rows (since we don't need to group those) val cleanedGroupDF = groupedDF.withColumn("group_id", when($"is_non_zero" === 0, null).otherwise($"group_id"))
Step 3: Generate Your Unique Group Key
Now that we have group_id identifying each continuous block, you can generate a key based on run_key values in the block. Here are a few options:
Option 1: Use the First run_key Value in the Block
This is simple and ensures uniqueness as long as each block has at least one non-zero value:
val groupWindow = Window.partitionBy($"group_id") val resultDF = cleanedGroupDF.withColumn("unique_group_key", first($"run_key").over(groupWindow))
Option 2: Use the Average of run_key in the Block
If you prefer a value representative of the entire block:
val resultDF = cleanedGroupDF.withColumn("unique_group_key", avg($"run_key").over(groupWindow))
Option 3: Just Use group_id Directly
The group_id itself is already a unique identifier for each block, so you can skip generating a separate key if it fits your use case.
Example Output
Let's say your original data looks like this:
| id | run_key |
|---|---|
| 1 | 0 |
| 2 | 5 |
| 3 | 6 |
| 4 | 0 |
| 5 | 7 |
After processing, your resultDF will look like this (using Option 1):
| id | run_key | is_non_zero | group_id | unique_group_key |
|---|---|---|---|---|
| 1 | 0 | 0 | null | null |
| 2 | 5 | 1 | 1 | 5 |
| 3 | 6 | 1 | 1 | 5 |
| 4 | 0 | 0 | null | null |
| 5 | 7 | 1 | 2 | 7 |
Now you can use unique_group_key to group and compare values across different non-zero blocks!
内容的提问来源于stack exchange,提问作者Ajar

