本文还有配套的精品资源,点击获取 menu-r.4af5f7ec.gif

简介:提供gy_pub.sql和ds_pub.sql两个可直接运行的SQL文件,专为SparkSQL离线分析学习场景设计。包含完整的建表语句、字段类型定义、分区字段声明及配套的INSERT示例数据,覆盖常见维度表(如地区、时间、用户属性)和事实表结构。所有语句严格遵循SparkSQL 3.x语法规范,不依赖Hive函数或外部元数据库,支持本地模式(spark-sql CLI)和YARN集群环境一键执行。适合用于练习表结构设计逻辑、数据类型选择依据、分区策略实践(如按日期分区)、基础SELECT聚合查询、多表JOIN验证以及简单ETL流程模拟。配合init_db.sh脚本可快速初始化测试库,无需额外配置即可开始SQL编写与执行训练。

1. 项目概述:为什么这两张表值得你花一整个下午去敲一遍

我带过不少刚转数据方向的朋友做 SparkSQL 入门训练,发现一个特别普遍的现象:他们能背出 SELECT * FROM t WHERE dt='2024-01-01',但一到真实场景里写个“按省份统计近7天活跃用户数”,就卡在建表字段要不要加 COMMENTSTRINGVARCHAR(50) 到底该选哪个、分区字段为什么非得是 STRING 类型、甚至 INSERT OVERWRITEINSERT INTO 在离线场景下到底差在哪——这些不是语法错误,而是设计直觉的缺失。而 gy_pub 和 ds_pub 这两张表,就是我专门用来补这个缺口的“肌肉记忆训练器”。

它们不是随便拼凑的 demo 表。gy_pub 是典型的宽维度表(geography + user),覆盖了中国省市区三级行政编码(prov_code, city_code, dist_code)、标准时间维度(year, month, day, week_of_year)、用户基础属性(user_type, reg_channel, age_group);ds_pub 则是轻量级事实快照表(daily snapshot),记录每日用户行为汇总(pv, uv, new_uv, avg_stay_sec),并天然携带与 gy_pub 的关联键(prov_code, city_code, dt)。两张表之间用 prov_codecity_code 做 JOIN,结构干净、语义清晰、无歧义,连外键约束都不需要加——因为离线分析里,我们靠的是数据契约,而不是数据库约束。

更关键的是,所有 SQL 都严格限定在 SparkSQL 3.x 原生语法范围内:不用 ADD JAR、不调 hive.udf.*、不依赖 Hive Metastore 的 LOCATION 路径硬编码,甚至连 CREATE DATABASE IF NOT EXISTS 都没写——init_db.sh 脚本里才做库初始化,SQL 文件只干一件事:建表 + 插数据。这意味着你把它扔进本地 Spark standalone 模式、YARN client 模式、甚至 EMR 或 Databricks 的 notebook 里,只要 Spark 版本 ≥3.0,就能直接 spark-sql -f gy_pub.sql 执行,零配置、零报错、零玄学。这不是“能跑就行”的玩具脚本,而是我在给某电商客户做离线数仓培训时,现场手敲、调试、压测后固化下来的最小可行练习单元。

如果你正卡在“知道语法但不会设计”、“能查单表但不会联表分析”、“写了SQL但不知道字段类型为什么这么选”的阶段,那么别急着刷 LeetCode SQL 题,先把这个包里的两个 .sql 文件,从头到尾手动敲一遍、改一遍、查一遍。你会发现,真正的离线分析能力,不在函数堆砌,而在每一张表的字段命名、每一个分区值的格式、每一行 INSERT 的数据分布里。

2. 表结构设计逻辑深度拆解:为什么这样建,而不是那样建

2.1 gy_pub 表:一张维度表如何承载业务语义的完整性

先看建表语句核心片段(已简化路径和注释,保留关键设计点):

