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

Spark Structured Streaming 2.3.0分组取最大行:SQL及无指定算子实现问询

解决Spark Structured Streaming 2.3.0的分组取最大行问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 03:54:30