首页 / 资讯中心 / 文章详情

Spark数据存储与读取优化实战:从格式选型到分区调优

Spark数据存储与读取优化实战:从格式选型到分区调优 ★ FEATURED ARTICLE
在Spark里做数据存储和读取说简单也简单df.read.parquet、df.write.csvAPI一行就完事了说复杂也复杂格式怎么选、分区怎么定、小文件怎么治、谓词下推为什么没生效全是坑。我最早接手一个数据平台的时候每天凌晨跑批一个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作业不经过HiveParquet往往是更平滑的选择如果团队有大量Hive数仓表Spark读ORC也没问题但要留意Spark版本与Hive版本的兼容性。我碰到过Spark 3.0加Hive 2.3的ORC读取偶发空指针的情况最后升级Hive版本解决。Avro是行式存储格式带Schema演进能力字段新增、删除对下游友好。它更适合数据落地、Kafka消息序列化、跨团队数据交换场景。缺点是列式查询性能不如Parquet如果分析师经常做选择性查询Avro会吃亏。实际上我在一些实时链路里Kafka消息体导到数仓ODS层时先用AvroDWD层再转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转structraw 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 MetastoreSpark会拿到表对应的存储路径和格式按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.enabledtrue和spark.sql.adaptive.coalescePartitions.enabledtrue让写入前自动合并小分区。第三对已经存在的小文件目录定期跑一次合并作业用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再checkpointCheckpoint和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.enabledtrue。序列化切换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捋一遍收益往往比重构代码大得多。这篇里的每一个参数、每一个坑都是真金白银买来的希望能帮你少交一点学费。
阅读完成 · 觉得有帮助?
咨询建站