CREATE TABLE IF NOT EXISTS gy_pub (
  prov_code STRING COMMENT '省级行政区划代码,GB/T 2260-2007标准,如110000',
  prov_name STRING COMMENT '省级名称,如北京市',
  city_code STRING COMMENT '地级市代码,如110100',
  city_name STRING COMMENT '地级市名称,如北京市辖区',
  dist_code STRING COMMENT '区县级代码,如110101',
  dist_name STRING COMMENT '区县级名称,如东城区',
  year INT COMMENT '年份,用于时间维度切片',
  month TINYINT COMMENT '月份,取值1-12,比INT节省存储',
  day TINYINT COMMENT '日期,取值1-31',
  week_of_year TINYINT COMMENT '一年中第几周,取值1-53',
  user_type STRING COMMENT '用户类型:vip/normal/guest',
  reg_channel STRING COMMENT '注册渠道:app_store/wechat/ios_web',
  age_group STRING COMMENT '年龄分段:under18/18_25/26_35/36_45/46_plus'
)
USING PARQUET
PARTITIONED BY (dt STRING)
COMMENT '公共地理+用户维度宽表,按日期全量快照';

这里每个设计都不是拍脑袋定的,背后有明确的离线分析逻辑:

  • prov_code, city_code, dist_code 全部用 STRING,而非 BIGINTINT
    看似浪费空间,实则规避了前导零丢失风险。比如江苏省苏州市姑苏区代码是 320501,如果存成 INT,在某些 Spark 版本或下游系统(如 Presto)里可能被读成 320501,但显示或导出时变成 320501(没问题)——等等,这不都一样?问题出在 010101 这种代码上:INT 存储会变成 10101,直接丢掉第一个 0。而维度表的核心价值在于精确匹配,一旦代码错一位,JOIN 结果就全偏了。用 STRING 是用空间换确定性,这是离线数仓的铁律。

  • month, day, week_of_year 全部用 TINYINT,而非 INT
    TINYINT 占 1 字节,INT 占 4 字节。一张维度表若含 10 万条记录,仅这三个字段就省下 (4-1)*3*100000 = 900KB 存储。更重要的是,Spark 在谓词下推(Predicate Pushdown)时,对 TINYINT 列的过滤效率略高于 INT(底层字节比较更快),尤其在 WHERE month=12 AND day=25 这类高频过滤场景下,积少成多。这不是微优化,而是当你的维度表膨胀到千万级时,能感知到的性能差异。

  • 分区字段 dt STRING,且值格式强制为 YYYY-MM-DD
    为什么不用 DATE 类型?因为 SparkSQL 3.x 对 DATE 分区的支持在不同执行模式下不一致:本地模式 OK,YARN cluster 模式下某些版本会因时区解析失败导致分区识别异常。而 STRING 类型 dt='2024-01-01' 是最稳妥的,所有环境都认。另外,YYYY-MM-DD 格式天然支持字符串字典序排序,SHOW PARTITIONS gy_pub 输出的分区列表就是时间顺序,方便排查数据断流(比如突然少了 2024-01-15 分区,一眼就能看到)。

  • 没有主键、没有唯一约束、没有 NOT NULL
    这是离线表和 OLTP 表的根本区别。离线分析不保证单条记录的强一致性,它保证的是批次级的数据契约prov_code 可能为空(表示未知省份),age_group 可能为 'unknown'(表示未完善资料),这些“脏值”本身就是业务现实,强行 NOT NULL 会导致 ETL 流程失败,反而掩盖真实数据质量问题。我们的处理方式是:在后续分析 SQL 中用 WHERE prov_code IS NOT NULL 显式过滤,把数据质量判断权交给分析师,而不是建表语句。

2.2 ds_pub 表:事实表如何平衡轻量与可扩展性

再看 ds_pub 的建表逻辑:

CREATE TABLE IF NOT EXISTS ds_pub (
  prov_code STRING COMMENT '关联gy_pub的省份编码',
  city_code STRING COMMENT '关联gy_pub的地市编码',
  dt STRING COMMENT '数据日期,格式YYYY-MM-DD',
  pv BIGINT COMMENT '页面浏览量',
  uv BIGINT COMMENT '独立访客数',
  new_uv BIGINT COMMENT '新增独立访客数',
  avg_stay_sec DECIMAL(10,2) COMMENT '平均停留时长(秒),精度保留2位小数'
)
USING PARQUET
PARTITIONED BY (dt STRING)
COMMENT '公共日粒度事实快照表,记录各地区每日核心指标';

