
1. 先想清楚Spark里的数据到底是怎么流动的做大数据尤其是用Spark干活的人十有八九都经历过这么一幕数据从HDFS或者对象存储上过来CSV一坨、JSON一堆、还有几个不知道谁生成的Parquet目录然后你顺手敲一句spark.read.csv(...)读进来看到_c0这种列名心态就崩了。再往后oracle、mysql里的历史数据要用Spark读取Kafka里的实时流也要落到Spark里做分析存储和读取这扇门没过好后面所有计算都是空中楼阁。所以我一直觉得Spark的数据存储与读取不只是“API怎么调”的问题它直接决定了一个任务的IO开销、资源占用、以及下游分析的爽不爽。你哪怕算子写得再漂亮只要数据以垃圾格式摆在烂目录结构里跑起来照样是几分钟起步。这个主题值得掰开了讲。1.1 存储与读取在整个Spark任务里的位置一个Spark批处理任务剥掉花里胡哨的业务逻辑后本质上就三件事读数据、算数据、写数据。网上大部分教程都在讲中间那个“算”算子、RDD、DataFrame、SQL翻来覆去但实际生产环境里最耗时的往往是开头和结尾这两段IO。我见过不止一个任务计算逻辑就一个简单的过滤聚合却因为源表是几万个小文件光扫描元数据和启动task就耗掉大半时间。存储格式和读取方式还不仅仅是性能问题它直接和正确性挂钩。数据的类型丢了、字段漂移了、null被读成空字符串“”、分区目录少了一截……这些都是在“读”这一步埋下的雷。所以我把存储与读取总结成一句话读得对、读得快计算才有意义写得好、写得稳下游才不骂人。1.2 RDD、DataFrame、Dataset读写时我先选哪个初学者最容易纠结的就是Spark到底有几种数据抽象。其实从数据读写的角度答案很朴素能用DataFrame和Dataset就不要轻易用RDD。原因有三层。第一DataFrame底层有Catalyst优化器能自动做谓词下推、列裁剪、常量折叠而RDD是纯血统逻辑优化全靠你手写同样一份数据量性能差距经常是数量级的。第二DataFrame自带Schema信息读文件时可以直接落到强类型结构上配合columnar存储可以做到“我只读需要的列”而RDD读进来就是一堆Raw Bytes类型全靠自己解析。第三从存储角度讲DataFrame/Dataset天然对接数据源APIParquet、ORC这些列式格式都是为它服务的。Dataset是DataFrame的强类型版本编码器Encoder编译后直接生成字节码速度上和DataFrame几乎没差别做着复杂业务逻辑时类型安全更香。RDD也不是完全没有出场机会比如完全没有Schema的二进制文件、需要自己控制分区的场景或者写自定义数据源插件时RDD依然是兜底方案。但凡是能用spark.read.parquet、spark.table(xxx)解决的问题就别绕到RDD上去。1.3 存储格式的本质列式、压缩与Schema聊Spark的数据存储格式绕不开两个核心维度行式还是列式、有没有自描述Schema。行式存储就是CSV、JSON、Avro这类一行数据物理上挨在一起。它的优点是想读某一行时一次IO全拿齐缺点是分析型查询往往只关心少数几个字段比如一张20列的表你只想算两列的平均值行式存储也得把那20列全部扫进来IO浪费巨大。列式存储比如Parquet和ORC文件里同一列的数据连续存放。我可以加一个超市购物的类比行式存储就像你把所有商品按订单装进购物袋想知道本周番茄卖了多少斤你得把所有购物袋翻一遍列式存储是把番茄单独摆一个货架你直接去对应区域称重就行。列式存储还天生适合压缩同一列数据的类型单一、取值往往相似压缩率显著高于混杂的行数据。“自描述Schema”是另一个容易被忽略的点。Parquet文件里本身就带着列名、类型、甚至统计信息min/max/null countSpark读它时不需要你告诉它有哪些字段这一点对后续的schema演化、谓词下推、统计优化都非常重要。而CSV、JSON这种纯文本格式字段信息全靠你推断或者硬编码说丢就丢。理解了这两个维度后面格式选型就顺理成章了。2. 数据存储格式怎么选文本、CSV、JSON、Parquet、ORC很多人在写本项目或搭建数仓的时候拿到数据就df.write.format(csv)完事图省事。我强烈不建议你这样干至少先在本地建几个不同格式的文件对比一下体积和扫描时间实测一次你就知道差距了。下面逐个说下常用格式的脾气。2.1 CSV/文本最普适但也最容易让人翻车CSV是通用性最强的格式数据库导出的首选谁都能打开。但用Spark读CSV时坑是最多的。首先是类型问题。spark.read.csv()默认是不推断类型的读进来所有列都是字符串你要用.option(inferSchema, true)才帮你猜类型。可推断出来的类型也不稳定某字符串列里混了一个非数字整列就可能被当成string导致后面聚合全乱套。然后是特殊字符问题。数据里有逗号、换行、引号是CSV最经典的三座大山。字段里带换行的文件如果不用multiLine选项Spark会直接读错行数带引号的字段要配置quote和escape。还有编码问题国内经常收到GBK编码的Excel导出文件Spark默认UTF-8不配encoding就直接乱码。我的习惯是CSV只用来做临时浏览、或对接外部无法改变的源数据但凡自己能主导落盘格式一律不选CSV。2.2 JSON半结构化数据的双刃剑JSON和CSV一样是人见人爱的格式尤其接口返回、日志采集基本都是JSON。但JSON有两个隐藏问题嵌套结构和schema漂移。Spark读JSON的方式很灵活// 标准的一行一个JSON对象 val df spark.read.json(hdfs:///data/events/20240101.log) // 多行JSON对象需要用multiline否则解析会失败 val df2 spark.read.option(multiLine, true).json(hdfs:///data/big.json)一行一个JSON对象是Spark默认接受的格式也是最推荐的日志落盘方式。如果整个文件是一个大JSON数组或者一个嵌套对象必须开multiLine。每次读这种文件我都会先head -c 500看一眼开头几行判断它的结构再决定怎么读。JSON嵌套字段读出来之后在DataFrame里表现为StructType和ArrayType。嵌套一深写SQL时df.select(payload.info.user.id)这种路径就得一层层点下去一旦某条数据缺了一层查出来全是null很多新手就在这栽跟头。另外JSON是行式存储体积大、扫描慢如果你要对半结构化数据做反复的交互式分析我建议先用一次Spark读取清洗后转成Parquet再落盘后续分析就在Parquet上做。2.3 Parquet大数据场景的默认首选如果要给Spark的数据格式选一个“默认不改”的方案我的答案就是Parquet。它是列式存储、自带Schema、原生支持嵌套结构、压缩率高而且Spark对它有高度的优化引擎支撑HDFS上跑SparkParquet基本是约定俗成的标准。读Parquet是最省心的val df spark.read.parquet(hdfs:///warehouse/dws/trade_detail) // 会自动带上分区列, 自动推断Schema写Parquet也非常直接df.write.mode(overwrite) .partitionBy(dt, city) .parquet(hdfs:///warehouse/dws/trade_detail)数据进了Parquet之后列裁剪和谓词下推是自动发生的。比如你读一个20列100GB的Parquet只select了一列Spark实际扫的IO可能只有几个GB这个效率是CSV、JSON永远给不了你的。再加上Parquet自带列统计信息min/max这类谓词在文件级别就能过滤掉大量无关数据。凡是跑批任务的中间结果、结果表我都建议无脑用Parquet。2.4 ORC和Avro什么时候不选ParquetORC长得和Parquet很像也是列式、也带Schema它最大的粉丝团是Hive。如果你的数仓底层大量使用Hive的ACID表、事务表或者历史遗留的ORC表比较多那你直接用Spark读 ORC 就行性能差异和Parquet基本一个水平。从感情上讲Spark社区对Parquet的调优更细没有强Hive依赖的话不需要刻意选择ORC。Avro则完全是另一种思路行式存储Schema用JSON定义。它的优势在于schema演化机制非常成熟字段增减、默认值管理做得很规范所以在Kafka消息、数据采集通道这种“同一份数据要被多种引擎消费”的场景里Avro的兼容性更香。但如果是给Spark做离线分析Avro的扫描性能明显弱于Parquet别把它当成大表的分析存储用。2.5 一张表看懂格式选型格式存储方式自带Schema压缩表现扫描性能Spark读取便利度推荐场景CSV/文本行式无差差差需推断/处理转义外部交换、临时预览JSON行式无差差中嵌套处理麻烦接口/日志原始数据Parquet列式有优优优离线分析、中间/结果表首选ORC列式有优优优Hive数仓存量表Avro行式有中中中多引擎消费、流式通道这张表的结论很直接分析型大数据任务默认Parquet存量的Hive生态用ORC兼容跟外部系统做交换才考虑CSV和JSON。3. 数据读取实操从文件、目录到Hive表格式选好了下一步就是怎么读。很多人的Spark入门是从spark.read.format(csv).load(path)开始的但真的搞清楚每个参数含义的人不多。这一节我会把读取姿势、Schema处理、目录发现这几个关键点挨个讲透。3.1 五种高频读取姿势Spark的DataFrameReader设计得很统一一类数据的读取姿势基本是一致的// 读CSV显式给出每个option才能避免翻车 spark.read .option(header, true) .option(delimiter, ,) .option(inferSchema, true) .option(multiLine, true) .csv(hdfs:///data/csv/2024/) // 读JSON: 默认一行一个对象, 大文件整体对象要multiline spark.read.json(hdfs:///data/json/2024/) // 读Parquet: 最省事, 什么都可以不配 spark.read.parquet(hdfs:///warehouse/ods/order_detail) // 读文本: 拿到的是整行字符串的一列value spark.read.text(hdfs:///data/logs/2024/) // 通用写法: format load, 适合程序里动态传格式 spark.read.format(parquet).load(hdfs:///warehouse/ods/order_detail)Python/PySpark写法几乎一样把spark换成session即可df (spark.read .option(header, true) .option(inferSchema, true) .csv(hdfs:///data/csv/2024/))这些方法背后都是同一个DataFrameReadercsv/json/parquet/text只是format(...)的语法糖。生产代码里我更喜欢用format(xxx).load(path)这种写法因为数据源类型经常由配置文件决定改成动态传参时不用改动调用逻辑。3.2 Schema到底该自己定义还是靠推断前面提过CSV默认不推断类型JSON和Parquet自带类型信息。真正需要你纠结的是“要不要自己定义Schema”。inferSchema看起来很智能实际在生产环境是隐患。一是它需要额外扫描一遍数据来探测类型大文件上明显拖慢读取速度二是探测结果不稳定比如某列99%都是数字但中间混了一个字符串Spark宁可把整列推断成string也不会报错等下游做sum时才原地爆炸三是null和空字符串的处理完全看心情一个列全是空值它可能给你推断成string也可能推断成null类型。我的原则很简单生产环境显式定义Schema一个字段都不许懒。比如写一个自己的传入文件规范import org.apache.spark.sql.types._ val schema StructType(Array( StructField(order_id, StringType, nullable false), StructField(user_id, LongType, nullable true), StructField(amount, DoubleType, nullable true), StructField(dt, StringType, nullable true) )) val df spark.read .option(header, true) .schema(schema) .csv(hdfs:///data/csv/2024/)有了显式SchemaCSV里混入脏数据时解析会直接报错或按要求置null而不是悄悄吞掉类型。更重要的是显式Schema让上下游协作有了“数据协议”别人看到代码就知道每一列是什么、能不能为空。省一时之事后面补坑的成本高十倍。3.3 读取路径与分区目录的细节Spark读一个目录时会自动做“分区发现”。也就是说你写到/warehouse/dws/trade_detail/dt2024-01-01/citybeijing/这种结构读取时只要指向父目录Spark就会自动把dt和city作为两列读进来你甚至不需要拿到文件名就能知道数据属于哪个分区。这是parquet分区目录最爽的一点。但有三个细节要注意。第一分区目录下通常有_SUCCESS这种隐藏文件Spark默认会过滤掉不用担心第二如果你读的是一个大的混合目录只想要部分文件可以用pathGlobFilterspark.read .option(pathGlobFilter, *.parquet) .parquet(hdfs:///warehouse/ods/mixed)第三目录嵌套层级比较随意时可以用recursiveFileLookup让它递归找文件但要注意这会让分区发现失效。还有个高发坑分区字段类型。目录名是dt20240101这种字符串但你想让它当日期类型用有时需要配合spark.sql.sources.partitionColumnTypeInference.enabled或显式Schema来兜底。我用得最多的方案是在写表时把分区列名设计成规范格式读取时永远配合显式Schema分区类型就不会乱。3.4 从Hive、JDBC、Kafka等源头读取数据文件读取只是存储读取的一部分实际工作里更多数据是放在Hive表、MySQL/Oracle、Kafka里的。Hive读取是最简单的只要SparkSession开启了enableHiveSupportval df spark.table(dwd.trade_detail) // 直接读metastore里的表 val df2 spark.sql(select * from dwd.trade_detail where dt2024-01-01)这里有一个性能要点从Hive表读数据时Spark会把过滤条件下推到Hive的元数据层尤其是分区表能只扫命中的分区。所以线上表一定要建成分区表否则全表扫描谁也救不了你。JDBC读取是接传统数据库的常见路径val jdbcDF spark.read .format(jdbc) .option(url, jdbc:mysql://host:3306/db) .option(dbtable, user_order) .option(user, xx) .option(password, xx) .option(partitionColumn, id) .option(lowerBound, 1) .option(upperBound, 100000000) .option(numPartitions, 16) .load()partitionColumnlowerBoundupperBoundnumPartitions四个参数配合Spark会把查询拆成16个并发分区去读否则就一个Driver连接单线程抽数大表能跑到天荒地老。Kafka则走format(kafka)读出来的是经过序列化的字节流配合from_json函数解析成结构化的DataFrame再落到Parquet表里。这些都是日常高频操作但每换一个数据源我都会重新确认一下资源分配一定不要一把梭。4. 数据写入实操模式、分区、压缩一个都不能少读取是基本功写入才是真正体现水准的地方。同一份数据一个会写的人和不会写的人产出的表在下游使用体验是天壤之别。写错一次模式可能把整张线上表覆盖清空分区选错可能让下游每次查询都全表扫描这些坑我都是真金白银踩出来的。4.1 四种写模式怎么选Spark的write有一个mode参数四个值模式行为适用场景errorifexists默认目标路径已存在直接报错防止误覆盖首次落盘overwrite删除目标路径后重写全量刷新维度表、结果表append追加到已有目录增量日志、日分区写入ignore目标存在则静默跳过流程幂等保护df.write .mode(overwrite) .partitionBy(dt) .parquet(hdfs:///warehouse/dws/trade_detail)我见过最惨的翻车现场是有人写增量任务时图省事用了overwrite结果把整年的历史分区全删了业务方第二天就找上门。所以增量场景老老实实用append覆盖场景优先指向一个新的临时目录确认无误后通过Hive命令切换分区“上线”这是成熟数仓的常规套路。另外overwrite在分区表上比较特殊老版本Spark默认是“删整个表目录再重写”后来出现spark.sql.sources.partitionOverwriteModedynamic之后可以做到只覆盖匹配到的分区我会在下面细讲。4.2 分区与分桶写入时就要想好怎么被读partitionBy是写入时最应该花心思的配置。它的本质是把数据按分区列的值拆成一个个子目录让下游读取时通过目录过滤来避免扫全表。df.write .mode(overwrite) .partitionBy(dt, city) .parquet(hdfs:///warehouse/dws/trade_detail)分区列的选择有三个原则低基数不要选user_id这种上亿取值的列、经常作为过滤条件、分区数量可控。一个常见误区是分区字段选得太多太碎比如按dt city channel platform四层分区每天几万个分区目录Spark读表时光是列分区的开销就够喝一壶。一般来说离线数仓里 dt 一级业务维度就够用了。分桶bucketing则更进一步把同一分区内的数据按桶字段哈希后散成固定数量的物理文件。它的好处在于做join时两边都按相同字段分桶可以直接走bucket join跳过shuffle。但分桶要求严格控制桶的数量和排序规则落地管理成本高中小团队我建议先不要主动上等真正遇到join性能瓶颈再说。我的经验是分桶是用来“解决特定问题的药”不是默认配置。4.3 压缩与文件大小控制写Parquet时默认压缩是snappy这在大多数场景下是对的速度够快压缩率适中。但如果你的表很大、冷数据比较多、查询频率低可以换zstd压缩率比snappy高出一截读取时稍微慢一丢丢也值gzip压缩率更高但CPU开销大lz4则偏吞吐。// 全局配置Parquet压缩格式 spark.conf.set(spark.sql.parquet.compression.codec, snappy) // 也可以在write时针对单个DataFrame指定 df.write.option(compression, zstd).parquet(hdfs:///warehouse/...)文件大小是另一个被忽视的魔鬼。一次写入如果生成10万个几十KB的小文件HDFS的NameNode压力先不说下游Spark读它时要为每个文件起一个task光task调度和元数据开销就能把集群拖垮反过来如果只生成一个1GB的“大文件”单任务扫描时会因为并行度不足而跑得异常慢。我一般把目标文件控制在64MB到256MB之间写入前用df.repartition(n)或df.coalesce(n)控制分区数n根据总数据量估算。写成几个文件直接决定了你后续读它的并行度。4.4 动态分区与覆盖误伤的细节前面提到过分区表上overwrite的坑这里单独说下。假设你有一张按dt分区的结果表每天跑批写入当天分区如果你用df.write.mode(overwrite).partitionBy(dt).parquet(hdfs:///...)在老版本Spark里这个操作会把整个目录删掉再写等于把过去所有历史分区一起干掉了。正确做法是开启动态分区覆盖spark.conf.set(spark.sql.sources.partitionOverwriteMode, dynamic) df.write.mode(overwrite).insertInto(dwd_trade_detail)开启之后就只覆盖匹配到的分区其他分区安然无恙。这种“写入行为语义”的差异文档里写得晦涩很多深耕两三年的人也未必知道但生产环境用错一次就是事故。我自己的习惯是凡是涉及分区表写入先在测试环境对比一次整体覆盖和动态覆盖的文件删除日志确认行为符合预期再上生产。5. 性能优化与常见问题排查实录读和写的API讲完了这一节我把这几年排障遇到的高频问题集中梳理一遍。这些问题单看报错信息都让人头大但底层原因其实就那么几类小文件、数据倾斜、内存配置、分区发现异常。5.1 小文件爆炸从源头防和事后治理小文件问题在Spark任务里出现的频次可能比你想的高得多。它的典型症状是任务明明数据量不大但stage的数量几千个运行时间大部分耗在task调度上。产生原因主要是上游写数据时分区设置不合理、每批append的增量文件太小、或者读取时过度使用repartition把数据切得太碎。从源头防是最省心的。写数据前先估算总大小目标输出文件数 总数据量 / 128MB然后用coalesce或repartition调整分区数再写。比如总数据量大约10GB我想生成约80个文件就df.repartition(80)。事后治理也有成熟手段。如果存量表已经碎成几千个小文件可以用Spark3的AQE自动合并spark.conf.set(spark.sql.adaptive.enabled, true) spark.conf.set(spark.sql.adaptive.coalescePartitions.enabled, true)或者手动重写一遍读出来coalesce(64)再写回新目录然后重建分区。治理小文件的过程本质是“读写平衡”读是便宜的重写写是昂贵的IO你要权衡划算与否。我的建议是高频使用的核心表定期治理低频历史表就让它维持原样别为了优雅牺牲太多计算成本。5.2 数据倾斜读进来容易算起来难数据倾斜是Spark性能问题的头号元凶而且它往往在存储层就埋下了种子。最典型的是分区字段取值不均一张订单表按city分区北京一个分区的数据量是其他城市总和的十倍Spark扫描这个分区时被分到的task天然就慢整条任务的短板效应立刻显现。倾斜的处理分两个层面。存储层面我建议分区字段选择“查询经常用且基数不高但分布均匀”的维度如果某个值的数据量实在太大干脆拆出“异常值分区”单独管理。计算层面Spark3的AQE开了之后可以自动检测并拆分红块spark.conf.set(spark.sql.adaptive.enabled, true) spark.conf.set(spark.sql.adaptive.skewJoin.enabled, true)但AQE只能处理join阶段的倾斜group by/聚合倾斜还是得靠业务层加盐、两阶段聚合或者广播小表避免shuffle。排查倾斜最实用的办法Spark UI里看stage中各task的运行时间和shuffle读写量如果个别task耗时是平均的10倍以上基本就是倾斜实锤。5.3 Spark内存相关的典型报错内存类的报错新人见得最多。一类是java.lang.OutOfMemoryError: Java heap space跑在Executor上多半是Executor内存给小了。另一类是Container killed by YARN for exceeding memory limits这个更坑看着像OOM但其实是物理内存超限往往是因为你虽然设了spark.executor.memory但没给足spark.executor.memoryOverhead堆外内存爆了被YARN干掉。做存储与读取时内存问题最常见于两种情况。第一种是读入超大文件时默认每个分区128MB如果你的Executor只有1GB内存一个分区数据进来再叠加中间算子的buffer很容易直接打爆。这时可以适当调小spark.sql.files.maxPartitionBytes比如设成64MB让每个task吃更少数据。第二种是collect()到大Driver这个操作无论何时都要警惕几十GB的表往Driver上搬不死才怪。真要看样例数据用limit(100)加collect就够了或者落成一个临时Parquet表慢慢看。还有一个Spark读取时的高频暗坑spark.sql.shuffle.partitions默认200很多人做聚合时没想过这个值结果一个小数据集硬生生shuffle成200个分区文件碎得一地。我会根据数据规模把spark.sql.shuffle.partitions设成 48/96/200 这种档位跟数据量和集群核数匹配。5.4 排查一次慢任务的完整思路每次被拉去救火分析任务慢在哪我都会走一套固定流程这里分享给大家。第一步看Spark UI的Event Timeline确认瓶颈在哪个stage。如果stage input大小巨大且只有一两个task说明分区数不够如果几百个task里个别特别慢八成是倾斜如果shuffle write和shuffle read之间的间隙特别长说明shuffle本身有问题。第二步看SQL执行计划用df.explain(formatted) // Spark3可看详细计划重点看Scan节点下的Filters谓词下推有没有生效、分区裁剪有没有自动加上、有没有出现非必要的Exchangeshuffle节点。如果一张表的扫描没有按分区裁剪多半是过滤条件写法问题比如把分区字段包在函数里成了where substr(dt,1,10)2024-01-01优化器无法把它转成对目录的过滤。第三步直接看日志。分析任务失败时先翻executor日志和YARN的app日志别急着改代码。很多看起来复杂的报错最后都是磁盘空间不足、权限问题、或者写路径不存在这类问题排掉一切自然恢复。这套流程能解决八成以上的“任务慢”和“任务挂”。技术人最忌讳的就是不看UI不查日志凭感觉一顿乱改参数。6. 实操总结踩过这些坑之后我的几个习惯写到这里我把这些年做Spark数据存储与读取沉淀下来的几个个人习惯总结一下。不是标准答案但照着做大概率能少走弯路。第一中间结果和结果表统一用Parquet加snappy除非遇到明确的特殊需求才换格式。这样整个数据链路的行为是可预测的谁接手都不会因为格式问题吵架。第二生产环境读任何外部文件全部显式声明Schema。宁可写50行StructType也不让Spark替我猜类型。这是花一小时省一整周的事。第三写分区表之前花30秒确认一下 mode 和 partitionOverwriteMode尤其是增量任务。我吃过的“全表被覆盖”的亏全部源于没检查这两行配置。第四每次写完表记录文件数和总大小。如果发现单个文件小于10MB或大于512MB立刻回头调整分区或repartition不要等到下游反馈再补救。第五Spark3之后我把 AQE 常开adaptive.enabled、coalescePartitions.enabled、skewJoin.enabled这三个配置能提升不少默认场景的稳定性。除非调试特殊SQL否则我没有关闭它的理由。最后再分享一个小技巧在写任何Spark任务前先到命令行hdfs dfs -ls -R看一遍输入目录的结构和文件大小分布再决定怎么读、读完后怎么写。这30秒的“侦察”对后续的任务设计和参数选择帮助巨大也是我觉得Spark数据工程师最值得养成的职业习惯之一。