资讯动态

Spark数据存取底层逻辑与读写调优:从文件格式到分区裁剪

发布时间:2026/10/9 3:35:53 来源:尧图企业网站定制
1. 先还原一次读数据慢的排查Spark存取的底层逻辑1.1 那个26分钟的作业问题出在读的姿势前阵子帮同事排查一个离线数仓任务作业本身逻辑非常简单从Parquet表读订单明细过滤最近7天数据按业务域聚合后写回结果表。可这个任务每次运行都要26分钟比其他同等规模的任务慢了一个量级。我打开Spark UIStage 0的Input拉到了2.3GB任务数小三千个。而真正要算的近7天分区数据量只有300MB出头。问题不在计算在于读取方式太粗放了——他直接把整张表的根目录作为数据源加载Spark只能把全部分区先扫描一遍然后再用过滤条件把不用的分区丢掉。最便宜的优化被他漏掉了读之前先做分区裁剪。这类问题在Spark数据存取里非常典型。很多人把精力放在调shuffle、调内存却忽略了一个事实Spark读数据的路径上从文件列表拉取、块扫描、任务切分到数据解码每个环节都可能成为瓶颈。存储与读取从来不是给它一个路径就能跑那么简单它决定了任务的下限。1.2 读写的三条路径与各自的入口Spark的数据存取大体分三类多数人日常只用到其中一类面试或做方案时容易漏掉另外两类。数据源类型读取入口写入入口典型场景文件系统HDFS、S3、本地spark.read.parquet/json/csv、sc.textFiledf.write.format(...).save()、rdd.saveAsTextFile离线ETL、日志解析、数据湖原始层表Hive Metastore / Spark Catalogspark.table(db.tbl)、spark.sql(select ...)df.write.saveAsTable、INSERT INTO数仓建模、统一元数据管理外部系统JDBC、Kafka、NoSQLspark.read.format(jdbc)、spark.readStream.format(kafka)df.write.jdbc(...)、writeStream.format(kafka)业务库抽取、实时入仓、维表关联传统写法里SparkContext.textFile读文本文件很直观但新项目我基本都建议直接用SparkSession统一入口——它内部维护了SparkContext同时把DataFrameReader、Streaming等读入口都收拢在一起代码更干净也避免API混用带来的序列化和类型问题。1.3 文件与分区数据是如何变成并行度的Spark读取文件时会把目标路径下的文件列表拉出来再按照一定的规则切成若干个Partition。每个Partition对应一个Task由Executor上的一个核心去执行。存储的物理形态直接决定了任务的并行度。这里有个容易忽略的小文件合并机制。Spark读文件时并不是一个文件一个分区而是参考spark.sql.files.maxPartitionBytes默认128MB和spark.sql.files.openCostInBytes默认4MB来做分组读取每个文件都有固定的打开成本因此一堆几个MB的小文件会被合并成较大的分区避免生成成千上万个Task。反过来如果文件超过128MB就可能被拆成多个分区来读。写入端同理——最终输出文件数量由最后一个Stage的分区数决定。一个阶段有200个分区写出就会产生200个文件。这个对称关系是所有小文件问题的根源后面专门展开。2. RDD、DataFrame与Dataset三种API的存取方式完全不同2.1 从API设计看三种抽象很多入门者把RDD、DataFrame、Dataset当成三种都能存能读的API这没错但它们的设计目标差异很大选错会让代码又啰嗦又慢。维度RDDDataFrameDataset是否带Schema否是是类型安全运行时运行时编译期Catalyst优化不参与完整参与完整参与典型序列化方式Java/KryoTungsten二进制行编码器Encoder适用场景非结构化数据、底层自定义算子绝大多数离线分析复杂业务逻辑、需要强类型校验RDD的核心是函数式变换map、flatMap、reduceByKey这些算子很好用但没有字段名、没有类型推断Spark的Catalyst优化器看到RDD基本无从下手。DataFrame本质是带Schema的分布式行集合优化器可以做谓词下推、列裁剪、常量折叠。Dataset则是强类型版本你可以把DataFrame转成Dataset[CaseClass]写代码时编译器能帮你查错。我的建议很直接日常ETL和分析优先DataFrame业务模型复杂、字段易变时用DatasetRDD只用于读取原始的非结构化数据或者实现自定义数据源时兜底。2.2 同一份JSON用三种方式读一次拿最常见的JSON日志来对比。假设HDFS上有一批事件日志每条记录有dt、api、level、cost四个字段。RDD方式val rdd sc.textFile(hdfs:///data/events/20240101/*.json) val parsed rdd.map { line // 自己解析JSON没有Schema推断也没有类型安全 // 通常引入第三方JSON库或手写处理 (extractField(line, api), extractField(line, cost).toLong) }RDD读取只负责把文本行拉到内存后续解析逻辑全得自己写字段类型也只能自己强转。如果日志字段增加代码得跟着改还容易在运行期抛异常。DataFrame方式val df spark.read.json(hdfs:///data/events/20240101) df.filter($level ERROR) .groupBy(api) .agg(sum(cost))读进来就有Schema字段类型自动推断groupBy(api)这些操作走Catalyst优化。同样的逻辑代码量差一个量级。Dataset方式case class Event(dt: String, api: String, level: String, cost: Long) val ds spark.read.json(hdfs:///data/events/20240101).as[Event] ds.filter(_.level ERROR) .groupByKey(_.api) .agg(sum(_.cost))可以像操作本地集合一样写_.level编译期就会校验字段名。但要注意as[CaseClass]的转换过程中会经过Encoder序列化复杂嵌套类型的性能不一定比DataFrame的二进制行好所以不要为了强类型而牺牲全部性能。2.3 写出的两种模型save与write的语义差异RDD的写出口比较老派saveAsTextFile、saveAsObjectFile、saveAsSequenceFile。其中saveAsObjectFile我强烈不建议在生产使用——它依赖Java序列化类结构一变就废下游还没法用其他工具直接读。DataFrame的写出口统一是df.write配合SaveMode控制写入语义ErrorIfExists目标已存在就报错默认行为防误写最安全Overwrite覆盖写HDFS上会先删目录再写Append追加写常用于增量入仓Ignore目标已存在则静默跳过Overwrite看起来很省事实际很危险。如果你不小心把路径写成某张表的根目录它会先把整个目录删掉再写新的数据恢复基本没戏。我在生产环境只用Append和显式指定子目录的Overwrite根目录的覆盖操作必须二次确认。3. 文件格式与压缩从JSON到Parquet的选型与参数3.1 列式存储为什么成为生产首选处理大数据的人应该都有体会生产分析表几乎不会用JSON或CSV存清一色Parquet或ORC。原因是列式存储有三个直接红利。第一是列裁剪。一张订单表50个字段你只需要region和amount两列聚合。Parquet按列组织数据读取时可以只把涉及到的列块读出来行式存储却要把整行都搬到内存再丢弃不需要的字段。列越多差距越大。第二是谓词下推。Parquet文件内部按行组Row Group组织每个行组的列块都带有min/max统计信息。Spark扫描时发现某个行组的dt范围与过滤条件不匹配整个行组直接跳过IO省一大截。这种下推在行式存储里做不到这么细。第三是压缩比。同一列的数据类型一致、取值规律性强压缩算法发挥空间大。枚举值字段、时间戳字段压缩下来往往只有原始大小的三分之一甚至更低。所以我的归档原则是能用Parquet绝不用文本格式ORC看周边生态Avro留给明确的Schema演进场景。3.2 各格式读取API的常用option// CSV第一行表头自动推断类型 spark.read.option(header, true) .option(inferSchema, true) .csv(hdfs:///data/csv/orders) // JSON开启multiline允许一条记录跨多行 spark.read.option(multiline, true) .json(hdfs:///data/events) // Parquet合并多个文件的schema差异 spark.read.option(mergeSchema, true) .parquet(hdfs:///data/parquet/orders) // ORC与Avro spark.read.format(orc).load(hdfs:///data/orc/orders) spark.read.format(avro).load(hdfs:///data/avro/events)几个option要特别注意。CSV默认inferSchemafalse不开启的话所有列都是StringType后续sum、avg全得手动cast。但开启后Spark会额外做一轮采样推断对超大文件带来一次额外的读取开销所以更推荐直接用schema参数手工定义字段类型。JSON的multiline只适用于标准JSON对象跨行的情况代价是Spark需要把一个文件的完整内容读入内存再解析超大JSON文件慎开。Parquet的mergeSchema解决的是同一目录下不同文件Schema不一致的问题开启后读取时会合并所有文件的元数据。这个能力好用但有开销不是每个任务都值得开具体见后面的故障案例。3.3 压缩编码器怎么选文件格式定好后压缩编码器是第二步。常见选择如下编码器压缩比速度是否可分割生产建议gzip高中文本格式下不可分割归档冷数据bzip2最高慢可分割极少用lzo中快需要索引老Hadoop生态snappy中很快可分割生产默认首选zstd高较快可分割追求压缩比时选它可分割这一点容易被忽略。对于文本格式的gzip压缩文件单个文件只能由一个Task读取如果一个1GB的日志文件gzip压缩后变成200MB最终只有一个Task在处理这个文件并行度直接崩溃。snappy和zstd没有这个问题。Parquet场景下内部按行组和page组织数据配合snappy或zstd是主流。配置方式// 全局配置 spark.conf.set(spark.sql.parquet.compression.codec, zstd) // 或单次写入指定 df.write.option(compression, zstd).parquet(hdfs:///data/out)我的习惯是在线分析链路用Parquet加snappy稳定且快离线归档冷数据用zstd压缩比高节省存储成本。3.4 JSON读取的两个高频坑JSON虽然在生产存储里不受待见但它是日志和接口数据最常见的原始格式读取坑也最多。第一个坑是multiline。很多人从接口平台下载的JSON文件都是pretty打印的每条记录占好几行。直接用spark.read.json读会把每一行当成一个独立JSON对象去解析结果要么报错要么读出来的字段全是null。解决方法是加option(multiline, true)但前面说了大文件要谨慎。第二个坑是字符编码。Spark的文本类数据源默认按UTF-8处理遇到GBK等编码的文件读出来直接乱码。基本没有直接在Spark里优雅处理GBK的办法我的做法是在数据源头规范编码或者先用转换工具统一成UTF-8再入Spark。另外JSON的Schema推断不是免费的。格式复杂、嵌套深的JSON会对读取性能有明显影响字段类型不稳定时还会推断出错。生产环境里如果JSON只是中间态我会尽快在读取时用.option(samplingRatio, 0.1)或手动指定schema避免每次启动都在元数据上耗时。4. JDBC、Kafka等外部数据源的读写要点4.1 JDBC分区读参数与全表扫描问题从关系型数据库抽数最怕的就是写一个没有分区策略的spark.read.jdbc。默认情况下Spark会用一个Task去执行整条SQL数据量大时数据库压力大、Spark侧并行度又是0。正确写法是显式指定分区参数val df spark.read.format(jdbc) .option(url, jdbc:mysql://...) .option(dbtable, (select id, name, update_time from users where update_time 2024-01-01) t) .option(user, etl_user) .option(password, ***) .option(partitionColumn, id) .option(lowerBound, 1) .option(upperBound, 10000000) .option(numPartitions, 8) .option(fetchsize, 1000) .load()底层逻辑是Spark按照lowerBound到upperBound的区间结合numPartitions生成多个子查询每个Task各查一段。这里有两个非常实际的注意点。第一partitionColumn必须是数值或时间类型字符串类型直接报错。第二这个参数只负责把任务切分均匀本身不做过滤真正的数据裁剪要写到dbtable子查询的where条件里否则全表数据还是会被拉出来。写回数据库同理。df.write.jdbc时可以设置batchsize默认1000和isolationLevel。大表写入建议把batchsize调到5000到10000但不要无脑调大数据库端事务和网络带宽都会成为瓶颈。4.2 与Kafka对接从offset到checkpointKafka接入在实时链路里是标配。读取端的标准姿势val stream spark.readStream .format(kafka) .option(kafka.bootstrap.servers, broker1:9092,broker2:9092) .option(subscribe, ods_user_behavior) .option(startingOffsets, earliest) .load()这里的startingOffsets只在首次从无checkpoint状态启动时生效。一旦你配置了checkpointLocationoffset就由checkpoint里的长期记录接管重启后不会丢数据也不会重复大量消费。写入端大多数人会写错字段名stream.selectExpr( cast(key as string) as key, cast(value as string) as value ).writeStream .format(kafka) .option(kafka.bootstrap.servers, broker1:9092) .option(topic, dws_user_behavior) .option(checkpointLocation, /spark/checkpoint/user_behavior) .outputMode(append) .start()Kafka sink要求的输出列必须叫key和value很多新手直接输出业务字段名启动就报错。另外checkpoint目录千万别乱删删了等于丢失offset记录重启后会按startingOffsets重新消费可能造成重复或丢失。4.3 写到NoSQL与消息系统时的缓冲控制HBase、Elasticsearch这类系统批量写入时缓冲区和并发数的控制很关键。写HBase如果用逐条put方式吞吐很难看正确思路是按Region分布做预分区配合批量Buffer写入。比如对DataFrame做一次repartition按行键前缀打散让每个Task写的数据尽量落在连续的Region上减少跨Region的随机写。Elasticsearch则要关注batch.size.bytes和batch.retry.count一次写入的文档数过多会压垮ES集群。我的经验是从小批量开始压测逐步加大到集群响应延迟可接受的阈值不要照搬网上的推荐值。5. 缓存、分区粒度与落盘调优5.1 什么时候cache有用什么时候是纯负担cache和persist是最容易被用错的API。很多人不管什么数据都先cache一下结果任务反而变慢。StorageLevel的选择存储级别含义适用场景MEMORY_ONLY只放内存不序列化小表高频复用MEMORY_ONLY_SER内存放序列化对象省空间大对象复用GC压力大时MEMORY_AND_DISK内存放不下落磁盘中等数据量复用MEMORY_AND_DISK_SER序列化存储内存优先数据复用频繁但空间紧张DISK_ONLY全落磁盘内存完全不够时兜底真正的使用场景是同一份DataFrame会被多个Action复用且复用次数大于一次——比如迭代算法、多轮join、反复进行不同维度的聚合。我见过一个ML训练任务每轮迭代都重新读源表和做特征拼接数据量5GB内存够但重算代价很大。加上.cache()后三轮迭代从40分钟降到12分钟。反过来的例子也很多一份数据读出来只是一个Action用到比如过滤完直接写结果缓存反而增加了序列化和GC成本。尤其是内存紧张时缓存的大块数据可能挤掉shuffle过程中需要的内存导致频繁spill到磁盘整体更慢。我的规则是先确认Action次数再决定是否缓存缓存用完了记得unpersist()不要一直占着内存等GC。5.2 写文件的数量由什么决定coalesce与repartition写过Spark作业的人十有八九遇到过任务运行没问题结果一看输出目录几千个小文件每个几百KB。这个现象几乎都是同一个原因——写出前没有控制分区数。写入文件数等于最后一个Stage的分区数。比如上游某次shuffle之后产生了默认200个分区你直接df.write就会输出200个文件。如果一个分区里的数据量只有几MB那就妥妥是小文件。减少分区用coalesce增加或重新分布用repartition。coalesce尽量不触发shuffle只是合并相邻分区适合从200降到50这类收拢操作repartition是全量shuffle按指定列或指定分区数重新打散适合数据分布不均时使用。如果目标是Hive分区表我一般会这样处理df.repartition(col(dt)) .write .mode(overwrite) .partitionBy(dt) .parquet(/warehouse/dws/orders)这里的思路是按写出的分区字段做一次repartition让同一个dt的数据尽量落在同一个Task里每个分区目录只产生少量文件。同时配合spark.sql.adaptive.coalescePartitions.enabledtrue让Spark在shuffle后自动合并过小的分区从源头减少小文件数量。如果表已经写坏了事后补救可以用Hive的ALTER TABLE table_name CONCATENATE它能把分区内多个小文件合并成更大的文件但一次只处理一个分区数据量大的话比较慢。5.3 分区裁剪与文件级过滤读分区表的第一个目标就是减少扫描量。Hive风格分区表在HDFS上的目录结构是dt2024-01-01/regioncn/xxx.parquetSpark读取时如果能命中where条件会直接从目录层面把无关分区过滤掉。select count(*) from dws_orders where dt 2024-01-01这条SQL里Spark不需要扫描其他日期的目录这就是分区裁剪。判断有没有生效可以用explain看执行计划找PartitionFilters和PushedFilters这两段explain select count(*) from dws_orders where dt 2024-01-01执行计划里出现PartitionFilters: [isnotnull(dt#123), (dt#123 2024-01-01)]就是裁剪成功。有一个低级错误要提醒不要在分区字段上套函数。比如where substr(dt, 1, 7) 2024-01Spark没法在目录层面做等值匹配只能把全部分区拉出来再过滤裁剪直接失效。分区粒度也要把握好。按天分区是多数离线数仓的默认选择按小时分区适合数据量极大且查询实时性要求高的场景但分区数膨胀后Hive Metastore的元数据请求都会变慢文件碎片化问题也会加剧。分区不是越细越好。6. 三个真实故障的完整排查链路6.1 多行JSON读出一堆null有一次同事报障某个JSON文件用Python打开完全正常但Spark读出来全是null。他把文件发给我看是pretty打印的格式每条记录占据了五行数组字段还跨行。根因很明确——Spark的JSON数据源默认把每一行当作一个独立JSON对象处理遇到这种一个对象跨多行的文件单行解析必然失败不报错就算好的更多时候是静默丢数据或产出null字段。排查链路是这样的先用spark.read.json读同一个路径.printSchema()看推断出的字段名发现字段全对但值为null再去看原始文件的物理换行结构确认是pretty格式后加上multilinetrue重新读取。val df spark.read.option(multiline, true).json(hdfs:///data/events/access.log)修完之后数据正常。但我也跟同事强调multiline模式会把整个文件读入内存做解析如果单个JSON文件超过几个GB宁可写个预处理脚本把JSON改成一行一条也不要硬开。6.2 Parquet新旧schema合并冲突另一个案例上游数据团队给订单表新增了一个coupon_amount字段之后新任务往同一张表的Parquet目录里写数据。结果下游旧任务读取时报错或者新字段在旧文件里读出来全是null。原因是Parquet读取时默认mergeSchemafalse。同一个目录下旧文件没有新字段新文件有新字段Spark扫描时发现schema不一致就会按旧schema读读不到新字段自然给null极端情况下直接报类型冲突。排查时我先确认了不是字段名拼写问题然后定位到新旧文件的schema差异最后在读取端开启schema合并val df spark.read .option(mergeSchema, true) .parquet(/warehouse/dws/orders)这个方案有效但代价是读取时要额外扫描文件的footer元数据做全量schema合并对元数据量大的目录有明显开销。所以我的实际建议是两段式紧急修复用mergeSchema解决长期方案是统一表结构变更流程让下游任务的schema提前对齐而不是靠读取端反复合并。6.3 小文件爆炸的前因后果还有一个高频故障某张Hive表越写越慢查询启动就要花十几秒点开Spark UI发现密密麻麻几千个TaskInput数据总量却不到2GB。典型的链路是这样的——上游任务在shuffle后没有控制分区数量默认200个分区再叠加动态分区写入目标表有500个分区每个分区都可能由多个Task写入最后产出一万多个小文件。读的时候Spark虽然会按openCostInBytes合并小文件但文件数量太多元数据拉取和任务调度本身就成了瓶颈。修复分为两步。第一步治标对已有的小文件目录做一次重写合并spark.read.parquet(/warehouse/dws/orders) .repartition(col(dt)) .write .mode(overwrite) .partitionBy(dt) .parquet(/warehouse/dws/orders_tmp)第二步治本在写任务的源头调整并行度。spark.sql.shuffle.partitions默认200只适合中小数据量大表要结合数据量估算让每个shuffle分区约有64MB到128MB的数据同时打开自适应分区合并让Spark自动收拢过小的分区。这件事给所有人的教训是小文件不是某一天突然出现的而是每一次shuffle、每一次动态分区写入都在积累。在写任务里提前控制分区数远比事后清理成本低。7. 被实测验证过的存取习惯与可复用模板7.1 一套常见的生产读写模板最后给出一套我自己反复在用的模板你可以直接抄再按业务调整val spark SparkSession.builder() .appName(dws_etl_template) .enableHiveSupport() .config(spark.sql.adaptive.enabled, true) .config(spark.sql.adaptive.coalescePartitions.enabled, true) .config(spark.sql.parquet.compression.codec, zstd) .getOrCreate() // 1. 读取能用分区裁剪就用分区裁剪别拉全表 val df spark.read.format(parquet) .load(/warehouse/ods_orders) .where(dt between 2024-01-01 and 2024-01-31) // 2. 处理只有被复用多次的结果才cache用完要unpersist val agg df .filter($status PAID) .groupBy(dt, region) .agg(sum(amount).as(sales_amount)) .cache() // 3. 写出按目标查询模式分区控制文件大小 agg.repartition(col(dt)) .write.mode(overwrite) .format(parquet) .partitionBy(dt, region) .saveAsTable(dws_orders_sales)有几个点需要解释。repartition(col(dt))放在写之前是为了让每个dt的数据尽量聚到同一个Task里避免一个分区被几十个Task各写一块。如果数据量特别大单个dt下还是会有多个文件这没关系只要每个文件保持在64MB到256MB的量级读取效率就能接受。saveAsTable会同时写数据和元数据省去手动建Hive表的步骤。但要注意它默认在Spark自带的Hive Metastore里建表如果你的数仓已经有统一元数据服务确认配置指向同一个Metastore再用。7.2 我保留的几个习惯踩过足够多的坑之后我现在做Spark存取的决策已经变成条件反射了。第一个习惯不管读什么数据先确认Input规模和文件数量。看Spark UI里的Storage和Task分布再决定要不要加分区过滤、要不要先合并小文件。大多数任务慢其实都是读的姿势不对而不是计算逻辑有问题。第二个习惯生产存储格式几乎只用Parquet或ORC压缩用snappy或zstd。JSON和CSV只用于临时排查和对外交换绝不让它们成为数仓的主要存储格式。第三个习惯写之前先算分区数。数据总量除以期望的单文件大小得出目标分区数再决定用coalesce还是repartition。绝不依赖默认的200个shuffle分区。第四个习惯面试或者写方案时如果有人问Spark的数据存储与读取方式我一般从三个层面回答——文件系统、表、外部数据源这三种入口RDD/DataFrame/Dataset三种抽象的适用差异以及分区裁剪、列式存储、压缩编码这些决定性能的物理因素。这样答基本不会冷场也确实是日常干活最重要的几个维度。存取是Spark所有计算的地基地基没打稳上面跑再好的业务逻辑都是白搭。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