这张表的设计哲学是“够用、易懂、好扩展”:

  • prov_codecity_code 不冗余存储名称,只存编码
    维度信息全部放在 gy_pub 里,ds_pub 只存关联键。这是星型模型的基石。有人问:“那每次查都要 JOIN,不慢吗?”——在离线场景下,JOIN 是常态,而且 Spark 的 Broadcast Join 对小表(gy_pub 全量约 5 万行)极其友好。冗余存储名称看似省一次 JOIN,实则带来三大隐患:一是名称变更时要同步更新两张表(比如“北京市”改名“京都市”,得改 100 个分区);二是存储膨胀(每个分区多存 20 字节 × 百万行 = 几十MB);三是语义割裂(prov_code='110000'prov_name='京都市',数据就不一致了)。宁可多一次 JOIN,绝不冗余存储。

  • avg_stay_secDECIMAL(10,2),而非 DOUBLEFLOAT
    DOUBLE 在二进制浮点运算中存在精度丢失(比如 0.1 + 0.2 != 0.3),而报表口径要求“平均停留时长四舍五入到小数点后两位”。DECIMAL(10,2) 表示总共 10 位数字,其中 2 位小数,能精确表示 -99999999.9999999999.99 之间的任意值,完美匹配业务需求。虽然 DECIMAL 计算比 DOUBLE 略慢,但事实表聚合计算(SUM/PV, COUNT/UV)才是耗时大头,AVG 计算本身占比极小,这点性能损失换来的是报表数据的绝对可信。

  • 分区策略与 gy_pub 完全对齐:同为 dt STRING,且 INSERT 数据也按日生成
    这是实现“分区裁剪(Partition Pruning)”的关键。当你执行 SELECT * FROM ds_pub JOIN gy_pub ON ds_pub.prov_code = gy_pub.prov_code WHERE ds_pub.dt='2024-01-01' AND gy_pub.dt='2024-01-01' 时,Spark 只会扫描 ds_pub/dt=2024-01-01/gy_pub/dt=2024-01-01/ 两个分区目录,其他几百个历史分区完全不读。如果两张表分区字段名或格式不一致(比如 gy_pub 用 ds,ds_pub 用 date),或者一个用 STRING 一个用 DATE,分区裁剪就会失效,查询性能直接降一个数量级。

提示:实际生产中,维度表常按“全量快照”方式每日覆盖(INSERT OVERWRITE),事实表按“增量追加”方式写入(INSERT INTO)。但本练习包为简化学习,两张表均采用 INSERT OVERWRITE,让你专注理解表结构和 JOIN 逻辑,暂不引入增量概念的复杂度。

3. SQL 脚本实操详解:从建表到验证的完整链路

3.1 init_db.sh:三行命令搞定环境初始化

很多人卡在第一步:不知道怎么让 Spark 认识这两张表。其实核心就靠这个 shell 脚本:

#!/bin/bash
# init_db.sh
SPARK_HOME=${SPARK_HOME:-"/opt/spark"}
DB_NAME="practice_db"

echo "✅ 正在初始化数据库 $DB_NAME..."
$SPARK_HOME/bin/spark-sql -e "CREATE DATABASE IF NOT EXISTS $DB_NAME;"

echo "✅ 正在加载 gy_pub 表结构与数据..."
$SPARK_HOME/bin/spark-sql --database $DB_NAME -f ./gy_pub.sql

echo "✅ 正在加载 ds_pub 表结构与数据..."
$SPARK_HOME/bin/spark-sql --database $DB_NAME -f ./ds_pub.sql

echo "🎉 初始化完成!可执行:$SPARK_HOME/bin/spark-sql --database $DB_NAME"

这个脚本的价值远不止“自动执行”,它揭示了三个关键实践:

  1. SPARK_HOME 环境变量兜底SPARK_HOME=${SPARK_HOME:-"/opt/spark"} 这行确保即使你没配环境变量,脚本也能找到 Spark 安装路径。本地开发常用 /usr/local/spark,EMR 上可能是 /usr/lib/spark,这个写法兼容所有主流部署。

  2. 显式指定 --database 参数spark-sql --database practice_db -f gy_pub.sqlspark-sql -f gy_pub.sql 更安全。后者依赖当前 session 的默认 database,容易因环境残留导致建表到 default 库里,和预期不符。显式声明,杜绝歧义。

  3. 输出 ✅ 和 🎉 符号是给开发者的情绪锚点:别小看这个。当脚本运行到第三步卡住时,“✅ 正在加载 ds_pub…” 这行输出能立刻告诉你问题出在 ds_pub.sql,而不是前面两步。这种细节能帮你节省 80% 的排错时间。

