1. 先还原一次"读数据慢"的排查:Spark存取的底层逻辑
1.1 那个26分钟的作业,问题出在读的姿势
前阵子帮同事排查一个离线数仓任务,作业本身逻辑非常简单:从Parquet表读订单明细,过滤最近7天数据,按业务域聚合后写回结果表。可这个任务每次运行都要26分钟,比其他同等规模的任务慢了一个量级。
我打开Spark UI,Stage 0的Input拉到了2.3GB,任务数小三千个。而真正要算的近7天分区,数据量只有300MB出头。问题不在计算,在于读取方式太粗放了——他直接把整张表的根目录作为数据源加载,Spark只能把全部分区先扫描一遍,然后再用过滤条件把不用的分区丢掉。最便宜的优化被他漏掉了:读之前先做分区裁剪。
这类问题在Spark数据存取里非常典型。很多人把精力放在调shuffle、调内存,却忽略了一个事实:Spark读数据的路径上,从文件列表拉取、块扫描、任务切分到数据解码,每个环节都可能成为瓶颈。存储与读取从来不是"给它一个路径就能跑"那么简单,它决定了任务的下限。
1.2 读写的三条路径与各自的入口
Spark的数据存取大体分三类,多数人日常只用到其中一类,面试或做方案时容易漏掉另外两类。
| 数据源类型 | 读取入口 | 写入入口 | 典型场景 |
|---|---|---|---|
| 文件系统(HDFS、S3、本地) | spark.read.parquet/json/csv、sc.textFile | df.write.format(...).save()、rdd.saveAsTextFile | 离线ETL、日志解析、数据湖原始层 |
| 表(Hive Metastore / Spark Catalog) | spark.table("db.tbl")、spark.sql("select ...") | df.write.saveAsTable、INSERT INTO | 数仓建模、统一元数据管理 |
| 外部系统(JDBC、Kafka、NoSQL) | spark.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",这没错,但它们的设计目标差异很大,选错会让代码又啰嗦又慢。
| 维度 | RDD | DataFrame | Dataset |
|---|---|---|---|
| 是否带Schema | 否 | 是 | 是 |
| 类型安全 | 运行时 | 运行时 | 编译期 |
| Catalyst优化 | 不参与 | 完整参与 | 完整参与 |
| 典型序列化方式 | Java/Kryo | Tungsten二进制行 | 编码器(Encoder) |
| 适用场景 | 非结构化数据、底层自定义算子 | 绝大多数离线分析 | 复杂业务逻辑、需要强类型校验 |
RDD的核心是函数式变换,map、flatMap、reduceByKey这些算子很好用,但没有字段名、没有类型推断,Spark的Catalyst优化器看到RDD基本无从下手。DataFrame本质是"带Schema的分布式行集合",优化器可以做谓词下推、列裁剪、常量折叠。Dataset则是强类型版本,你可以把DataFrame转成Dataset[CaseClass],写代码时编译器能帮你查错。
我的建议很直接:日常ETL和分析优先DataFrame;业务模型复杂、字段易变时用Dataset;RDD只用于读取原始的非结构化数据,或者实现自定义数据源时兜底。
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默认inferSchema=false,不开启的话所有列都是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到checkpoint
Kafka接入在实时链路里是标配。读取端的标准姿势:
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状态启动时生效。一旦你配置了checkpointLocation,offset就由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.enabled=true,让Spark在shuffle后自动合并过小的分区,从源头减少小文件数量。
如果表已经写坏了,事后补救可以用Hive的ALTER TABLE table_name CONCATENATE,它能把分区内多个小文件合并成更大的文件,但一次只处理一个分区,数据量大的话比较慢。
5.3 分区裁剪与文件级过滤
读分区表的第一个目标就是减少扫描量。Hive风格分区表在HDFS上的目录结构是dt=2024-01-01/region=cn/xxx.parquet,Spark读取时如果能命中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-01',Spark没法在目录层面做等值匹配,只能把全部分区拉出来再过滤,裁剪直接失效。
分区粒度也要把握好。按天分区是多数离线数仓的默认选择,按小时分区适合数据量极大且查询实时性要求高的场景,但分区数膨胀后Hive Metastore的元数据请求都会变慢,文件碎片化问题也会加剧。分区不是越细越好。
6. 三个真实故障的完整排查链路
6.1 多行JSON读出一堆null
有一次同事报障:某个JSON文件用Python打开完全正常,但Spark读出来全是null。他把文件发给我看,是pretty打印的格式,每条记录占据了五行,数组字段还跨行。
根因很明确——Spark的JSON数据源默认把每一行当作一个独立JSON对象处理,遇到这种"一个对象跨多行"的文件,单行解析必然失败,不报错就算好的,更多时候是静默丢数据或产出null字段。
排查链路是这样的:先用spark.read.json读同一个路径,.printSchema()看推断出的字段名;发现字段全对但值为null,再去看原始文件的物理换行结构;确认是pretty格式后,加上multiline=true重新读取。
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读取时默认mergeSchema=false。同一个目录下,旧文件没有新字段,新文件有新字段,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发现密密麻麻几千个Task,Input数据总量却不到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所有计算的地基,地基没打稳,上面跑再好的业务逻辑都是白搭。