Spark Scala实现按条件填充首行ID至后续行(新增列对比)
Got it, let's tackle this problem step by step. The core idea is to create a grouping logic based on non-zero fn values, then propagate the right id through each group until we hit the next non-zero fn row. Here's a practical implementation:
Step 1: Sort Your Data First
Since our logic depends entirely on the order of rows (based on id), we need to sort the DataFrame first to ensure correct processing order.
import org.apache.spark.sql.functions._ import org.apache.spark.sql.expressions.Window // Assume your input DataFrame is named `df` with columns `id` (Int) and `fn` (Int) val sortedDf = df.orderBy("id")
Step 2: Get the Next Row's ID
We use the lead window function to fetch the id of the next row—this will be our new "base" ID whenever we encounter a non-zero fn value.
val withNextId = sortedDf.withColumn( "next_id", lead(col("id"), 1).over(Window.orderBy("id")) )
Step 3: Mark Trigger Points for ID Switch
Create a column that only stores the next row's id when fn is non-zero. These are our trigger points where we need to switch the propagated ID.
val withTriggerId = withNextId.withColumn( "trigger_new_id", when(col("fn") =!= 0, col("next_id")).otherwise(null) )
Step 4: Propagate the ID Through Each Group
Use the last window function (with ignoreNulls=true) to carry forward the most recent trigger ID. For the initial group (before any non-zero fn), we fall back to the first id in the dataset using coalesce.
val resultDf = withTriggerId.withColumn( "new_id", coalesce( last(col("trigger_new_id"), ignoreNulls = true).over(Window.orderBy("id")), first(col("id")).over(Window.orderBy("id")) ) ).drop("next_id", "trigger_new_id") // Clean up intermediate columns we don't need
Step 5: Inspect the Result
You can now view the output to compare the original id with the generated new_id:
resultDf.show()
Example Output
For input data like this:
| id | fn |
|---|---|
| 1 | 0 |
| 2 | 0 |
| 3 | 5 |
| 4 | 0 |
| 5 | 0 |
| 6 | 3 |
| 7 | 0 |
The result will look like this:
| id | fn | new_id |
|---|---|---|
| 1 | 0 | 1 |
| 2 | 0 | 1 |
| 3 | 5 | 1 |
| 4 | 0 | 4 |
| 5 | 0 | 4 |
| 6 | 3 | 4 |
| 7 | 0 | 7 |
Handling Edge Cases
- First row has non-zero
fn: The first row uses its ownid, then switches to the next row'sidfor all subsequent rows. - No non-zero
fnvalues: Every row will use the firstidin the dataset. - Last row has non-zero
fn: The final row keeps the current group'sid(since there's no next row to switch to).
内容的提问来源于stack exchange,提问作者vrreddy1234