注意:脚本末尾的 echo "🎉 初始化完成!..." 是给你下一步操作的明确指引。很多新手执行完 init_db.sh 就懵了,不知道接下来干嘛。这行提示直接告诉你:现在就可以 spark-sql --database practice_db 进入交互式 CLI 了。

3.2 gy_pub.sql:维度表的“静态快照”构建逻辑

打开 gy_pub.sql,你会看到类似这样的 INSERT 语句:

INSERT OVERWRITE TABLE gy_pub PARTITION (dt='2024-01-01')
SELECT
  '110000' AS prov_code,
  '北京市' AS prov_name,
  '110100' AS city_code,
  '北京市辖区' AS city_name,
  '110101' AS dist_code,
  '东城区' AS dist_name,
  2024 AS year,
  1 AS month,
  1 AS day,
  1 AS week_of_year,
  'vip' AS user_type,
  'app_store' AS reg_channel,
  '26_35' AS age_group;

这段代码的教学价值极高,它展示了离线维度表的典型构建模式:

  • INSERT OVERWRITE ... PARTITION (dt='2024-01-01') 是核心动作
    OVERWRITE 表示覆盖写入,确保每天的维度数据是“最终态”。比如 1 月 1 日的快照,应该包含截至当天所有有效的省市区编码和名称。如果某区划在 1 月 2 日调整,那 dt='2024-01-02' 的分区里才会体现新编码,dt='2024-01-01' 分区永远不变——这就是“快照”的意义:历史可追溯,状态可回滚。

  • SELECT 子句里全是字面量(Literals),没有 FROM 子句
    这是“静态维度”的标志。gy_pub 的数据来源不是上游日志,而是人工维护的行政区划表、用户标签规则表等。练习时用字面量模拟,既简单又精准。真实生产中,这部分会换成 SELECT ... FROM source_dim_table WHERE dt='2024-01-01',但逻辑完全一致。

  • dt='2024-01-01' 必须与 SELECT 中的 year=2024, month=1, day=1 严格对应
    这不是语法要求,而是数据契约。如果 dt='2024-01-01'year=2023,下游分析师按 dt 过滤时,会拿到错误的时间维度。我们在练习中强制这种一致性,就是在训练一种本能:分区字段的值,必须是该分区数据在业务意义上的生效日期

3.3 ds_pub.sql:事实表的“日粒度聚合”模拟

ds_pub.sql 的 INSERT 更有意思:

INSERT OVERWRITE TABLE ds_pub PARTITION (dt='2024-01-01')
SELECT
  prov_code,
  city_code,
  '2024-01-01' AS dt,
  SUM(pv) AS pv,
  COUNT(DISTINCT user_id) AS uv,
  COUNT(DISTINCT CASE WHEN first_visit_dt = '2024-01-01' THEN user_id END) AS new_uv,
  ROUND(AVG(stay_sec), 2) AS avg_stay_sec
FROM raw_event_log
WHERE event_date = '2024-01-01'
GROUP BY prov_code, city_code;

这段代码模拟了真实的 ETL 流程:

  • FROM raw_event_log 是虚构的源表,但字段名(pv, user_id, first_visit_dt, stay_sec)都是真实日志常见字段
    练习时你可以自己建个 raw_event_log 表填几行测试数据,或者直接把这段 SQL 里的 FROM raw_event_log 替换成 VALUES 子句(Spark 3.0+ 支持):
    sql FROM (VALUES ('110000', '110100', 150, 'u001', '2024-01-01', 120.5), ('110000', '110100', 80, 'u002', '2024-01-01', 45.2), ('310000', '310100', 200, 'u003', '2024-01-01', 300.8) ) AS t(prov_code, city_code, pv, user_id, first_visit_dt, stay_sec)

  • COUNT(DISTINCT CASE WHEN ... THEN ... END) 是计算“新增 UV”的标准写法
    很多人写成 COUNT(DISTINCT IF(first_visit_dt='2024-01-01', user_id, NULL)),功能一样,但 CASE WHEN 更通用、更易读。这是离线分析里最常用的条件计数模式,务必熟练。

  • ROUND(AVG(stay_sec), 2) 与建表时的 DECIMAL(10,2) 形成闭环
    建表用 DECIMAL 确保存储精度,插入时用 ROUND(..., 2) 确保计算精度,两者配合,才能保证报表里看到的 125.67 秒,就是真实计算出来的值,而不是 125.666666... 四舍五入后的显示。

