在Spark里做数据存储和读取,说简单也简单,df.read.parquet、df.write.csv,API一行就完事了;说复杂也复杂,格式怎么选、分区怎么定、小文件怎么治、谓词下推为什么没生效,全是坑。我最早接手一个数据平台的时候,每天凌晨跑批,一个ETL任务光读取阶段就要跑40分钟,后来换了存储格式、调了分区策略,直接把时间砍到10分钟以内。所以这篇东西,我不打算复述官方文档,而是从一个实际跑过生产任务的人的角度,把Spark的数据存储与读取方式里最关键、最容易被忽略的细节拆开讲清楚。
这篇文章适合谁?刚上手Spark、想搞明白为什么读Parquet比读CSV快的人;遇到任务跑得慢但说不清卡在读还是卡在算的人;准备大数据面试、被问到“Spark支持哪些数据源”“如何选择存储格式”时想要系统性答案的人。我会把RDD/DataFrame/Dataset的差异、存储层选型、文件格式对比、实际代码调优、缓存与Checkpoint的用法、常见报错排查全部串起来讲,每一个点都会说明“为什么这么做”,而不是只丢结论。
1. 动手之前,先搞清楚Spark的存储逻辑
1.1 三种核心抽象,为什么读写行为完全不一样
提到Spark的存储与读取,绕不开RDD、DataFrame、Dataset这三兄弟。RDD是最底层的抽象,它只告诉Spark“我有这些分区,每个分区里是一堆Java/Scala对象”,Spark不知道里面的字段结构,所以读数据的时候只能靠用户自己手动解析。
DataFrame和Dataset则是带Schema的分布式表结构。DataFrame是Row对象的集合,编译期不检查类型,运行时报错;Dataset是强类型,用case class绑定结构,编译期就能发现问题。这三者在读写上的差异体现在:RDD读数据往往意味着全量序列化和反序列化对象,没有列剪枝、没有谓词下推,性能天然吃亏;DataFrame/Dataset走Catalyst优化器,会做逻辑计划和物理计划优化,读Parquet的时候能自动跳过不需要的列。
所以我的建议是,除非你在做非常底层的自定义算子,否则别直接用RDD去读数据,直接指定Schema用DataFrame或Dataset去读,后面要做的优化空间大得多。还有一个容易被忽略的点:Dataset在shuffle过程中走的编码器比Java序列化高效不少,尤其当你用强类型方式读数据再转DataFrame做map操作时,编码器的性能优势非常明显。我之前在百GB级日志解析任务里对比过,用Dataset读JSON再做flatMap,全程比RDD的sc.textFile(...).map(...)快了一倍以上,GC占用也小很多。
1.2 存储层选型:HDFS、对象存储还是本地盘
Spark本身不存数据,它只是一个计算引擎,数据放在哪里决定了路径的写法、I/O的瓶颈在哪、以及读取时的容错方式。业界最常见的三块存储:
- HDFS:大数据生态的老底座。路径写
hdfs://nameservice/data/...,块存储配合副本机制保证容错,在配合YARN调度时数据本地性最好。缺点是NameNode会成为瓶颈,目录结构频繁改动、小文件过多都会压垮元数据服务。 - S3/OSS/COS这类对象存储:云上标配。路径写
s3a://bucket/...,读取时会先列目录再拉对象,对“目录遍历”特别敏感。分区越多、目录层级越深,list操作越慢。优点是弹性扩容、成本低,缺点是延迟比HDFS高一个量级,而且没有数据本地性概念。 - 本地文件系统:只适合小规模demo和单机调试。
file:///...读取本地盘速度最快,但一到分布式环境就失效,因为每个executor所在的机器不一定有那个文件。
选型上没有绝对最优。我的经验是:如果公司已经有Hadoop集群,数据仓库落地首选HDFS;如果跑在云上或者数据湖方案走对象存储,那就用S3,但一定要学会合并小文件、控制分区层级深度。尤其要注意,对象存储的rename操作极其昂贵,Spark写数据时先写临时目录再rename到目标目录,这套机制在HDFS上还行,在S3上经常引发“目录已存在”或“rename失败”的诡异报错,后面我会专门讲。
1.3 序列化与内存模型:读写之外的隐藏变量
数据读写不只是“从磁盘拉进来再扔出去”,中间还隔着一层序列化与内存表示。Spark默认的Java序列化太慢、产物太大,生产环境基本都用Kryo。如果你要用RDD保存对象或做shuffle,记得设置spark.serializer为org.apache.spark.serializer.KryoSerializer,并且提前注册自定义类,否则Kryo每次遇到新类都会触发路径查找,性能直接打七折。
内存这块,Spark 2.x之后统一管理execution和storage内存。读进来的数据会以内存页的方式缓存在storage区,如果你调persist(MEMORY_ONLY),这些页直接复用。理解了这一点,你就明白为什么“读一遍再做多次action”比“每做一次action就重新读一次文件”要快得多——前者吃内存,后者吃磁盘I/O。内存不够的时候,可以开堆外内存,但堆外内存没法做压缩的哈希表优化,用不用得看具体场景。
2. 文件格式选型:读写性能的分水岭
2.1 Parquet为什么是默认王者
Parquet是Apache基金会下的列式存储格式,设计目标是高效压缩和高效扫描。它的核心优点有三条:
第一,列式存储。查询只需要读涉及的列,比如一张表有50个字段,你只想读其中3个,Parquet能够跳过剩余47列的数据块,这在宽表场景下收益巨大,这个特性叫列剪枝。第二,自带Schema。Parquet文件内嵌元数据,Spark读取时能推断出正确的类型,比CSV猜类型靠谱得多。第三,内置压缩与编码。默认snappy压缩,dictionary encoding对重复值多的字段压缩率非常夸张。我见过一个用户行为表,源JSON是120GB,转成Parquet且开启zstd压缩之后只有不到15GB,查询性能从分钟级降到秒级。
生产中的推荐配置:存储层用Parquet + snappy,兼顾压缩率与解压速度;如果字段重复度非常高、且对读取吞吐有极致要求,用zstd。不要为了“压缩到极致”用gzip,解压速度太慢,会让CPU成为瓶颈。
2.2 ORC与Avro:什么场景才值得换
ORC是Hive生态里非常成熟的列式格式,在Hive SQL下性能和压缩率常常超过Parquet。但要注意,Spark对ORC的支持依赖Hive的ORC SerDe,如果你用的是纯Spark作业不经过Hive,Parquet往往是更平滑的选择;如果团队有大量Hive数仓表,Spark读ORC也没问题,但要留意Spark版本与Hive版本的兼容性。我碰到过Spark 3.0加Hive 2.3的ORC读取偶发空指针的情况,最后升级Hive版本解决。
Avro是行式存储格式,带Schema演进能力,字段新增、删除对下游友好。它更适合数据落地、Kafka消息序列化、跨团队数据交换场景。缺点是列式查询性能不如Parquet,如果分析师经常做选择性查询,Avro会吃亏。实际上我在一些实时链路里,Kafka消息体导到数仓ODS层时先用Avro,DWD层再转Parquet,兼顾了写入灵活性和查询效率。
2.3 JSON/CSV的隐性成本,别让方便成为债
JSON和CSV是“能用但不划算”的典型代表。Spark官方支持这些格式,但它们的共同问题是:无法列剪枝,必须全量读取后再解析;行式存储,压缩率差;Schema推断依靠抽样,容易出错。CSV读进来的字段全部是字符串,要手动转类型;JSON则依赖路径推断,嵌套结构复杂时解析慢得离谱。
我的判断标准是这样的:数据量小于几百MB、临时分析一次性使用,用JSON/CSV无所谓;数据量达到GB级别或者每天定时产出给下游用,请立刻转Parquet。别让“方便”成为长期技术债,很多集群I/O飙升、任务OOM,源头就是躺着大量CSV和JSON。
3. 核心实操:数据读取的细节与参数调优
3.1 读Parquet/ORC时要不要手动指定Schema
读取Parquet最标准的姿势是:
df = spark.read.format("parquet").load("hdfs://nameservice/user/hive/warehouse/ods.db/orders")但有个细节很多人不知道:Spark读取Parquet时虽然会自动推断Schema,但推断过程需要读取文件尾部元数据,在大目录下会有一段额外的元数据读取开销。如果文件数量特别多,建议直接指定Schema,跳过推断:
from pyspark.sql.types import StructType, StructField, StringType, LongType schema = StructType([ StructField("order_id", StringType()), StructField("user_id", LongType()), StructField("amount", StringType()) ]) df = spark.read.schema(schema).parquet("hdfs://.../orders")指定Schema还有一个好处:避免了“原本字段是Long类型,读出来因为底层存储不一致变成Decimal”的坑。Parquet写入时会按列的原始类型编码,但如果上游用Hive改过表结构、字段以String类型落地,Spark读出来的类型会和你预期的不一样。预先声明Schema可以强制转换,省去后续cast。
列剪枝这个优化你不需要手动做,Catalyst会自动把用不到的列从Parquet读取计划里剔除。但你要知道一个反例:如果你先读全表再select,列剪枝依然生效,因为优化器会下推投影。真正不会下推的情况是你在RDD层面自己解析Parquet,或者用了一些旁路SDK绕过了DataFrame API,那才可能读到全量列。
3.2 读JSON:多行、嵌套与Schema推断的三连坑
很多人在搜“spark中读取json”,这里展开讲。读取JSON有两个常见坑:多行JSON和复杂嵌套。
多行JSON指的是整个文件是一个大JSON数组,或者每个对象跨多行。Spark默认只按行解析,遇到这种情况会报错“Each element in the array must be in a single line”。解法是加multiLine参数:
df = spark.read.option("multiLine", True).json("hdfs://.../data.json")复杂嵌套则建议先读成字符串再预处理。我经常遇到某个字段是一串JSON字符串,而不是真正的嵌套类型。这时先用spark.read.text把整行读进来,再用from_json转struct:
raw = spark.read.text("hdfs://.../nested.json") from pyspark.sql import functions as F df = raw.select(F.from_json(F.col("value"), schema).alias("data"))顺带一提,JSON的Schema推断会抽样部分行,如果文件很大,抽样仍可能触发额外的扫描开销。建议手动指定schema,否则遇到类型推断错误还得返工。
3.3 读Hive表、JDBC与Kafka的注意事项
除了文件,Spark数据读取的场景里Hive、JDBC、Kafka占大头。
读Hive表是spark.sql("select * from ods.orders"),内部走的是Hive Metastore,Spark会拿到表对应的存储路径和格式,按Hive表的SerDe配置去读。关键点在于:如果Hive表是TextFile格式且没有分桶,Spark读它会退化成一个全量扫描加反序列化过程,性能极差;把Hive表改成Parquet分区表或者用Spark写回一张Parquet表,性能完全不同。我自己优化过的一个案例,源Hive表200GB的文本格式,查询要6分钟;转成Parquet分区表后,同样查询只要45秒。
读JDBC需要注意两个参数:partitionColumn和numPartitions。Spark读MySQL/Oracle时,如果不指定分区列,只会用一个partition去读,数据量大时直接变成单点瓶颈。正确做法是:
df = spark.read.format("jdbc").option("url", "jdbc:mysql://...") .option("dbtable", "orders") .option("user", "root").option("password", "***") .option("partitionColumn", "id") .option("lowerBound", 1) .option("upperBound", 1000000) .option("numPartitions", 10) .load()读Kafka时,Spark Structured Streaming默认会给每个分区一个task,消费offset并转成DataFrame。注意value字段是二进制,需要自己from_json。这里有个生产教训:Kafka topic的partition数不要比executor的core数大太多,否则大量task调度和网络连接会拖垮作业,一般按executor core的2到3倍设置topic分区。
3.4 写数据时如何控制分区与小文件
写入是另一个大坑。很多人写完数据后不看产出,结果第二天下游任务因为小文件太多跑不动。
写数据时最常用的两个动作:partitionBy和bucketBy。partitionBy按指定列做目录分区,比如:
df.write.format("parquet").mode("overwrite").partitionBy("dt").save("hdfs://.../orders")分区列的选择直接影响读取效率。分区字段太少(比如只有dt),每个分区下可能堆几十GB数据,读取时并行度不够;分区字段太多(比如按城市加日期加小时),会产生海量小目录,元数据压力大。经验是:选择查询过滤最频繁的一到两个字段作为分区列,让每个分区的数据量在128MB到1GB之间比较合适。
小文件治理是个长期话题。Spark写文件时会为每个输出partition生成一个文件,如果你有1000个partition且每个最后只有1MB数据,那就生成了1000个1MB文件。解法有三个:
第一,写完后用coalesce或repartition控制最终输出文件数:
df.coalesce(10).write.format("parquet").save("hdfs://...")第二,开启Spark 3.x的adaptive query execution,设置spark.sql.adaptive.enabled=true和spark.sql.adaptive.coalescePartitions.enabled=true,让写入前自动合并小分区。
第三,对已经存在的小文件目录,定期跑一次合并作业,用read读全目录后再按目标大小重写。这个作业在数据仓库里要设成常态化任务,否则文件数只会越来越多。
4. 缓存与Checkpoint:让重复读取不再慢吞吞
4.1 cache的存储级别怎么选才不坑
数据读取进来后如果会被多个action反复使用,用persist缓存是标配。cache()等价于persist(MEMORY_ONLY)。但MEMORY_ONLY在数据集超过内存时会直接丢弃块并重新计算,导致雪崩式重算。这时候要看业务特征选级别。
RDD场景:MEMORY_ONLY适合数据量可控且有内存余量的场景;MEMORY_AND_DISK更稳妥,内存放不下就溢写磁盘,避免重算。DataFrame场景:默认的存储级别其实就是MEMORY_AND_DISK,而且DataFrame cache以列式内存格式存储,效果比RDD好很多。
使用persist的关键是“用完要清”。在长任务里,缓存块占用的是统一内存,不清的话后续stage可能因内存不足把缓存淘汰掉,等于白缓存。我一般在循环处理多个数据集时,处理完立即df.unpersist()。
一个很隐蔽的问题:如果你对同一个DataFrame调用两次collect,第一次没cache,第二次会重新读源文件。人们总觉得Spark“记住了”数据,其实没有,除非显式cache,否则每次action都重算血缘。这个“重算”在生产环境是性能杀手,尤其是源文件在对象存储上时,重复读的成本高得离谱。
4.2 checkpoint的正确姿势:先cache再checkpoint
Checkpoint和cache是两回事。cache把数据留在内存或磁盘,但不切断血缘;checkpoint会砍断血缘,把中间结果保存成一个物理文件。典型的用处有两个:长血缘链的容错恢复,以及循环迭代算法里防止血缘爆炸。
用法:
spark.sparkContext.setCheckpointDir("hdfs://.../checkpoint") df = df.checkpoint()注意:checkpoint会触发一次额外计算,因为是先计算再落盘。所以正确的组合往往是先cache再checkpoint——先缓存一份在内存,checkpoint再从缓存落盘,避免重新计算。这个组合在迭代场景里能省下不少时间。
还有一个细节:checkpoint的目录不能和业务表目录混在一起,否则后续清理数据时误删会出大问题;而且checkpoint会生成一堆随机UUID目录,调度系统清垃圾的时候要排除掉。
5. 常见问题与排查技巧实录
5.1 本地路径和HDFS路径分不清
刚开始用Spark的人最容易犯的错:在集群提交任务,却写了file:///data/xx.csv,结果每个executor都在自己的本地盘上找文件,不是报FileNotFound就是结果张冠李戴。排查方法很简单:用hdfs dfs -ls确认文件是否在HDFS上,路径前缀写hdfs://;如果是本地调试Spark Shell,再允许file://。
反过来也要注意:本地Spark Shell读取hdfs://nameservice/xxx时,如果未配置HA的nameservice,会报UnknownHostException,这时候要么直接用active namenode的地址,要么检查core-site.xml与hdfs-site.xml是否放到了Spark的conf目录。
5.2 小文件过多导致读取阶段卡死
现象:执行计划显示读取阶段有上万个task,但每个task才处理几十KB,整个job卡在file scan上。原因通常是上游写入未控制文件数,或者源表经历了太多次merge未压缩。
解决路径:先调大spark.sql.files.maxPartitionBytes,比如到256MB,让读取阶段按文件大小合并分区;再对下游做一轮repartition压缩任务数。长期方案是数仓任务里统一设置写入文件目标大小,我用的是“目标文件数等于数据总量除以1GB”,严格控制每个输出文件在512MB到1GB之间。
5.3 Schema推断错误引发运行期异常
CSV和JSON因为缺少强Schema,经常推断错类型。现象是读取时id列因为前100行是数字被推断为Long,第1000行出现字母,运行到一半抛NumberFormatException。
解法:读入时设置option("inferSchema", "false"),强制指定schema;或者读入后做两次cast。我给数据平台定过一个规矩:所有落地到仓库的表一律显式声明Schema,不依赖自动推断。自动推断只允许在临时探索性分析里出现。
5.4 缓存与OOM的恶性循环
很多OOM往往不是数据量真的超了,而是缓存级别选择不当。比如对一个大DataFrame执行persist(MEMORY_ONLY),内存不够就把已有块淘汰,后续action又触发重算,重算又耗内存,形成恶性循环。排查时用Spark UI的Storage标签页看缓存块的Removal次数,Removal次数高就说明存储级别要调整。
我自己遇到过的坑是堆内内存设得太大,留给task执行的内存不足,调度频繁失败。最后把spark.memory.storageFraction从默认0.5调低到0.3才稳住,默认给storage的空间收一点,execution更宽裕,任务稳定性明显提升。
6. 两个生产案例,把上面的理论串起来
6.1 200GB JSON日志批处理优化实战
某业务线每天同步一批订单JSON文件约200GB,同步后Spark需要读取并做字段规范化,之后join用户维表计算指标。改造前流程:直接读JSON全量、无分区、无缓存、默认Java序列化,每天跑批40分钟。
改造步骤:
- 将JSON文件转Parquet落地ODS层,snappy压缩,按天分区。
- 读取时指定Schema,禁止自动推断。
- join时的用户维表persist(MEMORY_AND_DISK)。
- 开启AQE,设置
spark.sql.adaptive.coalescePartitions.enabled=true。 - 序列化切换Kryo并注册类。
最终跑批时间从40分钟降至10分钟,I/O等待显著下降。这个案例说明一个问题:很多时候性能瓶颈根本不在计算逻辑,而在存储格式和读取方式上。
6.2 对象存储小文件风暴的治理记录
某次云上环境,ELT任务每小时生成2000个几百KB的Parquet文件,一周后目录下文件数接近30万。下游读取时Spark光列文件列表就花了3分钟,再拉数据又慢又贵。处理方案是每天凌晨跑一个合并任务,把所有小时文件按天分区读入,再用coalesce合并成每512MB一个文件。合并后30万文件变成不到300个,下游读取时间从15分钟降到40秒。
这个案例提醒大家,读数据不只是“打开文件”,列目录、文件数量、数据本地性都会深深影响性能。数据湖和数仓的日常运维里,一定要有专门监控文件数的自动化脚本,否则集群会被自己的小文件拖垮。
我之前操作中还发现,对象存储的list操作会随着目录深度恶化,所以建议的分区设计到天为止,不要轻易做小时级别分区。如果业务非要小时粒度,可以在表内用一个小时字段做过滤,而不是直接创建二级目录。
我自己这几年用Spark,最大的体会是:大多数任务变慢,根本不是计算逻辑的问题,而是存储与读写环节埋了雷。你去看那些跑得飞快的作业,往往都不是代码写得花哨,而是存储格式、分区、Schema、缓存这些基本功做得扎实。尤其是当你面对一个“别人留下”的作业时,先别急着改业务逻辑,把数据读写的整个过程用Spark UI捋一遍,收益往往比重构代码大得多。这篇里的每一个参数、每一个坑,都是真金白银买来的,希望能帮你少交一点学费。