1. 两个输出算子到底差在哪:从一次 HDFS 写失败说起
saveAsHadoopFile和saveAsNewAPIHadoopFile是 Spark 里把 PairRDD 落到 HDFS 的两个经典算子,都来自PairRDDFunctions。名字只差一个NewAPI,但底层走的是 Hadoop 两套完全不同的输出体系:前者基于org.apache.hadoop.mapred(老 API),后者基于org.apache.hadoop.mapreduce(新 API)。很多人在本地跑 Spark 写 HDFS 时,代码编译通过、任务也提交了,结果要么报ClassNotFoundException,要么输出目录里只有_SUCCESS没有part-*,要么压缩参数设了却不生效——根因基本都落在“用错了 API 家族”或“配置项名字对不上”。
这篇聚焦源码差异和踩坑点,同时把本地配置链路走通:用 TaoToken 统一 Key/API 通道管理模型调用凭据,避免在多个脚本里散落硬编码。适合正在写 Spark 批处理、需要把结果稳定落到 HDFS 的同学,也适合排查输出失败但日志只给一行Job aborted的场景。下面先讲清楚两个算子的源码行为,再给可复制的core-site.xml和接入骨架,最后用对比验证动作定位问题。
2. 前置准备:TaoToken 统一 Key 与本地 Spark 环境
在动算子之前,先把凭据和环境理顺。我习惯把模型调用、编码辅助这类需要 Key 的通道统一走 TaoToken,这样 Spark 脚本里不出现明文密钥,换环境只改一处。TaoToken 官网入口是 https://taotoken.net/?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= ,API 基址是 https://taotoken.net/api (这个不加 UTM)。
你需要准备的东西不多:一个可用的 Spark 3.x 本地环境(spark-shell或spark-submit都行)、一个能连的 HDFS(伪分布式即可)、以及 TaoToken 控制台里创建好的 API Key。创建 Key 的入口在 https://taotoken.net/console?utm_source=taotoken_aicg_blog_end&utm_content=console&utm_campaign=rewrite ,进去后新建一个 Key,复制出来备用。如果你后面要跑长期编码任务或 Agent 流程,可以看 Coding Plan: https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_content=coding-plan&utm_campaign=rewrite 。
环境变量建议这样设,避免写进代码:
export TAOTOKEN_API_KEY="sk-你的key" export TAOTOKEN_BASE_URL="https://taotoken.net/api"HDFS 侧确认core-site.xml里fs.defaultFS指向你的 NameNode,例如hdfs://leen:8020。这一步不对,后面两个算子都会在setOutputPath阶段直接抛IllegalArgumentException: Can not create a Path from an empty string或连接超时。
3. 可复制配置:core-site.xml 与 TaoToken 接入骨架
先给 HDFS 客户端配置。把下面内容放到$HADOOP_CONF_DIR/core-site.xml,Spark 会通过sc.hadoopConfiguration读到:
<configuration> <property> <name>fs.defaultFS</name> <value>hdfs://leen:8020</value> </property> <property> <name>hadoop.tmp.dir</name> <value>/opt/hadoop/tmp</value> </property> <property> <name>io.file.buffer.size</name> <value>131072</value> </property> </configuration>TaoToken 接入骨架用一个独立配置文件承载,Spark 侧只读环境变量,不落盘明文:
// TaoTokenConfig.scala object TaoTokenConfig { val apiKey: String = sys.env.getOrElse("TAOTOKEN_API_KEY", "") val baseUrl: String = sys.env.getOrElse("TAOTOKEN_BASE_URL", "https://taotoken.net/api") def headers: Map[String, String] = Map( "Authorization" -> s"Bearer $apiKey", "Content-Type" -> "application/json" ) }如果你要在 Spark 任务里调用模型做数据校验或生成摘要,用这个 headers 走 HTTP 即可,Key 不进入 RDD 闭包序列化,避免Task not serializable。需要看接口细节时翻接入文档: https://taotoken.net/doc?utm_source=taotoken_aicg_blog_end&utm_content=doc&utm_campaign=rewrite 。
4. 源码实例详解:saveAsHadoopFile 与 saveAsNewAPIHadoopFile
4.1 saveAsHadoopFile 的老 API 路径
saveAsHadoopFile接收的是JobConf和org.apache.hadoop.mapred.OutputFormat。源码里它做了四件事:设置outputKeyClass/outputValueClass、设置OutputFormat、处理压缩 codec、设置输出路径后调用saveAsHadoopDataset。压缩那段值得注意,它同时设了mapred.output.compress、mapred.output.compression.codec和mapred.output.compression.type,并且把CompressionType固定成BLOCK。
import org.apache.hadoop.mapred.TextOutputFormat import org.apache.hadoop.io.{Text, IntWritable} val rdd1 = sc.makeRDD(Array(("A",2),("A",1),("B",6),("B",3),("B",7)), 2) rdd1.saveAsHadoopFile( "hdfs://leen:8020/test01", classOf[Text], classOf[IntWritable], classOf[TextOutputFormat[Text, IntWritable]] )带压缩的版本多传一个 codec 类:
rdd1.saveAsHadoopFile( "hdfs://leen:8020/test02", classOf[Text], classOf[IntWritable], classOf[TextOutputFormat[Text, IntWritable]], classOf[org.apache.hadoop.io.compress.GzipCodec] )跑完用hadoop fs -du -h看,test01下是part-00000、part-00001,test02下变成part-00000.gz、part-00001.gz。分区数决定文件数,两个分区就是两个 part 文件,这点和saveAsNewAPIHadoopFile一致。
4.2 saveAsNewAPIHadoopFile 的新 API 路径
saveAsNewAPIHadoopFile用的是org.apache.hadoop.mapreduce.OutputFormat和Configuration。源码里它先NewAPIHadoopJob.getInstance(hadoopConf),再setOutputKeyClass/setOutputValueClass/setOutputFormatClass,最后把路径写进mapred.output.dir后调saveAsNewAPIHadoopDataset。关键差异:它没有 codec 参数,压缩必须自己在Configuration里设。
import org.apache.hadoop.mapreduce.lib.output.TextOutputFormat import org.apache.hadoop.io.{Text, IntWritable} val rdd1 = sc.makeRDD(Array(("A",2),("A",1),("B",6),("B",3),("B",7)), 2) rdd1.saveAsNewAPIHadoopFile( "hdfs://leen:8020/test03", classOf[Text], classOf[IntWritable], classOf[TextOutputFormat[Text, IntWritable]] )要压缩就改hadoopConf:
val hadoopConf = sc.hadoopConfiguration hadoopConf.set("mapred.output.compress", "true") hadoopConf.set("mapred.output.compression.codec", "org.apache.hadoop.io.compress.GzipCodec") hadoopConf.set("mapred.output.compression.type", "BLOCK") rdd1.saveAsNewAPIHadoopFile( "hdfs://leen:8020/test04", classOf[Text], classOf[IntWritable], classOf[TextOutputFormat[Text, IntWritable]], hadoopConf )对比表如下,方便你按场景选:
| 维度 | saveAsHadoopFile | saveAsNewAPIHadoopFile |
|---|---|---|
| API 家族 | org.apache.hadoop.mapred | org.apache.hadoop.mapreduce |
| 配置对象 | JobConf | Configuration |
| 压缩参数 | 方法参数传 codec | 手动 set 到 Configuration |
| OutputFormat | mapred.OutputFormat | mapreduce.OutputFormat |
| 路径设置 | FileOutputFormat.setOutputPath | mapred.output.dir |
5. 验证请求与成功结果
写完先别急着上集群,本地spark-shell跑一遍最快。启动后把上面rdd1和两个算子分别执行,然后验证:
hadoop fs -ls hdfs://leen:8020/test01 hadoop fs -ls hdfs://leen:8020/test02 hadoop fs -ls hdfs://leen:8020/test03 hadoop fs -ls hdfs://leen:8020/test04成功标志是每个目录下有_SUCCESS加若干part-*。再看内容:
hadoop fs -cat hdfs://leen:8020/test01/part-00000应该输出A 2这类 key-value 文本。压缩目录用hadoop fs -cat会看到乱码,这是正常的,用hadoop fs -text能自动解压:
hadoop fs -text hdfs://leen:8020/test02/part-00000.gz如果你在任务里接了 TaoToken 做辅助校验,验证模型通道是否通,可以直接用模型对话入口测一次: https://taotoken.net/model-chat?utm_source=taotoken_aicg_blog_end&utm_content=model-chat&utm_campaign=rewrite 。返回正常说明 Key 和 base URL 都对,再回到 Spark 侧排查输出问题就不会混淆两件事。
6. 本篇常见错排查
报ClassNotFoundException: org.apache.hadoop.mapred.TextOutputFormat:说明你用了新 API 的 OutputFormat 类去配老算子,或者反过来。检查 import:老算子配mapred.TextOutputFormat,新算子配mapreduce.lib.output.TextOutputFormat,两者不能混。
输出目录只有_SUCCESS没有 part 文件:多半是 RDD 为空,或者outputKeyClass/outputValueClass设成了接口类型。源码里这两个类会写进 JobConf,Hadoop 序列化时需要具体类。用rdd1.count()先确认数据量。
压缩设了不生效:saveAsNewAPIHadoopFile没有 codec 参数,必须在Configuration里设mapred.output.compress=true,且 codec 类名要写全限定名。老算子如果传了 codec 但没生效,检查是不是传成了None。
Task not serializable:把 TaoToken 的 Key 或 HTTP client 放进了 RDD 闭包。正确做法是只在 Driver 侧读环境变量,闭包内不引用不可序列化对象。
路径报Can not create a Path from an empty string:fs.defaultFS没配或core-site.xml没被 Spark 读到。确认sc.hadoopConfiguration.get("fs.defaultFS")有值。
推测执行导致数据丢失警告:源码里对Direct输出提交器有警告。如果你开了spark.speculation=true且用了自定义 committer,换成FileOutputCommitter更稳。
排障时如果怀疑是凭据通道问题,去 API Keys 页面核对 Key 状态: https://taotoken.net/api-keys?utm_source=taotoken_aicg_blog_end&utm_content=api-keys&utm_campaign=rewrite 。接入细节对不上就看文档: https://taotoken.net/doc?utm_source=taotoken_aicg_blog_end&utm_content=doc&utm_campaign=rewrite 。长期跑编码和 Agent 任务的话,Coding Plan 入口在 https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_content=coding-plan&utm_campaign=rewrite ,Claude Code 相关配置参考 https://taotoken.net/claude-code-anthropic?utm_source=taotoken_aicg_blog_end&utm_content=claude-code-anthropic&utm_campaign=rewrite 。
最后留一个我踩过的坑:saveAsHadoopFile的压缩参数里CompressionType被源码写死成BLOCK,如果你业务上需要RECORD级别压缩,老算子改不了,得换新算子自己设mapred.output.compression.type。这个差异在源码里不显眼,但线上排查时很关键。