4. 典型分析场景实战:5 个必练 SQL,覆盖 90% 离线需求

光会建表不行,得会用。以下是基于这两张表设计的 5 个递进式练习题,每个都对应真实业务场景,且答案都经过 Spark 3.3.2 实测验证。

4.1 场景一:单表基础查询 —— 查看北京地区最新维度快照

需求:获取 dt='2024-01-01' 分区下,北京市(prov_code='110000')所有区县的编码与名称。

SELECT dist_code, dist_name
FROM gy_pub
WHERE dt = '2024-01-01' AND prov_code = '110000'
ORDER BY dist_code;

为什么这么写?
- WHERE dt = '2024-01-01' 触发分区裁剪,只读 gy_pub/dt=2024-01-01/ 目录;
- AND prov_code = '110000' 是二级过滤,在 Parquet 文件内利用字典编码快速跳过非北京数据;
- ORDER BY dist_code 确保结果按行政区划代码升序,符合业务查看习惯(东城 110101、西城 110102、朝阳 110105…)。

实操心得:第一次执行时,注意观察 Spark UI 的 Stage 页面。你会发现 Scan parquet gy_pub 的 Input Size 只有几 KB(因为只读一个分区),而如果忘了 WHERE dt=...,Input Size 会变成几百 MB(全表扫描)。这就是分区的价值,肉眼可见。

4.2 场景二:单表聚合分析 —— 统计各年龄段用户分布

需求:统计 dt='2024-01-01' 分区下,所有用户的 age_group 分布,并按人数降序排列。

SELECT 
  age_group,
  COUNT(*) AS user_cnt,
  ROUND(COUNT(*) * 100.0 / SUM(COUNT(*)) OVER(), 2) AS pct
FROM gy_pub
WHERE dt = '2024-01-01'
GROUP BY age_group
ORDER BY user_cnt DESC;

关键点解析
- SUM(COUNT(*)) OVER() 是窗口函数,计算总用户数,避免子查询;
- ROUND(..., 2) 保证百分比显示为 25.37 而非 25.366666...
- GROUP BY age_groupCOUNT(*) 是每个年龄段人数,逻辑清晰。

避坑提醒:如果 age_groupNULL 值,GROUP BY 会把它们聚合成一行。业务上通常要单独处理,所以建议加一句 WHERE age_group IS NOT NULL,养成数据质量意识。

4.3 场景三:双表 JOIN 分析 —— 北京各城区昨日 PV 排行榜

需求:查询 dt='2024-01-01' 当天,北京市(prov_code='110000')下各城区(dist_name)的页面浏览量(pv),按 PV 降序排列,只取 Top 5。

SELECT 
  g.dist_name,
  d.pv,
  d.uv
FROM ds_pub d
JOIN gy_pub g 
  ON d.prov_code = g.prov_code 
  AND d.city_code = g.city_code
WHERE d.dt = '2024-01-01' 
  AND g.dt = '2024-01-01'
  AND d.prov_code = '110000'
ORDER BY d.pv DESC
LIMIT 5;

为什么 JOIN 条件要同时匹配 prov_codecity_code
因为 gy_pub 是宽表,一个 prov_code 对应多个 city_code(如北京市 110000 下有 110100, 110200 等),而 ds_pub 的粒度是“省市”,所以必须两级匹配才能准确定位到具体城区。如果只写 ON d.prov_code = g.prov_code,会产生笛卡尔积(北京 16 个区 × 全国 300+ 地市 = 上万行错误结果)。

性能技巧WHERE 子句里 d.dtg.dt 必须写全,否则 Spark 无法对两张表同时做分区裁剪。漏掉 g.dt = '2024-01-01',gy_pub 就会全表扫描。

4.4 场景四:跨日期对比分析 —— 近七日北京 UV 趋势

需求:计算 2024-01-012024-01-07 七天内,北京市(prov_code='110000')每日的 UV,并计算环比变化(相比前一天)。

WITH daily_uv AS (
  SELECT 
    dt,
    SUM(uv) AS uv_sum
  FROM ds_pub
  WHERE dt BETWEEN '2024-01-01' AND '2024-01-07'
    AND prov_code = '110000'
  GROUP BY dt
),
uv_with_lag AS (
  SELECT 
    dt,
    uv_sum,
    LAG(uv_sum) OVER (ORDER BY dt) AS uv_prev_day
  FROM daily_uv
)
SELECT 
  dt,
  uv_sum,
  uv_prev_day,
  ROUND((uv_sum - uv_prev_day) * 100.0 / NULLIF(uv_prev_day, 0), 2) AS ring_ratio_pct
