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_datecannot 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
相关产品推荐
相关产品推荐

