Spark Structured Streaming 2.3.0分组取最大行:SQL及无指定算子实现问询
Let's break down these two common streaming processing questions and walk through practical solutions tailored for Spark 2.3.0.
1. 使用原生Spark SQL选取每个分组的最大行
In Spark Structured Streaming 2.3.0, window functions are fully supported for stream processing, making this task straightforward. We'll use ROW_NUMBER() to rank records within each group and pick the top one (by my_timestamp in your example).
Step-by-Step Implementation
First, assume you've loaded your stream into a DataFrame streamDF with the schema matching your input:
val streamDF = spark.readStream .format("csv") .option("header", "true") .option("inferSchema", "true") .load("/your/stream/source/path")
Register the stream DataFrame as a temporary view to use raw Spark SQL:
streamDF.createOrReplaceTempView("transactions")
Write the SQL query to select the latest row per id:
val resultDF = spark.sql(""" SELECT id, amount, my_timestamp FROM ( SELECT id, amount, my_timestamp, -- Rank rows in each id group by timestamp descending ROW_NUMBER() OVER (PARTITION BY id ORDER BY my_timestamp DESC) AS row_rank FROM transactions ) ranked_transactions WHERE row_rank = 1 """)
Finally, output the result (e.g., to console):
resultDF.writeStream .outputMode("update") .format("console") .start() .awaitTermination()
This query will consistently output the latest record for each id as new data streams in.
2. 不使用order by或mapGroupWithState选取每个分组的最大行
If you can't use order by (even in window functions) or mapGroupWithState, we can solve this with an aggregate-then-join approach. The idea is to first find the maximum timestamp per group, then join back to the original stream to get the full record.
Step-by-Step Implementation
First, compute the maximum my_timestamp for each id:
val maxTimestampPerIdDF = streamDF.groupBy("id") .agg(max("my_timestamp").alias("latest_ts"))
Join this aggregated DataFrame back to the original stream to fetch the full matching record:
val resultDF = streamDF.join( maxTimestampPerIdDF, // Match on id and the latest timestamp streamDF("id") === maxTimestampPerIdDF("id") && streamDF("my_timestamp") === maxTimestampPerIdDF("latest_ts"), "inner" ).select(streamDF("id"), streamDF("amount"), streamDF("my_timestamp"))
Output the result the same way as before:
resultDF.writeStream .outputMode("update") .format("console") .start() .awaitTermination()
Note: If multiple records share the same maximum timestamp for an id, this method will return all of them. If you need a single record, you can add dropDuplicates("id") at the end, but this doesn't use order by so it's still compliant with your constraints.
Example Validation
For your input data:
id | amount | my_timestamp
1 | 5 | 2018-04-01T01:00:00.000Z
1 | 10 | 2018-04-01T01:10:00.000Z
2 | 20 | 2018-04-01T01:20:00.000Z
2 | 30 | 2018-04-01T01:25:00.000Z
2 | 40 | 2018-04-01T01:30:00.000Z
Both methods will output your expected result:
id | amount | my_timestamp
1 | 10 | 2018-04-01T01:10:00.000Z
2 | 40 | 2018-04-01T01:30:00.000Z
内容的提问来源于stack exchange,提问作者user1870400