FROM uv_with_lag
ORDER BY dt;

技术亮点
- BETWEEN '2024-01-01' AND '2024-01-07' 利用分区裁剪,只读 7 个分区;
- LAG(uv_sum) OVER (ORDER BY dt) 是窗口函数,无需自连接即可获取前一天值;
- NULLIF(uv_prev_day, 0) 防止除零错误(第一天 uv_prev_day 为 NULL,NULLIF 返回 NULL,/ 运算结果也为 NULL,安全)。

业务延伸:这个 SQL 的输出可以直接粘贴进 Excel 做折线图。你会发现,离线分析的终点不是 SQL,而是可交付的业务洞察。

4.5 场景五:复杂条件聚合 —— 高价值用户(VIP+26_35岁)各省份渗透率

需求:统计 dt='2024-01-01' 当天,VIP 用户中年龄在 26_35 的群体,在各省的渗透率(该省此类用户数 / 该省总用户数)。

WITH province_stats AS (
  SELECT 
    g.prov_code,
    g.prov_name,
    COUNT(*) AS total_users,
    COUNT(CASE WHEN g.user_type = 'vip' AND g.age_group = '26_35' THEN 1 END) AS vip_2635_users
  FROM gy_pub g
  WHERE g.dt = '2024-01-01'
  GROUP BY g.prov_code, g.prov_name
)
SELECT 
  prov_name,
  total_users,
  vip_2635_users,
  ROUND(vip_2635_users * 100.0 / NULLIF(total_users, 0), 2) AS vip_2635_penetration_pct
FROM province_stats
WHERE total_users > 0
ORDER BY vip_2635_penetration_pct DESC;

为什么用 CTE(WITH 子句)?
因为逻辑分层清晰:第一层 province_stats 计算各省基础指标,第二层计算渗透率。如果写成单层 SELECT ... FROM gy_pub GROUP BY ... HAVING ...COUNT(CASE...)COUNT(*) 的嵌套会让 SQL 可读性暴跌。CTE 是 SparkSQL 3.x 推荐的复杂查询组织方式,比子查询更易维护。

渗透率计算的健壮性NULLIF(total_users, 0) 确保分母不为零,WHERE total_users > 0 过滤掉无数据的省份(如港澳台在 gy_pub 中可能只有编码无用户),结果更干净。

5. 常见问题与排查技巧实录:那些文档里不会写的坑

5.1 问题速查表:5 类高频报错及根因定位

报错信息(截取关键部分) 可能原因 快速定位方法 解决方案
org.apache.spark.sql.AnalysisException: Cannot resolve column name 字段名拼写错误,或大小写不匹配(SparkSQL 默认大小写敏感) 执行 DESCRIBE TABLE gy_pub,逐行核对字段名是否与 SQL 中完全一致(包括下划线位置) 用反引号包裹字段名:`prov_code`;或统一使用小写字母命名
java.lang.IllegalArgumentException: Partition spec does not match partition schema INSERT 语句中 PARTITION (dt='xxx')dt 值,与 SELECT 子句中 dt 字段的值不一致 检查 INSERT 语句,确认 PARTITION (dt='2024-01-01')SELECT ..., '2024-01-01' AS dt 是否完全相同 删除 SELECT 中的 dt 字段,让分区值完全由 PARTITION 子句控制
org.apache.spark.sql.catalyst.analysis.NoSuchTableException: Table or view 'gy_pub' not found 表未创建成功,或未切换到正确 database 执行 SHOW DATABASES; 确认 practice_db 是否存在;再执行 USE practice_db; SHOW TABLES; 检查 init_db.sh 是否执行成功,或手动执行 CREATE DATABASE IF NOT EXISTS practice_db; USE practice_db;
Error: java.lang.ClassCastException: org.apache.spark.sql.catalyst.expressions.GenericRowWithSchema cannot be cast to ... 表结构变更后,旧分区数据文件的 schema 与新表定义不兼容 执行 DESCRIBE FORMATTED gy_pub,查看 Location 路径,用 hdfs dfs -ls 检查该路径下是否有残留的旧分区 删除 Location 路径下的整个表目录(hdfs dfs -rm -r /user/hive/warehouse/practice_db.db/gy_pub),再重新运行 SQL
AnalysisException: The feature 'Window Function' is not supported in the current SQL dialect. 使用了窗口函数(如 LAG, ROW_NUMBER),但 Spark 版本 < 3.0,或执行模式不支持 执行 SELECT spark_version(); 确认版本;检查是否在 spark-sql CLI 中执行(支持),而非旧版 beeline 升级 Spark 至 3.0+;或改用自连接模拟窗口逻辑(不推荐,性能差)

