前阵子有个做数据开发的同事找我吐槽说业务方扔给他一张几十万行的 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?useSSLfalseserverTimezoneAsia/ShanghairewriteBatchedStatementstrue) \ .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 里的参数容易被忽略。先说useSSLfalse很多开发库为了省事直接关掉 SSL但如果你连接的是云数据库或者公司开启了强制 SSL 的实例这参数就会被忽略。PostgreSQL 是另一个方案参数名不同。我后文排错部分还会专门讲 useSSL 和 sslmode 的混淆问题。URL 里还有一个很实用的是serverTimezoneAsia/Shanghai如果不加MySQL Connector/J 8.0 在解析 DATETIME 时可能因为时区不对把你读出来的时间全部偏移 8 小时。fetchsize1000也是重点。MySQL 驱动默认会把查询结果一次性全部拉到客户端内存里如果一张表有几百万行单分区读取时 Executor 端很容易 OOM。加了 fetchsize 之后JDBC 会按批次从数据库取数Spark 这边再逐批消费内存压力小很多。2.2 分区读取让 Spark 并行“拉数据”而不是单线程搬运读小表无所谓分区但如果一张订单表有 5000 万行你用默认方式读会发现任务只有一个 partition从头拉到尾要跑半小时数据库和 Spark 都很痛苦。JDBC 分区的核心机制是在读取时指定一个分区列Spark 会把这个列的值域拆成多个区间每个区间生成一条独立的 SQL 去查询orders_df spark.read.jdbc( urljdbc:mysql://localhost:3306/shop?useSSLfalse, tableorders, columnorder_id, lowerBound1, upperBound100000000, numPartitions8, 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( urljdbc_url, table(SELECT * FROM orders WHERE create_time 2025-01-01) t, columnorder_id, lowerBound1, upperBound50000000, numPartitions8, propertiesprops )分区列的选择直接影响查询是否均匀。如果 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-01MySQL 在存储引擎层就能完成过滤返回给 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)headertrue告诉解析器第一行是字段名。inferSchematrue是让 Spark 自动推断每列类型。multiLinetrue必须重点关注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 推断的代价不是大表的首选inferSchematrue用起来方便但扫描全量文件来推断类型的成本不可忽视。对几百 MB 的小文件无所谓但到了几 GB 甚至几十 GB 的 CSV自动推断会让读取时间明显变长因为它需要在真正解析数据之前额外做一次全文件扫描。更重要的坑是推断结果不稳定。同一个字段如果前面 100 万行都是空值Spark 可能推断成 nullType后面 join 的时候这个字段的类型和另一个表对不上报错或者给你一堆空值。日期类型尤其惨2025-01-01能推断成 date2025/01/01可能就变成了 string同样写法的不同类型在不同批次文件里出现直接导致下游处理不一致。如果你知道表结构我强烈建议直接手写 schemafrom 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, truncateFalse)这段 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 会把所有文件里出现的字段合并起来缺失字段统一补 nulldf spark.read \ .option(mergeSchema, true) \ .parquet(/warehouse/orders_history)但这里有一个度的问题。mergeSchema需要 Spark 读取所有文件列表和 schema 元数据文件越多这个操作越慢。如果历史分区特别多我一般建议不要在默认读取里长期开启 mergeSchema而是用一次性的 ETL 任务把旧数据重写成统一 schema之后正常读取。另外schema 变更的兼容性还有一个方向字段类型演进。Parquet 支持 int 到 long、float 到 double 之类的演进但如果你把一个 string 字段改成 intSpark 读取时大概率直接报错。所以设计 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( urljdbc:mysql://dbhost:3306/shop?useSSLfalseserverTimezoneAsia/Shanghai, tableorders, columnorder_id, lowerBound1, upperBound5000000, numPartitions6, 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 或者换一种关联方式而不是死磕广播。AQEAdaptive Query Execution在高版本 Spark 默认开启它能在运行时根据实际 shuffle 数据量动态调整 join 策略和分区数。在 ETL 脚本里我习惯显式开启并设置一个相对合理的spark.sql.shuffle.partitions比如 8 或 16避免每跑一次任务都按默认 200 个分区产生一堆小文件。5.4 写出与调度建议结果写出用coalesce(1)是为了尽量少生成文件在这个场景里结果集是“每月每层一行”数据量不超过几百行一个文件完全够。如果结果有上千万行强行 coalesce 反而会成为性能瓶颈那时候应该按分区字段写多个文件让每个分区目录下文件数可控。partitionBy(month)会把输出目录切成month2025-01-01、month2025-02-01这种分区结构下游查询时可以直接做分区裁剪。Parquet 加分区是数仓最常见的方式基本等于白送一级索引。调度这一层我在生产环境里会把这段代码包成一个 Spark 任务每天定时跑。MySQL 增量数据和历史 Parquet 全量合并后相当于每次任务都把“当月订单”和“历史全量”算一遍。如果历史越来越大全量合并不是一种好方案可以改成增量合并或者用拉链表来管历史。这个思路适合刚起步的报表项目量大了再演进。6. 排错实录多源读取途中见过的四个典型现场这一节写四个我真实遇到过的问题每一个都代表一大类同学会在项目里踩的坑。6.1 useSSL 与 sslmode 的混乱现场很多项目里会看到这种写法jdbc:mysql://10.0.0.10:3306/shop?useSSLfalse如果用的是 MySQL Connector/J 8.0.xx部分版本会抛出一句 warning说 useSSL 已经废弃建议用 sslMode。有人跟着文档改成jdbc:mysql://10.0.0.10:3306/shop?useSSLfalsesslmodeDISABLED然后连接直接失败了。原因很直接useSSL是 Connector/J 5.x 时代的参数sslMode是 8.0 之后引入的参数。两个参数同时出现某些版本驱动解析会冲突。正确做法是按驱动版本选择一种MySQL Connector/J 5.1.x用useSSLfalse。MySQL Connector/J 8.0.x用sslModeDISABLED如果公司数据库没开 SSL。另外还有个隐蔽的问题很多云数据库默认开了公网连接加密这时你直接用sslModeDISABLED反而连不上需要配置服务端证书路径。遇到这类问题不要盲目在网上抄参数先看驱动版本和数据库侧 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 驱动在useCursorFetchtrue时会使用游标方式边读边取。注意这个参数必须配合fetchsize一起用否则可能出现连接长期占用的问题。有个常见的坑是你设置了fetchsize1000但 MySQL 驱动只有在useCursorFetchtrue时才会真正生效否则 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 规范这些细节点上心里有数把这三类数据源揉在一起做分析就是一晚上的事情。
