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

Spark SQL按年月分组统计发布数及累计数报错排查

问题:PySpark SQL按月统计发布数并计算累计值

我有一张名为release_dates的表,数据如下:

+------------+
|release_date|
+------------+
|  1997-06-30|
|  1997-11-14|
|  1998-11-08|
|  1999-04-01|
|  1999-09-08|
|  1999-11-01|
|  2000-11-01|
|  2000-11-01|
|  2001-03-15|
|  2001-06-01|
|  2001-06-01|
+------------+

需要按月统计发布数量,同时计算相同粒度的累计发布数,期望结果如下:

+----------+-----------+------------+
|      date|nb_releases|cum_releases|
+----------+-----------+------------+
|1997-06-01|          1|           1|
|1997-11-01|          1|           2|
|1998-11-01|          1|           3|
|1999-04-01|          1|           4|
|1999-09-01|          1|           5|
|1999-11-01|          1|           6|
|2000-11-01|          2|           8|
|2001-03-01|          1|           9|
|2001-06-01|          2|          11|
+----------+-----------+------------+

尝试的SQL语句:

SELECT
  any_value(release_date) - DAY(any_value(release_date)) + 1 AS date,
  COUNT (*) AS nb_releases,
  SUM (COUNT(*)) OVER (ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS cum_releases
FROM 
  release_dates
WHERE
  release_date IS NOT NULL
GROUP BY
  YEAR (release_date),
  MONTH (release_date)
ORDER BY
  YEAR (release_date),
  MONTH (release_date)

执行后报错:

A column or function parameter with name release_date cannot be resolved

使用PySpark 3.5.3版本,希望仅通过SQL查询解决问题,不使用PySpark API。


错误原因

原SQL中,release_date未包含在GROUP BY的分组字段里,PySpark SQL的解析规则不允许在SELECT子句中直接引用未分组的原始字段(即便用any_value包裹,也会因为解析顺序问题触发错误)。

正确解决方案

方法1:用date_trunc生成月份起始日期(推荐)

利用PySpark内置的date_trunc函数直接将日期截断到月份,自动生成当月第一天,写法更简洁:

SELECT
  date_trunc('month', release_date) AS date,
  COUNT(*) AS nb_releases,
  SUM(COUNT(*)) OVER (ORDER BY date) AS cum_releases
FROM release_dates
WHERE release_date IS NOT NULL
GROUP BY date_trunc('month', release_date)
ORDER BY date;

方法2:先分组聚合,再计算累计值

先按年、月分组统计每月发布数,再基于分组结果构造月份起始日期并计算累计值:

WITH monthly_releases AS (
  SELECT
    YEAR(release_date) AS release_year,
    MONTH(release_date) AS release_month,
    COUNT(*) AS nb_releases
  FROM release_dates
  WHERE release_date IS NOT NULL
  GROUP BY YEAR(release_date), MONTH(release_date)
)
SELECT
  DATE_FROM_UNIXTIME(UNIX_TIMESTAMP(CONCAT(release_year, '-', release_month, '-01'), 'yyyy-MM-dd')) AS date,
  nb_releases,
  SUM(nb_releases) OVER (ORDER BY release_year, release_month) AS cum_releases
FROM monthly_releases
ORDER BY date;

补充说明

  • 窗口函数中ORDER BY date会自动按日期顺序累计,无需额外指定ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW(这是窗口函数的默认行为)。
  • date_trunc返回的是日期类型,完全符合期望结果中的格式要求。

内容的提问来源于stack exchange,提问作者johnc

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 07:04:55