Spark SQL/DataFrame统计指定日期范围每日活跃航班数
实现方案
核心思路
摒弃逐天展开全量有效期的冗余逻辑,采用事件增量累计法,性能远高于rpad+posexplode的实现,尤其适配航班有效期跨度大的场景,不会产生冗余中间数据:
- 把每条航班记录拆成两个事件:
start_date标记为+1(航班当日生效,活跃计数+1),end_date + 1天标记为-1(航班在end_date当日仍有效,次日失效,计数-1) - 补全统计窗口的边界日期,避免窗口内无事件时出现数据缺失
- 按日期排序对增量做累计求和,得到当日活跃航班数,最后关联统计区间内的连续日期补全无事件的日期即可
代码实现
1. Spark SQL 版本
-- 测试样例数据(实际使用替换为自身表即可) WITH flight_df AS ( SELECT 'r1' AS user, 'f1' AS flight_id, DATE '2022-05-08' AS start_date, DATE '2022-05-09' AS end_date UNION ALL SELECT 'r1' AS user, 'f2' AS flight_id, DATE '2022-05-10' AS start_date, DATE '2022-05-12' AS end_date UNION ALL SELECT 'r2' AS user, 'f3' AS flight_id, DATE '2022-05-07' AS start_date, DATE '2022-05-15' AS end_date ), -- 构造增减事件 event_df AS ( SELECT user, start_date AS stat_date, 1 AS delta FROM flight_df WHERE user = 'r1' UNION ALL SELECT user, DATE_ADD(end_date, 1) AS stat_date, -1 AS delta FROM flight_df WHERE user = 'r1' -- 补全统计窗口边界 UNION ALL SELECT 'r1' AS user, DATE '2022-05-08' AS stat_date, 0 AS delta UNION ALL SELECT 'r1' AS user, DATE '2022-05-11' AS stat_date, 0 AS delta ), -- 按日聚合增量,计算累计活跃值 daily_cnt AS ( SELECT stat_date, SUM(delta) OVER (ORDER BY stat_date ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS active_flight_cnt FROM (SELECT stat_date, SUM(delta) AS delta FROM event_df GROUP BY stat_date) t ) -- 生成连续日期,关联输出结果 SELECT date_col AS stat_date, LAST_VALUE(active_flight_cnt, true) OVER (ORDER BY date_col ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS active_flight_cnt FROM ( SELECT EXPLODE(SEQUENCE(DATE '2022-05-08', DATE '2022-05-10', INTERVAL 1 DAY)) AS date_col ) dates LEFT JOIN daily_cnt ON dates.date_col = daily_cnt.stat_date ORDER BY stat_date;
如果使用Spark 3.2+版本,最后一步的空值补全可以替换为更简洁的
WATERMARK语法。
2. Scala DataFrame API 版本
import org.apache.spark.sql.functions._ import org.apache.spark.sql.expressions.Window import java.sql.Date // 1. 测试数据(实际使用替换为自身DataFrame即可) val flightDF = Seq( ("r1", "f1", Date.valueOf("2022-05-08"), Date.valueOf("2022-05-09")), ("r1", "f2", Date.valueOf("2022-05-10"), Date.valueOf("2022-05-12")), ("r2", "f3", Date.valueOf("2022-05-07"), Date.valueOf("2022-05-15")) ).toDF("user", "flight_id", "start_date", "end_date") // 统计参数 val targetUser = "r1" val statStart = Date.valueOf("2022-05-08") val statEnd = Date.valueOf("2022-05-10") // 2. 构造增减事件 val eventDF = flightDF .filter(col("user") === targetUser) .select(col("start_date").as("stat_date"), lit(1).as("delta")) .union( flightDF .filter(col("user") === targetUser) .select(date_add(col("end_date"), 1).as("stat_date"), lit(-1).as("delta")) ) .union(Seq((statStart, 0), (Date.valueOf("2022-05-11"), 0)).toDF("stat_date", "delta")) .groupBy("stat_date") .agg(sum("delta").as("delta")) // 3. 计算累计活跃值 val dailyCntDF = eventDF .withColumn("active_flight_cnt", sum("delta").over(Window.orderBy("stat_date").rowsBetween(Window.unboundedPreceding, Window.currentRow))) .select("stat_date", "active_flight_cnt") // 4. 生成连续日期,关联补全结果 val dateRangeDF = spark.sql(s"SELECT EXPLODE(SEQUENCE(DATE '${statStart}', DATE '${statEnd}', INTERVAL 1 DAY)) AS stat_date") val resultDF = dateRangeDF .join(dailyCntDF, Seq("stat_date"), "left") .withColumn("active_flight_cnt", last("active_flight_cnt", ignoreNulls = true).over(Window.orderBy("stat_date").rowsBetween(Window.unboundedPreceding, Window.currentRow))) .orderBy("stat_date") // 输出结果 resultDF.show()
运行输出完全匹配预期:
+----------+------------------+ | stat_date|active_flight_cnt | +----------+------------------+ |2022-05-08| 1| |2022-05-09| 1| |2022-05-10| 2| +----------+------------------+
内容的提问来源于stack exchange,提问作者PSD
相关产品推荐
相关产品推荐

