前阵子有个做数据开发的同事找我吐槽,说业务方扔给他一张几十万行的 CSV 文件,让他跟 MySQL 里的订单明细对上,最后还要按月份拆到数仓目录里存成列式格式。他第一反应是用 Python 写脚本,结果需求一天变三次:CSV 里有的列名对不上要清洗,MySQL 那边一会儿要加时间条件,一会儿要改关联字段,数仓这边又希望输出按天分区。他改到第三版的时候跟我来了一句:早知道一开始就应该用 Spark 这种多数据源整合的方案。
这句话其实就是这篇文章的由来。Spark 最被低估的能力,不在于它做复杂聚合计算有多快,而在于它能用同一套 DataFrame 和 SQL 抽象,把 JDBC、CSV、Parquet 这些八竿子打不着的输入源拉在一起做 join、过滤、聚合。这篇文章我会把实际项目里围绕这三种数据源整合的配置、分区原理、编码问题、schema 演化、性能优化和踩坑记录一次性梳理出来。适合已经在用 Spark、但经常被“数据散落在 MySQL、文件导出、数仓目录”这类场景折磨的数据开发、数仓工程师和数据工程师。
1. 多源整合的核心价值:Spark 凭什么成为数据汇聚的中心
我嘴里说的“多数据源整合”,不是指把一堆文件拷贝到同一个目录下面,再用 shell 去处理,而是有一套真正统一的抽象。
在 Spark 的世界里,无论是 MySQL 的一张业务表,还是一个带表头的 CSV 文件,或者是一堆按天分区的 Parquet 数据,经过 DataFrameReader 读进来之后都是同一套 DataFrame。数据结构统一、数据操作统一、计算引擎统一。一个 join、一个 group by、一个 insert overwrite,可以完全跨数据源完成。你在 MySQL 表上能做的操作,在 CSV 上也能做;你在 CSV 上 register 成临时视图,然后用 SQL 去 join 一个 Parquet 表,写法和 join 两张 MySQL 物理表没有区别。
这一点在生产环境里的价值非常大。日常工作里数据大概率是“分裂”的:增量订单在 MySQL,历史明细在数仓的 Parquet 文件,活动名单是从 CRM 系统导出的 CSV。如果没有统一抽象,你就得写很多胶水脚本,先用 JDBC 查一遍,再解析 CSV,再写一个 Python 脚本把两边结果 merge 起来。只要业务字段变动一次,脚本就要改动对应解析逻辑,维护成本成倍上涨。
Spark 之所以能做到这种统一,底层靠的是 Catalyst 优化器。每个数据源读进来后,Spark 不关心它是来自 MySQL 还是文件,先把它映射成逻辑计划里的一个关系算子。优化器在处理过滤、join、聚合时,统一在这个逻辑计划上做规则优化和物理计划生成。你写的 DataFrame API 或者 SQL,最后都会变成一棵可执行的算子树。所以你在代码里用filter还是写WHERE,本质上都是作用在一张虚拟表上。
从成本角度讲,多源整合最大的收益是“少写一半 ETL”。我用一张表来对比这三种数据源在实际项目里的典型性格,方便后面对号入座:
| 数据源 | 典型场景 | 最大优势 | 最大的坑 |
|---|---|---|---|
| JDBC | 在线业务库、明细查询 | 数据实时、可下推过滤 | 并发拉取容易把数据库压垮 |
| CSV | 系统导出、外部交付 | 人能直接打开,跨系统方便 | 无 schema、编码乱、文件可能切得很碎 |
| Parquet | 数仓存储、分析查询 | 列式裁剪、压缩率高 | 人类无法直接查看,需要工具配合 |
这篇文章后面三个大节,就是按上表三个数据源展开。每说一种数据源,我都只会讲项目中真正高频用到的能力和坑,不讲废话。
2. JDBC 接入关系库:连接参数、分区读取与连接健壮性
JDBC 多源整合最常见的场景,是把 MySQL、PostgreSQL、Oracle 里的业务表周期性拉进数仓。网上搜“Spark JDBC”你会看到很多入门写法,但实际跑生产,有几个细节不搞定很容易翻车。
2.1 最基础的读取配置,但细节都在 URL 和 Driver 里
一个典型的 MySQL 读取长这样:
orders_df = spark.read \ .format("jdbc") \ .option("url", "jdbc:mysql://localhost:3306/shop?useSSL=false&serverTimezone=Asia/Shanghai&rewriteBatchedStatements=true") \ .option("dbtable", "orders") \ .option("user", "root") \ .option("password", "your_password") \ .option("driver", "com.mysql.cj.jdbc.Driver") \ .option("fetchsize", "1000") \ .load()很多人连着第一次就报 ClassNotFound,原因不是代码写错,而是com.mysql.cj.jdbc.Driver这个类没有打进 Spark 任务的 classpath。在生产环境里,我一般建议用spark-submit --jars mysql-connector-j-8.0.33.jar显式把驱动带进去,或者放到 Spark 的jars目录里。别指望在代码里临时--packages就能省事,离线集群经常下载不了东西。
URL 里的参数容易被忽略。先说useSSL=false,很多开发库为了省事直接关掉 SSL,但如果你连接的是云数据库或者公司开启了强制 SSL 的实例,这参数就会被忽略。PostgreSQL 是另一个方案,参数名不同。我后文排错部分还会专门讲 useSSL 和 sslmode 的混淆问题。URL 里还有一个很实用的是serverTimezone=Asia/Shanghai,如果不加,MySQL Connector/J 8.0 在解析 DATETIME 时可能因为时区不对,把你读出来的时间全部偏移 8 小时。
fetchsize=1000也是重点。MySQL 驱动默认会把查询结果一次性全部拉到客户端内存里,如果一张表有几百万行,单分区读取时 Executor 端很容易 OOM。加了 fetchsize 之后,JDBC 会按批次从数据库取数,Spark 这边再逐批消费,内存压力小很多。
2.2 分区读取:让 Spark 并行“拉数据”而不是单线程搬运
读小表无所谓分区,但如果一张订单表有 5000 万行,你用默认方式读,会发现任务只有一个 partition,从头拉到尾要跑半小时,数据库和 Spark 都很痛苦。
JDBC 分区的核心机制是在读取时指定一个分区列,Spark 会把这个列的值域拆成多个区间,每个区间生成一条独立的 SQL 去查询:
orders_df = spark.read.jdbc( url="jdbc:mysql://localhost:3306/shop?useSSL=false", table="orders", column="order_id", lowerBound=1, upperBound=100000000, numPartitions=8, properties={ "user": "root", "password": "your_password", "driver": "com.mysql.cj.jdbc.Driver", "fetchsize": "1000" } )这里的column必须是数值列或者时间戳列,要是日期时间字段。SPark 会把 lowerBound 到 upperBound 这个区间切成 numPartitions 份,然后每个 Executor 执行类似SELECT * FROM orders WHERE order_id >= ... AND order_id < ...这样的查询。要注意:lowerBound 和 upperBound 不是过滤条件,只是切分区间用的边界。你仍然可以在 SQL 层面加条件,比如把table参数写成子查询:
orders_df = spark.read.jdbc( url=jdbc_url, table="(SELECT * FROM orders WHERE create_time >= '2025-01-01') t", column="order_id", lowerBound=1, upperBound=50000000, numPartitions=8, properties=props )分区列的选择直接影响查询是否均匀。如果 order_id 是从 1 到 5000 万连续递增的,切成 8 个区间非常平均。但如果主键不是连续自增,或者中间有大量空洞,有的分区可能查出来几十万行,有的分区只查出来几百行,就会出现数据倾斜。多源整合场景里,我更倾向用一个时间列做分区,因为业务表大多数按时间写入,数据分布相对稳定。
2.3 把过滤条件下推,别让数据库白忙活
用 JDBC 读数据库时,一个很常见的误区是“先把整张表读进 Spark 再 filter”。这在逻辑上没错,但性能上很亏。
Spark Catalyst 在处理 JDBC 数据源时,会把 DataFrame 上的 filter 条件下推到数据源。比如你写了:
df = orders_df.filter("create_time >= '2025-01-01'")最终生成的 JDBC SQL 极可能是SELECT * FROM orders WHERE create_time >= '2025-01-01',MySQL 在存储引擎层就能完成过滤,返回给 Spark 的数据量大幅减少。这就是谓词下推(Predicate Pushdown)。
验证下推是否生效很简单,把 Spark 日志开到 INFO 级别,看物理计划里读取 JDBC 时的 SQL 语句;或者在 MySQL 侧开 general log,看 Spark 实际发了哪些 SQL 出来。如果发现 filter 没下推,通常是你在 Python 代码里用了 UDF 或者在 DataFrame 上做了复杂表达式,导致 Catalyst 无法识别。保持过滤条件用简单表达式,下推效率最高。
下推还有一个好处:可以在数据库侧直接走索引。如果过滤条件是主键或者有索引的字段,MySQL 的查询性能可能比你想象中高一个量级。ES 上没索引的字段,比如对某个状态值做过滤,即使下推了,数据库还得全表扫,但至少网络传输少了,这也是值得的。
2.4 一次 ETL 把数据库压垮的故事
这不是段子,是真实事故。有次我们做一次全量拉取,为了跑得快,把 numPartitions 设成了 32,让 32 个 Executor 同时从业务 MySQL 拉数据。凌晨跑批刚开始五分钟,值班监控就报了 MySQL 连接数超过 max_connections,大量查询排队,连线上应用都受了影响。
JDBC 直连数据库不等于访问本地 HDFS,数据库的并发承载能力是有限的。尤其业务库,每增加一个并发查询,都是一次额外的磁盘 IO 和内存消耗。合理的做法是:
- 分区数控制在 6 到 10 之间,最多不超过 MySQL 能同时接受的活跃查询数。
- 用 fetchsize 控制单次从 ResultSet 拉取的行数,避免大结果集撑爆内存。
- 抽数时间尽量安排在业务低峰期,例如凌晨 1 点到 6 点。
- 源库是主从架构时,优先读从库。
- 如果只是同步,考虑用 DataX 这类专用同步工具,而不是让 Spark 直连。
连接池这块也要注意。Spark JDBC 读取不会复用已有的数据库连接池,每个 partition 的查询都会新建连接。如果你还在 URL 里配置了很多无用的连接池参数,不仅不会生效,还可能引起 Driver 解析错误。连接参数保持精简,越容易排查问题。
3. 读写 CSV 不翻车的诀窍:编码、schema 推断与输出规范
CSV 是业务交付数据最常用的格式,Excel 能打开、能编辑、能导出,跨系统也方便。但在 Spark 里用 CSV 坑也不少,尤其是中文乱码、schema 推断、写出文件分片这几个点,几乎每个项目都要踩一遍。
3.1 读 CSV 的固定姿势:header、inferSchema、charset 一个都不能少
用 Spark 读 CSV 的标准写法如下:
user_tags_df = spark.read \ .option("header", "true") \ .option("inferSchema", "true") \ .option("delimiter", ",") \ .option("charset", "UTF-8") \ .option("multiLine", "true") \ .option("escape", "\\") \ .csv("/data/tags/2025/user_tags.csv")header=true告诉解析器第一行是字段名。inferSchema=true是让 Spark 自动推断每列类型。multiLine=true必须重点关注:CSV 里如果某个字段的值里包含换行符,而它又被双引号括起来了,不加 multiLine 会把一条记录拆成两行,解析结果全错。很多报表系统导出的 CSV 都有这种问题。
charset更是中文环境的高频坑。我们手里拿到的 CSV 文件来源五花八门,有从 SAP 导出的,有从老旧的 Windows 系统生成的,这些文件可能是 GBK 编码。你用默认的 UTF-8 去读,出来的中文全是乱码。遇到这种文件,把 charset 改成GBK就是最简单的解法:
df = spark.read.option("charset", "GBK").option("header", "true").csv("/path/to/gbk_file.csv")3.2 中文乱码和 BOM 的来龙去脉
乱码分两种。一种就是上面说的编码不匹配,文件本身是 GBK,你却用 UTF-8 解析。另一种是 BOM 问题。UTF-8 文件开头可能会有三个不可见字节EF BB BF,就是 BOM 头,部分 Windows 工具生成 CSV 时会自动加。Spark 在解析时,如果没正确处理 BOM,第一列列名前面会多一个\ufeff字符,导致你后面写 SQL 时列名对不上,甚至 join 时明明同名却关联不上。
判断方法很简单:打印 DataFrame 的 columns,看第一列名前面是不是有一个奇怪的前缀。处理方式是在读取之后做一次列名清洗:
from pyspark.sql.functions import col for c in df.columns: if c.startswith("\ufeff"): df = df.withColumnRenamed(c, c.replace("\ufeff", ""))如果你用 Spark 工具读了多个 CSV 文件,每次都要这样做,最好封装一个函数,读入后统一处理 BOM。这样比反复改源文件省事得多。
3.3 schema 推断的代价:不是大表的首选
inferSchema=true用起来方便,但扫描全量文件来推断类型的成本不可忽视。对几百 MB 的小文件无所谓,但到了几 GB 甚至几十 GB 的 CSV,自动推断会让读取时间明显变长,因为它需要在真正解析数据之前额外做一次全文件扫描。
更重要的坑是推断结果不稳定。同一个字段,如果前面 100 万行都是空值,Spark 可能推断成 nullType,后面 join 的时候这个字段的类型和另一个表对不上,报错或者给你一堆空值。日期类型尤其惨,2025-01-01能推断成 date,2025/01/01可能就变成了 string,同样写法的不同类型在不同批次文件里出现,直接导致下游处理不一致。
如果你知道表结构,我强烈建议直接手写 schema:
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, DateType schema = StructType([ StructField("user_id", IntegerType(), True), StructField("level", StringType(), True), StructField("tag_name", StringType(), True), StructField("create_time", DateType(), True) ]) df = spark.read \ .option("header", "true") \ .option("charset", "UTF-8") \ .schema(schema) \ .csv("/data/tags/2025/user_tags.csv")显式 schema 有几个直接好处:一是省去推断扫描,读取更快;二是类型稳定,不会因为数据内容变化导致同一字段在不同批次变成不同类型;三是在写 Parquet 之前就能把类型规范化,避免下游解析出错。
3.4 写出 CSV:面对“就要一个文件”的需求
Spark 写 CSV 时,默认是每个 partition 写一个文件。如果你读入时有 200 个分区,写出来就是 200 个 part-xxx.csv。但在很多交付场景里,对方就要一个文件,怎么办?
df.coalesce(1) \ .write \ .mode("overwrite") \ .option("header", "true") \ .option("charset", "UTF-8") \ .csv("/data/output/user_tags_output")coalesce(1)会把数据集中到一个分区后再写,这样最终只有一个 part 文件。但要注意,coalesce(1)是把所有数据拉到同一个 Executor 上去写,数据量很大的时候,这个 Executor 会成为瓶颈,可能直接 OOM。我的建议是:CSV 交付场景一般数据量不会太大,如果超过 1 GB 还非要用一个 CSV 交付,要么接受多个 part 文件,要么提前做一版聚合压缩,而不是硬撑。你还可以用repartition(1),它在极端情况下会触发一次 shuffle,但能保证数据分发均匀。
另外,Spark 写出的 CSV 在 Windows Excel 里打开可能乱码,因为部分旧版 Excel 默认按 ANSI/GBK 解析文本文件。要彻底解决,可以在 Codec 层把输出文件加上 UTF-8 BOM。Spark 自带的 CSV writer 不支持直接写 BOM,你可以改用coalesce(1)生成文件后,再写一个小脚本用 HDFS API 给文件头部补三个字节;或者干脆在代码里先用普通方式写,然后对第一个分区做一次额外处理。这个技巧在真实项目中很实用,但多数文章不会写。
3.5 用 SQL 直接过滤 CSV 数据
有很多开发者会习惯性地在 Python 里逐行读 CSV 再过滤,但在 Spark 里,你完全可以把 CSV 当成一张数据表来用:
user_tags_df.createOrReplaceTempView("user_tags") spark.sql(""" SELECT user_id, level, tag_name FROM user_tags WHERE level IN ('高价值', '中价值') AND tag_name IS NOT NULL """).show(20, truncate=False)这段 SQL 背后的执行逻辑,跟你查询一张 MySQL 表没有本质区别。它支持 WHERE、GROUP BY、JOIN、窗口函数,甚至可以直接和 JDBC 表 join。CSV 再也不是一个“只能用 pandas 读”的中间产物了。
这里多说一句:SQL 过滤 CSV 时,如果字段类型没有显式指定,Spark 会按字符串处理。字符串和数字比较时它会尝试隐式转换,但为了保险起见,建议先通过显式 schema 或cast转换,再参与条件过滤,避免因为类型隐式转换导致过滤结果不符合预期。
4. Parquet 的列存优势与 schema 演化
Parquet 在数仓里几乎是默认的存储格式,但在很多刚开始用 Spark 的人眼里,它只是“一种比 CSV 高级的文件格式”而已。这里我想把它的原理和实际收益说透,否则你只会用它,但不知道为什么该用它。
4.1 为什么同一份数据,Parquet 比 CSV 快好几倍
CSV 是行式纯文本,每一行都被完整存储。你要计算,就必须把整行读进内存,再在代码里把字符串解析成字段。Parquet 是列式存储,同一列的数据在物理上顺序排列在一起,并且每个列都有独立的统计信息和压缩编码。
最直观的时间对比:同样一份 10 GB 的用户行为日志,存成 CSV 可能要 10 GB,存成 Parquet 压缩完可能只要 2 GB 到 3 GB。查询的时候,如果你只需要user_id和action两列,Spark 只需要读取 Parquet 文件里这两列的数据块,CSV 则必须把全部 10 GB 都读一遍。在大宽表场景下,这种差距可以达到 10 倍甚至更大。
所以在多源整合里,我的习惯是:CSV 或 JDBC 进来的原始数据,如果后面还要反复查询,先落一遍 Parquet。这个动作本身花费几秒钟,但后续的每个查询都会受益。
4.2 列剪枝与谓词下推的实际效果
Parquet 的列剪枝(Column Pruning)很容易理解:Spark SQL 的物理计划会分析你最终需要哪些列,只从 Parquet 文件读取这些列对应的数据块。比如一个表有 50 个字段,你的 SQL 只 select 其中 3 个,那另外 47 个字段的数据块在读文件阶段就被跳过了。
谓词下推在此基础上更进一步。Parquet 文件每个行组都会记录列的 min/max 统计信息,Spark 在执行WHERE过滤时,会先根据这些统计信息跳过不符合条件的行组。比如你要查dt = '2025-01-01'的数据,而文件里某个行组的 dt 的 min 是 1 月 2 日,这个行组整个就不会被读取。这就是为什么在 Parquet 表上做分区字段过滤,性能可以快到几乎没有读取成本。
这一点对多源整合极其关键。你在 MySQL 上过滤是在数据库存储引擎里做,很灵活;你在 Parquet 上过滤则是在读取阶段做,而且是通过文件元数据跳过数据块实现的,速度非常快。所以 ETL 链路里,中间结果一旦转成 Parquet,后续所有层级的查询性能都会上一个台阶。
4.3 schema 演化:加列、mergeSchema 与兼容性
Parquet 文件自带 schema 信息,这既是优势也是要注意的坑。优势是 Spark 读取时不需要任何配置就能知道每一列的类型;坑是如果上游改了 schema,下游读旧文件可能遇到兼容性问题。
最常见的场景是:原始订单历史数据存成 Parquet,跑了一周后发现要新增一个coupon_amount字段。业务表里加了这个字段,但历史文件里没有。此时 Spark 读出来,旧文件这个字段会是 null,但如果你用mergeSchema选项,Spark 会把所有文件里出现的字段合并起来,缺失字段统一补 null:
df = spark.read \ .option("mergeSchema", "true") \ .parquet("/warehouse/orders_history")但这里有一个度的问题。mergeSchema需要 Spark 读取所有文件列表和 schema 元数据,文件越多,这个操作越慢。如果历史分区特别多,我一般建议不要在默认读取里长期开启 mergeSchema,而是用一次性的 ETL 任务把旧数据重写成统一 schema,之后正常读取。
另外,schema 变更的兼容性还有一个方向:字段类型演进。Parquet 支持 int 到 long、float 到 double 之类的演进,但如果你把一个 string 字段改成 int,Spark 读取时大概率直接报错。所以设计 schema 时,宁可一开始宽一点,用 string 存可能变化较大的字段,也不要为了省一点空间把字段类型卡得很死。数仓建模有个经验:能用 string 就用 string,数字全是 long,时间全是 timestamp,布尔用 boolean,数组用 array。这个原则放到 Parquet schema 设计里也能少踩很多坑。
5. 多源整合实战:一个订单分析场景的完整落地
前几节把三种数据源的原理分别讲了一遍,现在把它们组合起来,做一个完整度比较高的例子。这个例子是我在某电商数据分析项目中实际做过的一个简化版,你可以直接改改表名和字段名拿去复用。
5.1 场景设定:一份 CSV 标签、一张 MySQL 订单、一批历史 Parquet
假设业务方要一份按月份、按用户分层分组的订单统计结果,用于渠道运营分析。数据来源有三个:
- MySQL 库
shop中有一张orders表,存最近 3 个月的订单明细,字段有order_id、user_id、amount、create_time。 - 数仓 HDFS 上有一批历史订单 Parquet 文件,目录是
/warehouse/orders_history,字段和 MySQL 表基本一致,但按dt分区。 - 用户分层标签来源于一张 CSV 文件
/data/user_tags.csv,字段有user_id、level、tag_name。
需求:输出 2025 年每个月、每个用户层级的订单量和 GMV,把结果写回数仓的一个 Parquet 表。
5.2 核心代码:三条读取链路与一次 SQL 整合
代码直接用 PySpark 写,整体逻辑非常清晰:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("multi_source_orders_analysis") \ .config("spark.sql.adaptive.enabled", "true") \ .config("spark.sql.shuffle.partitions", "8") \ .getOrCreate() orders_mysql = spark.read.jdbc( url="jdbc:mysql://dbhost:3306/shop?useSSL=false&serverTimezone=Asia/Shanghai", table="orders", column="order_id", lowerBound=1, upperBound=5000000, numPartitions=6, properties={ "user": "root", "password": "your_password", "driver": "com.mysql.cj.jdbc.Driver", "fetchsize": "1000" } ) orders_history = spark.read.parquet("/warehouse/orders_history") schema = StructType([ StructField("user_id", IntegerType(), True), StructField("level", StringType(), True), StructField("tag_name", StringType(), True) ]) user_tags = spark.read \ .option("header", "true") \ .option("charset", "UTF-8") \ .schema(schema) \ .csv("/data/user_tags.csv")读取完三段数据之后,分别注册成临时视图:
orders_mysql.createOrReplaceTempView("orders_mysql") orders_history.createOrReplaceTempView("orders_history") user_tags.createOrReplaceTempView("user_tags") result = spark.sql(""" SELECT date_trunc('month', o.create_time) AS month, COALESCE(t.level, '未知') AS user_level, COUNT(DISTINCT o.order_id) AS order_cnt, SUM(o.amount) AS gmv FROM ( SELECT order_id, user_id, amount, create_time FROM orders_mysql UNION ALL SELECT order_id, user_id, amount, create_time FROM orders_history ) o LEFT JOIN user_tags t ON o.user_id = t.user_id WHERE o.create_time >= '2025-01-01' AND o.create_time < '2026-01-01' GROUP BY 1, 2 """) result.coalesce(1) \ .write \ .mode("overwrite") \ .option("compression", "snappy") \ .partitionBy("month") \ .parquet("/warehouse/dws/monthly_user_order")这几行代码值得拆开讲的内容不少。
把 MySQL 表和 Parquet 历史表通过UNION ALL合并,就完成了增量数据和历史数据的统一。这是多源整合最常见的用法——你不需要在应用层区分哪些订单在 MySQL、哪些订单在 Parquet,引擎统一处理。
date_trunc('month', o.create_time)是月份聚合的高频函数,比你自己拼字符串或截取日期更高效,也天然适配 timestamp 类型。如果你想做日期加减,Spark SQL 里date_add、date_sub、add_months都很好用,比如要统计“近 30 天”可以直接WHERE o.create_time >= date_add(current_date(), -30)。
LEFT JOIN这里的数据倾斜风险也要注意。如果某些user_id在 CSV 标签里有重复,join 后订单量会被放大;反之如果没有标签,会留下 null,需要在聚合时用COALESCE(t.level, '未知')兜底。真实业务里我遇到过标签表用户重复的问题,清洗逻辑必须提前做。
5.3 这个场景里的连接策略与优化选择
有人会问:user_tags 可能只有几十万行,MySQL orders 有上千万行,Spark 在做 left join 时需要 shuffle 吗?答案是:取决于 Spark 如何选择 join 策略。
Spark 默认有一个广播阈值,配置项是spark.sql.autoBroadcastJoinThreshold,默认 10 MB。如果右表大小低于这个阈值,Spark 会把它广播到所有 Executor,避免产生 shuffle。在这个例子里,几十万行的 CSV 标签,不出意外会在广播阈值内。但如果是几 GB 的标签表,Spark 就不得不走 sort merge join 了。
这里正好带出网上经常搜到的那句话:“left outer join 只能广播右侧”。这不是 Spark 的 bug,而是广播 join 的实现限制。左表作为驱动表,右表作为被广播的 build 表,能保证左表所有行都保留。你要是把左侧写成广播表,就没法在并行计算时完整保留右侧未匹配的行。所以项目里如果遇到“左表很小、右表很大”的 left join 场景,我会考虑把 SQL 改写成 right join 或者换一种关联方式,而不是死磕广播。
AQE(Adaptive Query Execution)在高版本 Spark 默认开启,它能在运行时根据实际 shuffle 数据量动态调整 join 策略和分区数。在 ETL 脚本里,我习惯显式开启并设置一个相对合理的spark.sql.shuffle.partitions(比如 8 或 16),避免每跑一次任务都按默认 200 个分区产生一堆小文件。
5.4 写出与调度建议
结果写出用coalesce(1)是为了尽量少生成文件,在这个场景里结果集是“每月每层一行”,数据量不超过几百行,一个文件完全够。如果结果有上千万行,强行 coalesce 反而会成为性能瓶颈,那时候应该按分区字段写多个文件,让每个分区目录下文件数可控。
partitionBy("month")会把输出目录切成month=2025-01-01、month=2025-02-01这种分区结构,下游查询时可以直接做分区裁剪。Parquet 加分区是数仓最常见的方式,基本等于白送一级索引。
调度这一层,我在生产环境里会把这段代码包成一个 Spark 任务,每天定时跑。MySQL 增量数据和历史 Parquet 全量合并后,相当于每次任务都把“当月订单”和“历史全量”算一遍。如果历史越来越大,全量合并不是一种好方案,可以改成增量合并或者用拉链表来管历史。这个思路适合刚起步的报表项目,量大了再演进。
6. 排错实录:多源读取途中见过的四个典型现场
这一节写四个我真实遇到过的问题,每一个都代表一大类同学会在项目里踩的坑。
6.1 useSSL 与 sslmode 的混乱现场
很多项目里会看到这种写法:
jdbc:mysql://10.0.0.10:3306/shop?useSSL=false如果用的是 MySQL Connector/J 8.0.xx,部分版本会抛出一句 warning,说 useSSL 已经废弃,建议用 sslMode。有人跟着文档改成:
jdbc:mysql://10.0.0.10:3306/shop?useSSL=false&sslmode=DISABLED然后连接直接失败了。原因很直接:useSSL是 Connector/J 5.x 时代的参数,sslMode是 8.0 之后引入的参数。两个参数同时出现,某些版本驱动解析会冲突。正确做法是按驱动版本选择一种:
- MySQL Connector/J 5.1.x:用
useSSL=false。 - MySQL Connector/J 8.0.x:用
sslMode=DISABLED,如果公司数据库没开 SSL。
另外还有个隐蔽的问题,很多云数据库默认开了公网连接加密,这时你直接用sslMode=DISABLED反而连不上,需要配置服务端证书路径。遇到这类问题,不要盲目在网上抄参数,先看驱动版本和数据库侧 SSL 策略。在 Spark 项目里,如果只是拉数仓同步,且网络环境是内网,我一般直接关 SSL 并限制白名单访问。
6.2 JDBC 流式读取的认知和实操
“JDBC 查询流式输出”是很多人在搜的词,尤其在数据同步场景。在 MySQL 驱动底层,所谓流式读取其实是指通过Statement.setFetchSize(Integer.MIN_VALUE)让驱动每次只从服务端拉取一部分数据,而不是把整个 ResultSet 一次性加载到 JVM。
在 Spark JDBC 读取里,对应参数是这个:
properties.setProperty("useCursorFetch", "true"); properties.setProperty("fetchsize", "1000");加了这个配置,Spark 读取大表时,每个 partition 的 JDBC ResultSet 不会一次性被拉完,而是按 1000 行一批次消费,这样 Executor 端内存占用能降下来。MySQL 驱动在useCursorFetch=true时会使用游标方式,边读边取。注意这个参数必须配合fetchsize一起用,否则可能出现连接长期占用的问题。
有个常见的坑是:你设置了fetchsize=1000,但 MySQL 驱动只有在useCursorFetch=true时才会真正生效,否则 fetchsize 被忽略。所以如果发现从 MySQL 拉数据时 Executor 内存占用异常高,先检查这两个参数是否都配了。
6.3 小文件把 NameNode 打爆
这是每天都在发生的问题。一个 CSV 目录里可能有几百个 part 文件,你用 Spark 读完之后,再write.parquet("/warehouse/output"),如果不对输出做 coalesce 或者 repartition,写入的 Parquet 文件数量大概率跟输入的分区数一致。几百个文件还好,如果源数据有几万个分区文件,写出来就可能生成几万个 Parquet 小文件。小文件过多会带来两个直接后果:一是 NameNode 内存被大量元数据占用,二是后续任何全表扫描任务启动成本都会高得离谱。
解决思路有三层:
- 在写出前先
repartition(分区数)或coalesce(目标分区数)控制并发度。 - 在 SQL 场景里,合理设置
spark.sql.shuffle.partitions,不要默认 200,很多 ETL 任务根本不需要 200 个 shuffle 分区。 - 使用 Spark 3.2 之后的 AQE,它会自动合并小分区。
6.4 LEFT JOIN 广播限制的“伪报错”
有人跑任务会看到类似下面这种日志:
Cannot broadcast the table that is larger than 8GB: ... LEFT SEMI/ANTI join cannot broadcast the left side这不是代码语法错误,而是 Spark join 策略选择失败。left join 场景下左侧不能作为广播表,Spark 只能尝试换 sort merge join。如果你的任务之前因为某个表小、想用 broadcast hint 加速,但在 left join 里把左侧标成 broadcast 了,就会出现这个限制。
处理办法有两个方向:
- 检查 SQL 逻辑,把广播 hint 放到右表。
- 如果确实左表小、右表大,可以用 RIGHT OUTER JOIN 改写,把右表变成驱动侧。例如原 SQL 是:
SELECT ... FROM small_table s LEFT JOIN big_table b ON s.id = b.user_id可以改写成:
SELECT ... FROM big_table b RIGHT OUTER JOIN small_table s ON s.id = b.user_id这样语义不变,但广播策略的适用性变了。实际项目里大数据量 join 不多不少总会有这种细节问题,明白原理之后,改起来不慌。
回到文章开头那句结论:Spark 多数据源整合,真正解决的是数据工程师的维护成本和执行效率问题。JDBC 负责接关系库,CSV 负责接外部交付,Parquet 负责把中间结果稳定下来。只要你在读取参数、分区策略、schema 规范这些细节点上心里有数,把这三类数据源揉在一起做分析,就是一晚上的事情。