做数据开发遇到最烦的事情之一,就是同一个需求要同时面对七八种数据格式。上周刚帮业务部门跑通一个周报自动化,数据一半在MySQL业务库里,一半是运营同事从后台导出来的CSV日志,还有一部分历史快照是数仓用Parquet落好的。三个数据源,三种格式,各自为政,最后要合成一张宽表。如果你也经常在JDBC、CSV、Parquet之间来回切换,这篇文章应该能帮你少走不少弯路——我会把Spark做多数据源整合时的连接方式、关键参数、踩坑点,以及一个完整的ETL示例都讲清楚。
顺便说一句,文里的代码我尽量用PySpark写,但JDBC的Driver配置、Parquet的Schema演进这些机制是语言无关的,你用Scala版Spark Shell或者Spark SQL跑,思路完全一样。
1. 为什么多数据源整合是Spark的主场
1.1 数据散落是常态,而不是意外
先聊一个现实问题:一家公司里,数据从来不会乖乖待在一个地方。业务系统为了保证事务一致性,数据在MySQL或者PostgreSQL里躺着,每天十几个接口在写;运营和产品同学导出的报表数据,为了图方便,直接存成CSV扔在共享盘或者对象存储上;数仓团队则会把清洗好的历史快照、统计中间表,统一用Parquet或者ORC落盘。
这套组合拳短期没问题,但一旦业务方要"把订单数据和广告日志放一起跑个ROI",事情就麻烦了。手动从MySQL导出CSV,再拿Excel做关联?第一周可以,第二周数据量翻倍,Excel直接卡死,第三周发现上周的CSV漏导了一天数据。这时候你就需要一个统一的计算层,能把不同存储位置、不同格式的数据,一次性拉到一个引擎里做关联计算——这正是Spark的核心使用场景。
1.2 Spark统一数据访问层到底统一了什么
Spark的DataFrameReader和DataFrameWriter设计得很聪明:不管底层是关系库、文件还是消息队列,对外暴露的都是同一套API。读数据就是spark.read.format(...).option(...).load(),写数据就是df.write.format(...).mode(...).save()。格式之间的差异被封装进了各种DataSource实现里,业务代码不用关心底层是JDBC连接还是文件扫描。
这种统一带来的直接收益是代码可维护性。我见过很多团队的ETL脚本,每个数据源一套独立的Python脚本,用pymysql连库、用pandas读CSV、用pyarrow读Parquet,三个脚本三个环境,光依赖冲突就够喝一壶。换成Spark之后,一个程序入口能覆盖所有数据源,还顺手解决了分布式处理的问题——数据量从几百万涨到几亿,脚本逻辑一行都不用改,只要集群资源跟得上。
2. 环境准备:版本组合与三类数据源的能力对照
2.1 版本组合与依赖引入
我这边用的是Spark 3.2.1搭配Scala 2.12、Java 8。这个组合比较稳健,PySpark、Spark SQL、Structured Streaming都能正常跑。如果你要连MySQL,需要准备MySQL的JDBC驱动(mysql-connector-java8.x或者com.mysql.cj.jdbc.Driver对应的驱动包),放到$SPARK_HOME/jars目录下,或者提交任务的时候用--jars参数带进去。
# 提交任务时携带JDBC驱动 spark-submit \ --master yarn \ --jars /opt/drivers/mysql-connector-java-8.0.30.jar \ --driver-memory 4g \ --executor-memory 8g \ --executor-cores 4 \ etl_report.py如果只是本地练习,直接用spark-shell --packages拉依赖也行,但生产环境我强烈建议手动管理驱动jar包,别让集群在线下载依赖,网络和版本都不可控。
2.2 JDBC、CSV、Parquet三类数据源的特征对比
| 数据源 | 典型来源 | 读取速度 | Schema支持 | 推荐场景 | 主要坑点 |
|---|---|---|---|---|---|
| JDBC | MySQL、PostgreSQL、达梦、GaussDB等 | 中等(受限于源库连接和查询能力) | 强,天然有表结构 | 业务明细查询、增量抽取 | 连接数控制、DDL不适配 |
| CSV | 运营导出、日志文件、三方系统 | 快,但全量扫描 | 弱,需推断或手动指定 | 一次性分析、临时数据 | 编码、脏数据、类型推断混乱 |
| Parquet | 数仓落地、HDFS/对象存储 | 很快,列式+谓词下推 | 强,自描述文件 | 数仓明细层、宽表落地 | Schema演进、小文件问题 |
这三者不是替代关系,而是互补。JDBC适合跟在线系统交互,CSV适合接临时交付的数据,Parquet适合做长期存储和频繁分析的底座。弄清楚各自的位置,你才知道什么时候该用什么。
3. JDBC直连关系库:连接、分区与参数调优
3.1 一个最基础的JDBC读取示例
先看最朴素的写法:
orders = spark.read \ .format("jdbc") \ .option("url", "jdbc:mysql://10.0.1.100:3306/business_db?useSSL=false&serverTimezone=Asia/Shanghai") \ .option("user", "etl_read") \ .option("password", "******") \ .option("dbtable", "orders") \ .option("driver", "com.mysql.cj.jdbc.Driver") \ .load()这段代码的核心就一个动作:Spark在Executor上建立JDBC连接,执行SELECT * FROM orders,把结果集转成DataFrame。但实际生产里很少直接读全表,更常用的写法是把查询条件放进dbtable的子查询里,让数据库先做一轮过滤和裁剪:
orders = spark.read \ .format("jdbc") \ .option("url", "jdbc:mysql://10.0.1.100:3306/business_db?useSSL=false") \ .option("dbtable", "(SELECT order_id, amount, status, create_time FROM orders WHERE create_time >= '2024-01-01') t") \ .option("user", "etl_read") \ .option("password", "******") \ .option("driver", "com.mysql.cj.jdbc.Driver") \ .load()两种写法差别很大。第一种是Spark把全表数据拉到计算层再过滤,浪费IO;第二种是数据库先做投影和过滤,只把结果集回传给Spark。尽量把过滤条件下推给数据库,这是JDBC整合的第一条铁律。
3.2 分区读取:parallelism和连接数的平衡艺术
单连接读大表是不可接受的。拿一张千万级的订单表来说,单线程全表扫描,数据库要跑几分钟,Spark这边Executor全闲着。解决办法是让Spark并行拉数据——靠partitionColumn、lowerBound、upperBound、numPartitions这四个参数。
orders = spark.read \ .format("jdbc") \ .option("url", "jdbc:mysql://10.0.1.100:3306/business_db?useSSL=false") \ .option("dbtable", "(SELECT * FROM orders WHERE create_time >= '2024-01-01') t") \ .option("partitionColumn", "id") \ .option("lowerBound", "1") \ .option("upperBound", "10000000") \ .option("numPartitions", "10") \ .load()Spark拿到这四个参数之后,会把[1, 10000000]这个区间均分成10份,每个Executor负责一个区间,执行类似WHERE id >= 1 AND id < 1000000这种范围查询。好处是扫描并行度直接变成10,坏处也很明显:数据库同时要抗10个连接。如果你把numPartitions设成50,数据库就要开50个会话,连接池稍微小点就报Too many connections。
我的经验是:numPartitions最好控制在数据库max_connections的十分之一以内,同时保证分区字段上有索引。没有索引的话,数据库每个分区都是全表扫描再过滤,等于把一张表扫了N遍,性能比单连接还差。
这里还有一个MySQL专属的坑。默认情况下,Spark读MySQL的fetchsize是不生效的,结果集可能一次性全拉进内存,大表直接OOM。解决办法是在URL后面加useCursorFetch=true,再配合fetchsize参数,让JDBC驱动用游标方式流式读取:
.option("url", "jdbc:mysql://10.0.1.100:3306/business_db?useSSL=false&useCursorFetch=true") .option("fetchsize", "1000")加了这个之后,Executor内存占用会明显下降,特别是做全量抽取的时候效果很直观。
3.3 写入方向:mode、batchsize与覆盖策略
JDBC不止用来读,数仓结果回写业务库也很常见。写数据走的是DataFrameWriter:
result_df.write \ .mode("append") \ .option("batchsize", "5000") \ .jdbc("jdbc:mysql://10.0.1.100:3306/business_db?useSSL=false", "report_daily", props)batchsize控制每个批次写入多少条,默认是1000。调大之后写得更快,但单批失败的回滚成本也变高了,我一般设在3000到5000之间。mode("overwrite")配合truncate选项有个细节值得注意:Spark的overwrite模式在JDBC里是先把目标表drop掉再重建,如果你只想清空数据而不动表结构,得加.option("truncate", "true"),这样Spark会先TRUNCATE TABLE再写入,速度也更快。
3.4 JDBC整合我踩过的几个坑
- Driver class not found:最常见。驱动jar没放进
$SPARK_HOME/jars,或者提交任务时忘了--jars。本地IDE能跑、集群上跑不了,九成是这个原因。 - 国产数据库Driver类名各不相同:达梦是
dm.jdbc.driver.DmDriver,神通是com.oscar.Driver,GaussDB是org.postgresql.Driver(兼容PG协议),别拿MySQL的Driver去套。 - PostgreSQL的schema限定:
dbtable要写成public.orders,否则会去search_path里找,找不到表就报错。 - 时区问题:MySQL连接串不指定
serverTimezone,读timestamp字段可能差8个小时。统一用Asia/Shanghai。
4. CSV读写:编码、脏数据与单文件输出
4.1 读取CSV必须显式声明的参数
CSV大概是Spark所有数据源里最"随性"的格式,不同系统导出来的CSV,分隔符、引号、换行规则全都不一样。所以读取CSV的时候,我的建议是所有关键参数全部显式声明,不要依赖默认值:
ad_logs = spark.read \ .option("header", "true") \ .option("inferSchema", "true") \ .option("sep", ",") \ .option("quote", "\"") \ .option("multiLine", "true") \ .option("encoding", "UTF-8") \ .csv("hdfs:///data/ads/2024/01/")header声明第一行是列名;sep声明分隔符,有些系统导出的是Tab分隔,这里就写"\t";multiLine处理字段值里自带换行的情况,如果一个字段被双引号包裹且内部有换行,没开这个参数就会把一行数据拆成两行,后面的列全部错位;quote指定转义字符,默认是双引号。
4.2 inferSchema是方便,但别滥用
inferSchema=true会让Spark先扫描一遍整个文件来推断每列的类型,省事是真省事,坑也真坑。比如一列数据前10000行全是数字,突然第10001行出现一个"N/A",Spark推断结果可能直接是string而不是double,后面做数值计算全崩。更麻烦的是,当你读一个目录下的多个CSV时,不同文件的类型推断结果可能互相冲突。
我的做法是:大文件或者要长期复用的CSV,手动指定Schema:
from pyspark.sql.types import StructType, StructField, StringType, DoubleType, TimestampType schema = StructType([ StructField("date", TimestampType(), True), StructField("ad_id", StringType(), True), StructField("spend", DoubleType(), True), StructField("impressions", DoubleType(), True), ]) ad_logs = spark.read \ .schema(schema) \ .option("header", "true") \ .csv("hdfs:///data/ads/2024/01/")显式Schema有几个好处:读文件不需要二次扫描,速度快;类型明确,后续算子不会做奇怪的隐式转换;还有一项隐藏能力——Spark读取目录时会自动合并所有文件的Schema来对齐列,显式指定之后,列顺序和缺失字段都在你控制之下。
4.3 编码与BOM:中文CSV的两大杀手
运营同事发来的CSV,十有八九是GBK编码,从Windows的Excel里直接导出来的。Spark默认按UTF-8读,GBK文件读出来全是乱码。解决办法是显式指定读取编码:
.option("encoding", "GBK")另外还有一个更隐蔽的问题:BOM头。UTF-8的CSV如果带BOM(文件开头有三个字节EF BB BF),Spark读完第一列列名会变成\ufefforder_id,关联计算时这个列名怎么都对不上。如果你发现读出来第一列列名莫名带了个隐藏字符,八成就是BOM。处理思路有两个:
# 方案一:一次性清洗列名 ad_logs = ad_logs.withColumnRenamed( ad_logs.columns[0], ad_logs.columns[0].lstrip("\ufeff") ) # 方案二:读取后统一重命名所有列(配合显式schema)这些隐藏字符问题,在Spark SQL里排查起来相当折腾,所以我建议写个通用的读取函数,把编码处理和BOM清洗统一封装进去,全项目复用。
4.4 脏数据处理与_corrupt_record列
CSV脏数据是躲不掉的。mode选项可以控制Spark对坏行的处理策略:
PERMISSIVE(默认):把坏行放到_corrupt_record列,不中断任务。DROPMALFORMED:直接丢弃坏行。FAILFAST:遇到坏行立刻报错。
生产上我建议用PERMISSIVE,因为直接丢弃和直接报错都太极端,先让任务跑完,再用filter(col("_corrupt_record").isNull)把脏行摘出来看,能保留数据审计的线索:
ad_logs_clean = ad_logs.filter("_corrupt_record IS NULL") ad_logs_bad = ad_logs.filter("_corrupt_record IS NOT NULL")4.5 写出单文件CSV的正确姿势
业务方经常提一个需求:"把结果导成一个CSV给我。"Spark默认情况下每个分区写一个文件,200个分区就是200个part-00000.csv,业务方根本没法用。这时候要用coalesce(1):
result_df.coalesce(1) \ .write \ .mode("overwrite") \ .option("header", "true") \ .csv("hdfs:///data/export/report_20240101.csv")注意,coalesce(1)是把数据全部拉到同一个分区,小数据量无所谓,几个GB的数据这么搞,单节点内存和网络都会成为瓶颈,甚至可能OOM。数据量大的时候,更合理的做法是保持多文件输出,同时附带一个_SUCCESS标记文件,或者直接把结果写进数仓而不是导出CSV。我一般跟业务约定:小于200MB的数据可以合并成单文件,再大的数据就走Parquet落地,谁也别为难谁。
5. Parquet列式存储:性能与Schema管理的优势
5.1 Parquet为什么快:列裁剪和谓词下推的真实效果
Parquet是列式存储格式,数据按列组织存放。这意味着Spark读取时只需要扫描查询涉及的列,而不是像CSV那样把整行读进来再丢列。我手头有一张3.2亿行的用户行为表,CSV格式68GB,同样数据转成Parquet+Snappy压缩只有14GB左右,空间少了将近80%。
更关键的是谓词下推。Parquet文件内部按行组(Row Group)划分,每个行组在文件尾部记录着每列的统计信息(min/max)。当Spark执行WHERE date = '2024-01-01'时,它先读元数据,跳过那些根本不包含目标日期的行组。实测下来,在几亿行的大表上做过滤查询,从CSV的分钟级直接降到秒级。这两个特性——列裁剪和谓词下推——是Parquet成为数仓主流格式的根本原因。
5.2 强类型Schema:CSV给不了的确定性
Parquet文件自带Schema,字段名、类型、是否为空全都写在文件里,读出来是什么就是什么,不需要像CSV那样推断。这给多源整合带来一个隐形好处:你的下游逻辑是确定的。
举个例子,业务同事用CSV发来一份客户数据,里面"年龄"字段时而是整数、时而混着几个"未知";但同数据从数仓Parquet落地的话,类型就是int,空值就是null。你的清洗逻辑只需要处理null这一种情况,而不是去猜某列到底是string还是double。在数据管道多级串联的场景里,这种确定性省掉的排查时间非常可观。
5.3 分区发现与目录结构
Parquet落地通常配合分区目录使用,典型结构是这样:
hdfs:///warehouse/dwd_order/ year=2024/ month=01/ part-0000-xxx.snappy.parquet part-0001-xxx.snappy.parquet month=02/ part-0000-xxx.snappy.parquetSpark读这个目录时能自动识别year和month作为分区列:
order_dwd = spark.read.parquet("hdfs:///warehouse/dwd_order/") # order_dwd 会自动包含 year 和 month 两列如果目录层级比分区字段深,比如还有一个day=15层,但你只想读到year和month这一级,需要显式指定basePath:
spark.read.option("basePath", "hdfs:///warehouse/dwd_order/") \ .parquet("hdfs:///warehouse/dwd_order/year=2024/month=01/day=15/")分区目录还带来一个性能红利——分区裁剪。你查询WHERE month = '01'时,Spark直接跳过整个month=02目录,连文件都不用打开。这也是Parquet落地表比JDBC直连更快的原因之一,源库再好的索引,也好不过这种物理层面的跳过。
5.4 Schema演进的两种处理方式
Parquet的强类型是个优势,但也会带来麻烦:上游表结构变了怎么办?比如原来订单表没有discount字段,这周加了,新旧数据混在同一个目录里。默认情况下Spark读的时候发现两个文件的Schema不一致,直接抛异常,任务失败。
处理方式有两种。写数据时开启mergeSchema:
df.write \ .mode("append") \ .option("mergeSchema", "true") \ .parquet("hdfs:///warehouse/dwd_order/")这样写进去的新文件会保留旧字段并补齐新字段,缺列的旧文件读取时对应列显示为null。另一种方式是在读取时设置spark.sql.parquet.mergeSchema=true,让读操作容忍Schema不一致。我的建议是:写侧开启mergeSchema来演进,读侧保持严格模式。读侧宽松容易掩盖上游的Schema变更问题,等真正出事的时候非常难排查。
6. 三源合一:一个真实报表任务的完整实现
6.1 场景定义
现在把前面这些技术点串起来,做一个完整的例子。需求是生成一张每日GMV宽表,字段包括订单ID、日期、客户ID、订单金额、广告花费、客户等级。数据来源:
- MySQL业务库
orders表:订单事实,包含order_id、amount、customer_id、create_time。 - 运营团队CSV文件:广告投放日志,包含
order_id、date、spend、customer_id,注意这个CSV是GBK编码,日期是字符串。 - Parquet数仓表
dim_customer:客户维度,包含customer_id、customer_level。
6.2 完整实现代码
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window from pyspark.sql.types import StructType, StructField, StringType, DoubleType, DateType spark = SparkSession.builder \ .appName("daily_gmv_report") \ .config("spark.sql.shuffle.partitions", "200") \ .enableHiveSupport() \ .getOrCreate() # 数据源1:MySQL订单表,只取近30天 orders = spark.read \ .format("jdbc") \ .option("url", "jdbc:mysql://10.0.1.100:3306/business_db?useSSL=false&useCursorFetch=true") \ .option("user", "etl_read") \ .option("password", "******") \ .option("driver", "com.mysql.cj.jdbc.Driver") \ .option("dbtable", "(SELECT order_id, customer_id, amount, DATE(create_time) AS dt FROM orders WHERE create_time >= DATE_SUB(CURRENT_DATE, 30)) t") \ .option("partitionColumn", "id") \ .option("lowerBound", "1") \ .option("upperBound", "50000000") \ .option("numPartitions", "8") \ .option("fetchsize", "2000") \ .load() orders = orders.withColumn("dt", F.to_date("dt")) # 数据源2:CSV广告日志,显式schema + GBK解码 ads_schema = StructType([ StructField("order_id", StringType(), True), StructField("date", StringType(), True), StructField("spend", StringType(), True), StructField("customer_id", StringType(), True), ]) ads = spark.read \ .schema(ads_schema) \ .option("header", "true") \ .option("encoding", "GBK") \ .csv("hdfs:///data/ads/2024/") ads_clean = ads \ .filter(F.col("_corrupt_record").isNull() if "_corrupt_record" in ads.columns else F.lit(True)) \ .withColumn("spend", F.col("spend").cast("double")) \ .withColumn("dt", F.to_date(F.col("date"), "yyyy-MM-dd")) \ .select("order_id", "customer_id", "dt", "spend") # 数据源3:Parquet客户维度 customers = spark.read.parquet("hdfs:///warehouse/dim_customer/") # 合并两份事实:订单金额 + 广告花费 fact = orders.select("order_id", "customer_id", "dt", "amount", F.lit(None).cast("double").alias("spend")) \ .unionByName( ads_clean.select("order_id", "customer_id", "dt", F.lit(None).cast("double").alias("amount"), "spend") ) # 去重:同一个order_id在两边都有数据时,保留金额非空的那条 window = Window.partitionBy("order_id").orderBy(F.col("amount").desc_nulls_last(), F.col("spend").desc_nulls_last()) fact_dedup = fact.withColumn("rn", F.row_number().over(window)).filter("rn = 1").drop("rn") # 关联维度,写出Parquet分区表 result = fact_dedup.join(customers, on="customer_id", how="left") result = result.withColumn("gmv", F.coalesce(F.col("amount"), F.col("spend"))) result.write \ .mode("overwrite") \ .partitionBy("dt") \ .option("mergeSchema", "true") \ .parquet("hdfs:///warehouse/ads/dws_gmv_daily/")这段代码里有几个细节值得说说。unionByName是Spark 3.0之后才有的,它不要求两个DataFrame的列顺序一致,而是按列名对齐,配合allowMissingColumns=True还可以容忍两边列集合不完全一样。去重用的是窗口函数row_number(),比dropDuplicates多了可控性——可以定义"保留金额非空的那条"这样更贴近业务的去重规则。最后的coalesce把订单金额和广告花费合成一个gmv字段,两个事实源口径不同,写在明面上比藏着好。
6.3 分区写出与小文件控制
写Parquet时按dt分区,这是数仓分层的标准做法。但分区写出有个隐患:如果最后一步的DataFrame有200个Task,每个分区目录下都会散落200个小文件,日积月累就是几千个几十KB的小文件,后面读起来元数据开销巨大。
控制办法是在写出前对目标分区列做一次repartition:
result.repartition(8, "dt") \ .write \ .mode("overwrite") \ .partitionBy("dt") \ .parquet("hdfs:///warehouse/ads/dws_gmv_daily/")这样每个dt分区下最多8个文件。Spark 3.x环境还可以开启自适应查询执行(AQE),设置spark.sql.adaptive.enabled=true和spark.sql.adaptive.coalescePartitions.enabled=true,让Spark在运行时自动合并过小的分区,从源头减少小文件。
6.4 可重跑性与数据校验
ETL任务最怕"跑一半挂了,重跑一次数据翻倍"。因为用了partitionBy("dt")分区表,重跑只需要保证写的是动态分区覆盖模式。Spark 3.0以上默认支持INSERT OVERWRITE动态分区覆盖,但用DataFrameWriter写Parquet时,要注意设置分区覆盖模式:
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")设成dynamic之后,mode("overwrite")只会覆盖被本批次数据命中的分区目录,不会把整个dws_gmv_daily目录清空重来。这样即使某一天的数据重跑,其他日期的分区不受影响。
任务跑完记得做三道校验:总数校验(今日行数和源表行数对比)、主键唯一性校验(count distinct order_id等于总数)、空值校验(核心字段空值率是否超过阈值)。三道都过了再写_SUCCESS标记,调度系统看到标记才认为任务成功。
7. 高频报错速查与调参建议
最后整理一份我在多数据源整合过程中实际遇到的高频问题,按现象、根因、解决方式来列,方便你出问题的时候直接对着查。
| 现象 | 根因 | 处理方式 |
|---|---|---|
java.sql.SQLException: No suitable driver found | 驱动jar不在classpath | 把jar放进$SPARK_HOME/jars,或者提交任务时加--jars |
Communications link failure或连接被拒 | numPartitions开太大,数据库连接数被打满 | 调小numPartitions,检查源库max_connections,设置connectTimeout |
| 读取MySQL大表OOM | 默认按结果集整体拉取 | URL加useCursorFetch=true,配合fetchsize=1000~5000 |
CSV第一列带\ufeff前缀 | UTF-8文件带BOM头 | lstrip("\ufeff")清洗列名,或者文件预处理去掉BOM |
| CSV中文乱码 | 文件是GBK编码,Spark按UTF-8读 | .option("encoding", "GBK") |
Job aborted due to stage failure: Failed to merge incompatible schemas | Parquet目录下新旧文件Schema不一致 | 写侧开mergeSchema=true,或统一清洗后重写历史分区 |
| 输出几百个几十KB的小文件 | 分区数太多,没有合并 | 写前repartition(8, "dt"),开启AQE自动合并 |
Task not serializable | 在RDD算子闭包里捕获了不可序列化的连接对象 | 改用foreachPartition,在每个分区内部创建和关闭连接 |
| 日期字段莫名其妙少了8小时 | MySQL连接串没指定时区 | URL加serverTimezone=Asia/Shanghai |
有一类问题经常被忽略,是源库侧压力。JDBC直连读取,你这边看着Spark跑得欢,数据库那边可能已经报警了。所以我建议所有JDBC读取都用只读账号,并且把连接超时、查询超时都设置好,宁可任务失败重试,也不要拖垮业务库。
再补充两个调参经验。第一,spark.sql.shuffle.partitions默认200,如果你的数据总量不大,join和去重之后会产生大量空Task,可以按数据量调小到50或者100;反过来数据量大时,200又不够,会造成单个Task处理数据过多。这个参数没有唯一正确答案,我习惯按"最终输出文件大小除以128MB"来估算。第二,多个JDBC源同时读的时候,给每个源单独设置numPartitions,而不是全局共用一个值,这样小表少开连接、大表多开连接,数据库压力分布才均匀。
整合多数据源这件事,做久了你会发现,真正的难点从来不是某个API不会写,而是你对每种数据源的脾气摸得不够透——CSV的编码、JDBC的连接数、Parquet的Schema,哪一个都在不经意间给你挖坑。把这些坑的位置记下来,下次再碰到三源合一的任务,你就能直接把精力放在业务逻辑上,而不是跟格式搏斗了。