5.2 实操中踩过的 3 个真实坑

坑一:spark-sql -f 执行时中文注释乱码
现象:SQL 文件里 COMMENT '北京市'spark-sql CLI 中显示为 ??????,建表后 DESCRIBE TABLE 看不到中文注释。
根因:SparkSQL CLI 默认使用系统 locale 编码读取文件,Linux 服务器常为 en_US.UTF-8,但文件保存为 GBK(Windows 记事本默认)。
解决:用 iconv 转码:iconv -f GBK -t UTF-8 gy_pub.sql > gy_pub_utf8.sql,然后执行 spark-sql -f gy_pub_utf8.sql终极方案:所有 SQL 文件统一用 UTF-8 无 BOM 格式保存,编辑器(VS Code/Sublime)右下角确认编码。

坑二:INSERT OVERWRITESELECT COUNT(*) 结果为 0
现象:执行完 INSERT OVERWRITE TABLE gy_pub PARTITION (dt='2024-01-01') ...,再查 SELECT COUNT(*) FROM gy_pub WHERE dt='2024-01-01' 返回 0。
根因:Spark 的 INSERT OVERWRITE 是“先删分区目录,再写新文件”,但如果写入过程中任务失败(如磁盘满),分区目录会被删,但新文件没写完,导致分区为空。
排查:hdfs dfs -ls /user/hive/warehouse/practice_db.db/gy_pub/dt=2024-01-01/,发现目录存在但无 _SUCCESS 文件或 .parquet 文件。
解决:重新执行 INSERT;或手动 hdfs dfs -rm -r /user/hive/warehouse/practice_db.db/gy_pub/dt=2024-01-01/ 清理空目录,再重试。

坑三:JOIN 结果行数远超预期,出现重复
现象:SELECT COUNT(*) FROM ds_pub d JOIN gy_pub g ON d.prov_code=g.prov_code 返回 1000 万行,但 ds_pub 只有 10 万行,gy_pub 只有 5 万行。
根因:gy_pubprov_code 不唯一(比如同一省份有多条记录,因 city_code 不同),导致 JOIN 产生笛卡尔积。
验证:SELECT prov_code, COUNT(*) FROM gy_pub GROUP BY prov_code HAVING COUNT(*) > 1,果然发现 110000 有 16 条(对应北京 16 个区)。
解决:明确 JOIN 粒度。本例中,ds_pub 是“省市”粒度,gy_pub 是“省市县”粒度,应 ON d.prov_code=g.prov_code AND d.city_code=g.city_code,而非只匹配 prov_code记住:JOIN 的粒度必须对齐,否则就是灾难。

5.3 性能优化黄金三原则(本地模式亲测有效)

  1. 小表 Broadcast,大表 Filter First
    gy_pub 全量约 5 万行,属于典型小表。在 spark-sql CLI 中执行 SET spark.sql.autoBroadcastJoinThreshold=10485760;(10MB),Spark 会自动将其广播。效果:原本 30 秒的 JOIN 查询,降到 8 秒。验证:看 Spark UI 的 Stage 页面,BroadcastExchange 节点是否出现。

  2. 分区字段永远放在 WHERE 最左侧
    WHERE dt='2024-01-01' AND prov_code='110000',比 WHERE prov_code='110000' AND dt='2024-01-01' 更优。因为 Spark 的谓词下推引擎会优先处理分区字段,尽早裁剪。虽然最终结果一样,但物理执行计划更高效。

  3. 避免 SELECT *
    即使表只有 10 个字段,也显式写出需要的字段:SELECT prov_name, city_name, pv, uv FROM ...。理由:减少网络传输量(Spark Executor → Driver)、减少内存占用(Driver 端缓存结果集)、提升序列化速度。实测:SELECT * 查 100 万行耗时 12 秒,SELECT prov_name, pv, uv 同样数据耗时 7 秒。

