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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.01 02:21:21