Spark 3.x系列写到这里,基础篇的RDD和核心算子已经聊完了,后台问得最多的一个问题就是:“我已经会写SQL了,还有必要专门学Spark SQL吗?”我的回答一直很直接:如果只打算在Spark里学一样东西,那一定是Spark SQL。原因很简单,自Spark 2.x之后,DataFrame和SQL就是官方主推的数据处理入口,RDD那套在大部分业务场景里已经被SQL和内置算子替代了。尤其是做数据仓库、ETL、报表分析的同学,日常90%的活儿用Spark SQL都能搞定,而且代码量少、可读性强、优化空间大。
这篇是系列第三篇,不打算把官方文档搬一遍,重点聊三件事:第一,Spark 3.x的SQL引擎到底比RDD时代强在哪,DataFrame API和SQL怎么选;第二,Spark 3.x带来的几个关键变化,比如ANSI模式和AQE;第三,高频但容易翻车的日期时间处理细节,特别是很多人都在搜的“spark sql 日期加年”“日期加减”“日期转月份”这几个场景,我会把正确写法和坑一次说清楚。适合刚接触Spark准备做离线数仓的同学,也适合写过一阵子Spark SQL但老在日期和分区上踩坑的工程师。
1. Spark SQL为什么值得单独写一篇
1.1 从RDD到DataFrame,Spark SQL解决的三个核心问题
先说个背景。早期Spark基本都是RDD API,但RDD有几个很实际的问题。第一,RDD存的是JVM对象,Python里来回转pickle,序列化开销感人;第二,RDD没有schema概念,Spark完全不知道你这个RDD里有几个字段、字段什么类型,想优化都没法下手;第三,也是最重要的,RDD API写数据分析逻辑真的很啰嗦。一个简单的汇总统计,用RDD写要map、reduceByKey加各种函数组合,换成Spark SQL就是几行SELECT的事。
DataFrame和DataSet的引入,本质上是把“数据描述”这件事做了升级。DataFrame带schema,Spark能知道每一列的类型,于是可以做列裁剪、谓词下推、常量折叠这些优化。比如你查一个100列的表,SELECT只取其中3列,Catalyst优化器会把不需要的列直接裁掉,底层Scan阶段就不会读多余的数据。这个能力RDD时代想都不要想。
还有一个很多人忽略的点:Spark SQL不只支持SQL语法,它和DataFrame API底层走的是同一套优化流程。你写的DataFrame算子会被翻译成逻辑计划,进Catalyst优化器,再生成物理计划执行。所以Spark SQL的优点不仅仅是“支持SQL”,而是“SQL和程序API共用一套优化引擎”,这也是它在2.x之后能稳定替代RDD做主流开发入口的根本原因。
1.2 DataFrame、Dataset和SQL,平时到底用哪个
很多新手会纠结这个问题,我的答案比较务实:看场景,不要只看性能。
先理清概念。DataSet是强类型API,DataFrame其实就是DataSet[Row]的别名,所以在Scala和Java里DataFrame是弱类型的DataSet,Python里只有DataFrame可用。SQL则是直接写字符串。
从执行角度看,三种方式性能差距不大,因为最终都会进Catalyst。真正的区别在设计意图上:
- 做即席查询、报表、离线ETL,直接用SQL,短小精悍,非开发人员也能看懂。
- 做复杂的程序化数据处理,比如多步流程、逻辑分支、循环处理,用DataFrame API更合适,因为代码写起来更结构化,也方便做单元测试。
- 追求编译期类型安全,用强类型DataSet,但这种场景其实不多,大多数数仓工程师都不会用到。
我实际带项目的习惯是:能写SQL的先写SQL,等SQL变得太长或者需要程序逻辑介入时,再换DataFrame API。把两套东西结合好,比单纯选哪一套重要得多。因为Spark SQL里的内置函数和DataFrame API是互通的,写SQL时积累的日期函数、聚合函数经验,切到DataFrame API照样用。
2. Spark 3.x给SQL引擎带来的关键变化
2.1 ANSI模式:严格SQL到底要不要开
Spark 3.0一个很重要的变化就是引入了ANSI SQL模式,通过spark.sql.ansi.enabled参数控制,默认是false。这个参数影响很大,我建议所有用Spark 3.x的人都要了解一下。
开启ANSI模式后,很多原来“容忍错误”的行为会变成直接报错。举几个最常碰到的:
- 整数除以0,默认行为是返回null,ANSI下直接抛异常。
- 类型转换失败,比如把字符串'abc'转成int,默认返回null,ANSI下直接报错。
- 日期时间非法值,比如to_date('2024-13-45'),默认返回null,ANSI下会拒绝执行。
听起来开ANSI会影响稳定性,但它有一个很大的好处:把数据问题提前暴露。默认模式下,脏数据被静默转成null,你的聚合结果悄悄变少,排查起来非常费劲。ANSI模式下作业直接失败,你立刻知道数据有问题,这是“fail fast”的典型思路。
我的建议是:在开发环境开启ANSI模式跑一遍,把所有因为脏数据导致的作业失败处理掉,要么清洗数据,要么写CASE WHEN做防御。等数据质量稳定了,再决定生产环境是否开启。如果数据来源不可控又没人清洗,那还是先别开,不然每天凌晨任务报警能把你烦死。
2.2 AQE和动态分区裁剪,3.x性能优化绕不开的两个能力
Spark 3.x性能上最值得关注的两个特性,一个是自适应查询执行(AQE),一个是动态分区裁剪(DPP)。
AQE在Spark 3.2之后默认开启,核心思路是让Spark根据运行时统计信息动态调整执行计划。它的能力集中在三块。
一是自动合并Shuffle分区。很多时候你设置了200个shuffle分区,但数据量就那么点,每个分区没多少数据,白白浪费资源。AQE会根据shuffle输出大小,自动把分区数降下来,减少下游task数量,作业整体能快不少。
二是自动切换Join策略。比如某个join的表实际数据量很小,AQE会在运行时把SortMergeJoin换成BroadcastHashJoin,避免一次大shuffle。这个优化在数据量预估不准的时候特别有用。
三是处理数据倾斜。AQE能自动检测倾斜分区,把大分区拆成多个小分区去join,避免某个task拖垮整个作业。这个能力在事实表join维表时非常实用。
动态分区裁剪是另一个我特别偏爱的优化。它解决的是这样一个场景:一张大事实表join一张维表,维表上有过滤条件,理论上这个过滤条件能推到大事实表的Scan阶段,提前裁掉不需要的分区,但常规优化器做不到。DPP在运行时从维表侧获取过滤结果,动态生成分区裁剪条件,下推到事实表Scan端。效果就是事实表少读大量数据,作业时间呈数量级下降。
要确认你的作业是否生效,用EXPLAIN看执行计划即可。如果看到DynamicPartitionPruning相关字样,说明DPP生效了。实际使用中,DPP对非广播join、分区表场景效果最明显。
3. 日期时间函数实操,按热搜词逐个拆解
3.1 日期加减:date_add、date_sub和add_months的边界行为
先来解决“spark sql 日期加减”这个高频需求。分布式数据开发里,日期加减大概是最常见的操作,没有之一。比如T-1跑批、滚动近30天订单、计算同环比,全部离不开日期运算。
Spark SQL里最常用的三个函数:
-- 加N天 SELECT date_add(date'2024-01-15', 7); -- 2024-01-22 -- 减N天,等价于date_add(..., -7) SELECT date_sub(date'2024-01-15', 7); -- 2024-01-08 -- 加N个月 SELECT add_months(date'2024-01-31', 1); -- 2024-02-29date_add和date_sub的参数很简单,一个日期,一个整数,天为单位。这里我特别想提醒的是add_months的边界处理。比如1月31日加1个月,结果不是2月31日,而是2月29日(2024年是闰年),Spark会做“月末收缩”。这个行为在数仓月末跑批、月底结算的场景里特别重要,如果你用date_add(date'2024-01-31', 31)去模拟加一个月,结果大概率是错的,因为不同月份天数不一样。
还要注意add_months反向操作也是按月收缩。比如add_months(date'2024-03-31', -1)的结果是2024-02-29,不是2024-02-31。这个细节在算上月末、上上月末时尤其容易出错,建议在代码注释里写清楚。
3.2 日期加年的正确写法,别再用date_add加365天
现在说“spark sql 日期加年”。这个我一直认为是文档没写透的东西,因为没有直接的add_years函数,导致很多人直接写date_add(日期, 365),短时间看着没问题,跨个闰年就翻车。
为什么不推荐date_add(..., 365)?因为一年不总是365天。闰年有366天。看个例子:
-- 2020年是闰年,从1月1日到2021年1月1日中间隔了366天 SELECT date_add(date'2020-01-01', 365); -- 2020-12-31,错了 SELECT add_months(date'2020-01-01', 12); -- 2021-01-01,正确所以日期加年应该统一用add_months,把年数乘12:
-- 加1年 SELECT add_months(order_date, 12); -- 加3年 SELECT add_months(order_date, 36); -- 减2年 SELECT add_months(order_date, -24);这个写法的另一个好处是可以正确处理2月29日这种边界日期。比如add_months(date'2024-02-29', 12)的结果是2025-02-28,Spark自动做了月末收缩,符合大多数业务对“一年后”的预期。
3.3 日期转月份的三种做法,对应三种不同场景
“spark sql 日期转存月份”这个热搜词,翻译成白话就是:把date/timestamp字段变成月份,用于按月聚合、按月分区、或者生成一个月份维度字段。这个需求在实际数仓开发里出现频率极高,常见的做法有三种,每种对应不同场景。
方法一,格式化字符串:
SELECT date_format(order_date, 'yyyy-MM') AS month_str FROM orders;返回类型是string,形如“2024-01”。适合做报表展示、按字符串分组。也是很多人第一直觉会想到的方式。但要注意,date_format底层走Java的DateTimeFormatter,在大表上做group by时性能相对一般。
方法二,截断到月初:
SELECT trunc(order_date, 'month') AS month_start FROM orders;trunc返回date类型,结果是2024-01-01这种月初日期。这个方式我强烈推荐在按月分区表里用,因为分区字段的数据类型是date,读取和剪裁都更高效。注意trunc的第二个参数是小写'month',写成'MM'或者'mon'都会报错或者得 不到预期结果。
方法三,生成月份整数:
SELECT year(order_date) * 100 + month(order_date) AS month_int FROM orders;结果形如202401,是int类型。适合做维度表的关联键,比如月度维表、考核月字段,效率和可读性都挺好。
从性能角度,如果在10亿级的大表上按月做聚合,我的排序是:trunc > year*100+month > date_format。trunc的计算路径更短,而且返回date类型对后续日期运算更友好。date_format也不是不能用,只是没必要在核心链路里承担太多格式化开销。
3.4 字符串、时间戳和日期互转,容易忽略的隐藏行为
日期时间处理的另一半是类型转换。实际开发里,数据从业务库同步过来,很多时间字段都是字符串,比如“2024-01-15 10:30:00”。这时候就要用到to_date、to_timestamp、date_format、unix_timestamp和from_unixtime这一组函数。
-- 字符串转date,默认格式yyyy-MM-dd,可指定格式 SELECT to_date('2024-01-15 10:30:00'); -- 2024-01-15 SELECT to_date('2024/01/15', 'yyyy/MM/dd'); -- 2024-01-15 -- 字符串转timestamp SELECT to_timestamp('2024-01-15 10:30:00', 'yyyy-MM-dd HH:mm:ss'); -- timestamp/date转字符串 SELECT date_format(timestamp'2024-01-15 10:30:00', 'yyyy-MM-dd HH:mm'); -- 字符串转unix秒级时间戳 SELECT unix_timestamp('2024-01-15 10:30:00', 'yyyy-MM-dd HH:mm:ss'); -- 时间戳转字符串,注意from_unixtime返回的是string,不是timestamp SELECT from_unixtime(1705314600, 'yyyy-MM-dd HH:mm:ss');这里有一个我在Spark 3.x踩过的坑:from_unixtime的返回值类型是string,不是timestamp。如果你拿它直接跟timestamp字段比较,Spark会做隐式转换,看起来结果对,但一旦时间精度需要到毫秒、或者涉及时区换算,结果就会变得诡异。更好的做法是用timestamp_seconds(1705314600)得到timestamp类型,后续操作不会出幺蛾子。
如果你用的Spark是3.4及以上版本,还有个新类型TIMESTAMP_NTZ,全称Timestamp without time zone。它的特点是不受session时区影响,适合存业务本地时间、排班表这类语义上本来就不带时区的时间。但要注意,TIMESTAMP_NTZ目前在一些外部数据源连接器里的支持还不完整,引入前先在目标环境验证一遍。
4. 日期处理的老坑与排查技巧实录
4.1 经典坑:yyyy和YYYY不是一回事
这个坑我几乎每隔一段时间就会在别人的代码里看到一次。在Java和Spark的日期格式串里,小写yyyy表示“日历年”,也就是我们通常理解的年份;大写YYYY表示“ISO周历年”,它跟所在周的年份是对齐的,不一定等于日历年。
每年跨年那几天就是中招的窗口期。比如2021年1月1日,用date_format格式化:
SELECT date_format(date'2021-01-01', 'yyyy-MM-dd'); -- 2021-01-01 SELECT date_format(date'2021-01-01', 'YYYY-MM-dd'); -- 2020-01-01我们看第二个结果,周历年取成了2020,再加上MM和dd还是01-01,整串日期就从2021-01-01变成了2020-01-01,如果拿这个字符串去关联分区,数据直接错位。
解决方案很简单:所有日期格式化统一用小写yyyy,不要用大写YYYY。如果你在写分组逻辑时按日期分组,建议养成习惯,写完格式化后拿跨年的日期做一次单测,这是成本最低的验证方式。
4.2 函数包裹分区字段,导致分区裁剪失效
这个属于性能排查里频率极高的问题。很多人写SQL时为了图方便,在WHERE条件里对分区字段做函数运算:
-- 反面教材:对分区字段用date_format,分区裁剪失效,全表扫描 SELECT * FROM orders WHERE date_format(order_date, 'yyyy-MM') = '2024-01'; -- 正确写法:用范围条件,分区裁剪生效 SELECT * FROM orders WHERE order_date >= date'2024-01-01' AND order_date < date'2024-02-01';原理很好理解:优化器在扫描阶段需要根据原始字段值决定读哪些文件,一旦字段被函数包了一层,它就无法在Scan阶段推算原值范围,只能全量读进来再过滤。这就是谓词下推失效。
所以写SQL的第一原则是:分区字段保持裸写,不要套函数。如果确实要按月过滤,用月初和月末范围;如果分区字段是月初日期,直接WHERE month_start = trunc('2024-01-15', 'month')也行,等号右边有函数没关系,左边干净就行。
4.3 时区不一致,日期函数结果对不上
还有一个非常隐蔽的坑是时区。Spark的timestamp在底层存储上是绝对时间戳,但展示和格式化时会按spark.sql.session.timeZone换算。如果你的Spark作业没设置过时区,默认跟集群环境走,而集群环境往往是UTC。
举个例子,上游同步过来的订单时间存成timestamp'2024-01-01 00:00:00'(UTC),业务系统时区是UTC+8。你在开发环境跑:
SET spark.sql.session.timeZone=Asia/Shanghai; SELECT date_format(timestamp'2024-01-01 00:00:00', 'yyyy-MM-dd'); -- 2024-01-01 SET spark.sql.session.timeZone=UTC; SELECT date_format(timestamp'2024-01-01 00:00:00', 'yyyy-MM-dd'); -- 2024-01-01这个例子看起来一样,但如果你拿到的timestamp是'2024-01-01 00:00:00'且它本身已经是UTC时间,在Asia/Shanghai时区下转成日期,结果可能变成2024-01-01没错,还是2023-12-31取决于具体值。真实业务里最常出现的情况是:同一条时间数据在不同任务里因为session timezone不同,计算出的日期不同,导致日分区数据串位。
我的建议是:所有时间字段统一口径。数仓内部时间字段统一用时间戳或统一时区的date字段;跨时区业务数据,在入仓时就转成目标时区的时间;每个作业里显式设置spark.sql.session.timeZone,不要依赖集群默认值。
4.4 月份运算在月末跑批中的实战注意点
最后说一个数仓里非常常见的场景:按月跑批,计算上一月、上月月初、上月月末这些边界。
-- 当前日期 SELECT current_date(); -- 比如2024-03-31 -- 上月末:当月1日减1天 SELECT date_sub(trunc(current_date(), 'month'), 1); -- 2024-02-29 -- 上月初:加-1个月,再取月初 SELECT trunc(add_months(current_date(), -1), 'month'); -- 2024-02-01 -- 本月月初 SELECT trunc(current_date(), 'month'); -- 2024-03-01这段SQL看起来简单,但坑在时序。比如你在3月31日跑批,current_date()是2024-03-31,上月末应该是2024-02-29,用date_sub(trunc(...), 1)没问题。如果某天任务从3月1日凌晨0点跑,此时current_date是2024-03-01,那上月末就是2024-02-29,逻辑依然成立。问题出在有人喜欢用“今天减天数”来算月份,比如date_add(current_date(), -30),跨大月小月就会偏。
另一个注意点是只算月份差时,不要直接用months_between然后取整:
-- 这种写法在月末场景容易踩坑 SELECT months_between(date'2024-02-28', date'2024-01-31'); -- 结果是0.9677...,不是1months_between返回的是基于每月天数的浮点差,不是简单的整数月份差。想算整月数,建议先trunc到月初再算:
SELECT months_between(trunc(date'2024-02-28', 'month'), trunc(date'2024-01-31', 'month')); -- 结果正好是1.0这个细节能帮你避免很多月末跑批的月末数据对不上的问题。
5. 常用日期操作速查表与个人建议
5.1 高频操作速查表
把前面提到的内容整理成一张速查表,方便直接当工具用。
| 需求 | 推荐写法 | 返回类型 | 注意事项 |
|---|---|---|---|
| 日期加N天 | date_add(dt, N) | date | N可为负数,等价于减 |
| 日期减N天 | date_sub(dt, N) | date | N可为负数 |
| 日期加N月 | add_months(dt, N) | date | 月末自动收缩 |
| 日期加N年 | add_months(dt, 12*N) | date | 别用date_add(dt, 365*N) |
| 日期减N年 | add_months(dt, -12*N) | date | 跨闰年也能正确处理 |
| 两个日期差多少天 | datediff(d1, d2) | int | d1 - d2 |
| 两个日期差多少月 | months_between(d1, d2) | double | 先trunc到月初再算 |
| 取月初 | trunc(dt, 'month') | date | 参数是字符串'month' |
| 取月末 | last_day(dt) | date | 返回该月最后一天 |
| 日期转月份字符串 | date_format(dt, 'yyyy-MM') | string | 用yyyy,别用YYYY |
| 日期转月份月初日期 | trunc(dt, 'month') | date | 适合做分区字段 |
| 日期转月份整数 | year(dt)*100+month(dt) | int | 适合关联维度表 |
| 字符串转日期 | to_date(s, 'yyyy-MM-dd') | date | 格式串需匹配 |
| 时间戳转日期date | to_date(timestamp_col) | date | 注意时区影响 |
| 时间戳转字符串 | date_format(ts, 'yyyy-MM-dd HH:mm:ss') | string | 受session时区影响 |
这张表是这么多年来攒下来的,基本覆盖了数据开发日常90%的日期处理需求。可以放进团队Wiki里当规范用。
5.2 给新手的几条实用建议
最后说几点实操建议,都是带团队时反复强调的东西。
第一,时间字段能date就别timestamp。如果业务上不需要时分秒,尽量用date类型,存储更省、比较更直观、不容易踩时区坑。需要精确到时分秒的场景单独用timestamp。
第二,分区字段的命名和类型要统一。我强烈建议按月分区表的分区字段直接用date类型(月初日期),不要用string类型存“2024-01”这种,因为date类型支持范围比较和算术运算,string类型很容易写出错误的时间过滤条件。
第三,写SQL之前先EXPLAIN。不是让你天天看执行计划,而是在新写的SQL或者SQL性能变慢时,EXPLAIN一眼就能看出分区裁剪有没有生效、DPP有没有生效、有没有全表扫描。这个习惯比任何调参都管用。
第四,日期格式化串统一用小写。yyyy、MM、dd、HH、mm、ss,在Spark 3.x里日期格式串大小写敏感,特别要记住YYYY这个坑,大写跟小写行为完全不一样。
结尾
写了这么多,想起之前带团队做数仓项目时,很多同学SQL基础挺扎实,但一到日期计算就各种玄学bug——跨年错一天、月末多一天、时区差八小时。其实Spark SQL的日期函数设计得很完善,只要掌握date_add、add_months、trunc、last_day这几个核心函数,再加上对边界情况的敬畏,90%的问题都能避免。
我个人在实际操作中的体会是:日期处理的bug往往不是函数不会用,而是对边界行为没有预期。所以我现在每写一个日期相关的SQL,都会习惯性先拿几个边界日期去验证再上线。这篇文章如果能在你下次写月份计算或日期转换时少踩一个坑,这篇就值了。Spark 3.x系列后面我准备接着聊Spark SQL里的窗口函数和复杂数据类型,继续补全这套实战拼图。