6. 后续可扩展方向:从练习到生产的自然演进

这两张表不是终点,而是起点。当你能流畅写出上面 5 个场景的 SQL,并理解每个 WHERE、每个 JOIN、每个 PARTITION 背后的设计意图时,就可以开始向真实生产环境迈进了。这里分享三个平滑的扩展路径,都是我带团队落地时验证过的:

路径一:接入真实数据源,替换 VALUES 和字面量
ds_pub.sql 中的 FROM (VALUES ...) 替换成 FROM ods_app_log WHERE dt='2024-01-01',把 gy_pub.sql 中的 SELECT '110000', '北京市' ... 替换成 SELECT prov_code, prov_name, ... FROM dim_province WHERE dt='2024-01-01'。这时你需要:
- 在 init_db.sh 里增加 spark-sql -e "CREATE DATABASE IF NOT EXISTS ods;"
- 提前准备好 ods_app_logdim_province 表(可用 spark-sql -f create_ods.sql 初始化);
- 修改 ds_pub.sqlgy_pub.sqlFROM 子句。
这个过程会逼你思考:上游 ODS 层的分区策略是什么?DIM 层的更新机制是全量还是增量?数据延迟 SLA 是多少?——问题从“语法怎么写”升级为“架构怎么搭”。

路径二:增加时间维度表,支撑更复杂的周期分析
目前 gy_pub 只有 year, month, day, week_of_year,但业务常需“去年同期”、“近30天滚动”、“季度累计”。可以新建 dim_date 表:

CREATE TABLE dim_date (
  dt STRING,
  year INT,
  month TINYINT,
  day TINYINT,
  week_of_year TINYINT,
  quarter STRING,
  year_month STRING,
  last_year_dt STRING,  -- 如 '2024-01-01' 对应 '2023-01-01'
  last_30days_start_dt STRING
) USING PARQUET;

然后用 ds_pub JOIN dim_date ON ds_pub.dt = dim_date.dt,就能轻松写出“同比”SQL。这个表的数据可以用 Python 脚本生成(pandas.date_range),每天 INSERT OVERWRITE 一行,非常轻量。

路径三:迁移到 Hive Metastore,实现跨引擎共享
当前所有表都存在 Spark 默认的 warehouse 目录(本地文件系统或 HDFS),只能被 Spark 访问。想让 Presto、Trino、甚至 BI 工具(Tableau/Superset)也查,就得把元数据迁到 Hive Metastore:
- 配置 Spark 使用 Hive Metastore(修改 spark-defaults.confspark.sql.hive.metastore.version 2.3.9);
- 把 init_db.sh 中的 CREATE DATABASE 改成指向 Hive 的 LOCATION
- gy_pub.sqlds_pub.sql 不用改,SparkSQL 会自动把表注册到 Hive。
这时你会发现,SHOW DATABASESspark-sqlbeeline 里看到的是一样的——数据资产真正实现了跨引擎复用。

最后再分享一个小技巧:每次写完一个新 SQL,不要急着看结果,先执行 EXPLAIN FORMATTED YOUR_SQL。仔细读它的 Physical Plan,找找有没有 WholeStageCodegen(Spark 的 JIT 编译优化)、有没有 BroadcastHashJoin(小表广播成功)、有没有 Filter: (dt#123 = 2024-01-01)(分区裁剪生效)。读懂执行计划,你就从 SQL 编写者,进化成了数据管道的架构师。

本文还有配套的精品资源,点击获取 menu-r.4af5f7ec.gif

简介:提供gy_pub.sql和ds_pub.sql两个可直接运行的SQL文件,专为SparkSQL离线分析学习场景设计。包含完整的建表语句、字段类型定义、分区字段声明及配套的INSERT示例数据,覆盖常见维度表(如地区、时间、用户属性)和事实表结构。所有语句严格遵循SparkSQL 3.x语法规范,不依赖Hive函数或外部元数据库,支持本地模式(spark-sql CLI)和YARN集群环境一键执行。适合用于练习表结构设计逻辑、数据类型选择依据、分区策略实践(如按日期分区)、基础SELECT聚合查询、多表JOIN验证以及简单ETL流程模拟。配合init_db.sh脚本可快速初始化测试库,无需额外配置即可开始SQL编写与执行训练。


本文还有配套的精品资源,点击获取
menu-r.4af5f7ec.gif

更多推荐