资讯详情

资讯详情

PySpark多数据源整合实战:JDBC、CSV与Parquet一网打尽

做数据开发遇到最烦的事情之一就是同一个需求要同时面对七八种数据格式。上周刚帮业务部门跑通一个周报自动化数据一半在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支持推荐场景主要坑点JDBCMySQL、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?useSSLfalseserverTimezoneAsia/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?useSSLfalse) \ .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?useSSLfalse) \ .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后面加useCursorFetchtrue再配合fetchsize参数让JDBC驱动用游标方式流式读取.option(url, jdbc:mysql://10.0.1.100:3306/business_db?useSSLfalseuseCursorFetchtrue) .option(fetchsize, 1000)加了这个之后Executor内存占用会明显下降特别是做全量抽取的时候效果很直观。3.3 写入方向mode、batchsize与覆盖策略JDBC不止用来读数仓结果回写业务库也很常见。写数据走的是DataFrameWriterresult_df.write \ .mode(append) \ .option(batchsize, 5000) \ .jdbc(jdbc:mysql://10.0.1.100:3306/business_db?useSSLfalse, 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.DriverGaussDB是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分隔这里就写\tmultiLine处理字段值里自带换行的情况如果一个字段被双引号包裹且内部有换行没开这个参数就会把一行数据拆成两行后面的列全部错位quote指定转义字符默认是双引号。4.2 inferSchema是方便但别滥用inferSchematrue会让Spark先扫描一遍整个文件来推断每列的类型省事是真省事坑也真坑。比如一列数据前10000行全是数字突然第10001行出现一个N/ASpark推断结果可能直接是string而不是double后面做数值计算全崩。更麻烦的是当你读一个目录下的多个CSV时不同文件的类型推断结果可能互相冲突。我的做法是大文件或者要长期复用的CSV手动指定Schemafrom 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 BFSpark读完第一列列名会变成\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同样数据转成ParquetSnappy压缩只有14GB左右空间少了将近80%。更关键的是谓词下推。Parquet文件内部按行组Row Group划分每个行组在文件尾部记录着每列的统计信息min/max。当Spark执行WHERE date 2024-01-01时它先读元数据跳过那些根本不包含目标日期的行组。实测下来在几亿行的大表上做过滤查询从CSV的分钟级直接降到秒级。这两个特性——列裁剪和谓词下推——是Parquet成为数仓主流格式的根本原因。5.2 强类型SchemaCSV给不了的确定性Parquet文件自带Schema字段名、类型、是否为空全都写在文件里读出来是什么就是什么不需要像CSV那样推断。这给多源整合带来一个隐形好处你的下游逻辑是确定的。举个例子业务同事用CSV发来一份客户数据里面年龄字段时而是整数、时而混着几个未知但同数据从数仓Parquet落地的话类型就是int空值就是null。你的清洗逻辑只需要处理null这一种情况而不是去猜某列到底是string还是double。在数据管道多级串联的场景里这种确定性省掉的排查时间非常可观。5.3 分区发现与目录结构Parquet落地通常配合分区目录使用典型结构是这样hdfs:///warehouse/dwd_order/ year2024/ month01/ part-0000-xxx.snappy.parquet part-0001-xxx.snappy.parquet month02/ part-0000-xxx.snappy.parquetSpark读这个目录时能自动识别year和month作为分区列order_dwd spark.read.parquet(hdfs:///warehouse/dwd_order/) # order_dwd 会自动包含 year 和 month 两列如果目录层级比分区字段深比如还有一个day15层但你只想读到year和month这一级需要显式指定basePathspark.read.option(basePath, hdfs:///warehouse/dwd_order/) \ .parquet(hdfs:///warehouse/dwd_order/year2024/month01/day15/)分区目录还带来一个性能红利——分区裁剪。你查询WHERE month 01时Spark直接跳过整个month02目录连文件都不用打开。这也是Parquet落地表比JDBC直连更快的原因之一源库再好的索引也好不过这种物理层面的跳过。5.4 Schema演进的两种处理方式Parquet的强类型是个优势但也会带来麻烦上游表结构变了怎么办比如原来订单表没有discount字段这周加了新旧数据混在同一个目录里。默认情况下Spark读的时候发现两个文件的Schema不一致直接抛异常任务失败。处理方式有两种。写数据时开启mergeSchemadf.write \ .mode(append) \ .option(mergeSchema, true) \ .parquet(hdfs:///warehouse/dwd_order/)这样写进去的新文件会保留旧字段并补齐新字段缺列的旧文件读取时对应列显示为null。另一种方式是在读取时设置spark.sql.parquet.mergeSchematrue让读操作容忍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() # 数据源1MySQL订单表只取近30天 orders spark.read \ .format(jdbc) \ .option(url, jdbc:mysql://10.0.1.100:3306/business_db?useSSLfalseuseCursorFetchtrue) \ .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)) # 数据源2CSV广告日志显式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) # 数据源3Parquet客户维度 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, oncustomer_id, howleft) 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的列顺序一致而是按列名对齐配合allowMissingColumnsTrue还可以容忍两边列集合不完全一样。去重用的是窗口函数row_number()比dropDuplicates多了可控性——可以定义保留金额非空的那条这样更贴近业务的去重规则。最后的coalesce把订单金额和广告花费合成一个gmv字段两个事实源口径不同写在明面上比藏着好。6.3 分区写出与小文件控制写Parquet时按dt分区这是数仓分层的标准做法。但分区写出有个隐患如果最后一步的DataFrame有200个Task每个分区目录下都会散落200个小文件日积月累就是几千个几十KB的小文件后面读起来元数据开销巨大。控制办法是在写出前对目标分区列做一次repartitionresult.repartition(8, dt) \ .write \ .mode(overwrite) \ .partitionBy(dt) \ .parquet(hdfs:///warehouse/ads/dws_gmv_daily/)这样每个dt分区下最多8个文件。Spark 3.x环境还可以开启自适应查询执行AQE设置spark.sql.adaptive.enabledtrue和spark.sql.adaptive.coalescePartitions.enabledtrue让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或者提交任务时加--jarsCommunications link failure或连接被拒numPartitions开太大数据库连接数被打满调小numPartitions检查源库max_connections设置connectTimeout读取MySQL大表OOM默认按结果集整体拉取URL加useCursorFetchtrue配合fetchsize1000~5000CSV第一列带\ufeff前缀UTF-8文件带BOM头lstrip(\ufeff)清洗列名或者文件预处理去掉BOMCSV中文乱码文件是GBK编码Spark按UTF-8读.option(encoding, GBK)Job aborted due to stage failure: Failed to merge incompatible schemasParquet目录下新旧文件Schema不一致写侧开mergeSchematrue或统一清洗后重写历史分区输出几百个几十KB的小文件分区数太多没有合并写前repartition(8, dt)开启AQE自动合并Task not serializable在RDD算子闭包里捕获了不可序列化的连接对象改用foreachPartition在每个分区内部创建和关闭连接日期字段莫名其妙少了8小时MySQL连接串没指定时区URL加serverTimezoneAsia/Shanghai有一类问题经常被忽略是源库侧压力。JDBC直连读取你这边看着Spark跑得欢数据库那边可能已经报警了。所以我建议所有JDBC读取都用只读账号并且把连接超时、查询超时都设置好宁可任务失败重试也不要拖垮业务库。再补充两个调参经验。第一spark.sql.shuffle.partitions默认200如果你的数据总量不大join和去重之后会产生大量空Task可以按数据量调小到50或者100反过来数据量大时200又不够会造成单个Task处理数据过多。这个参数没有唯一正确答案我习惯按最终输出文件大小除以128MB来估算。第二多个JDBC源同时读的时候给每个源单独设置numPartitions而不是全局共用一个值这样小表少开连接、大表多开连接数据库压力分布才均匀。整合多数据源这件事做久了你会发现真正的难点从来不是某个API不会写而是你对每种数据源的脾气摸得不够透——CSV的编码、JDBC的连接数、Parquet的Schema哪一个都在不经意间给你挖坑。把这些坑的位置记下来下次再碰到三源合一的任务你就能直接把精力放在业务逻辑上而不是跟格式搏斗了。
觉得有用,分享给同行:

为您的企业打造数字门面

稳重轻奢商务风格,端正雅致视觉,长效耐看不易过时。

立即咨询 →