拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Spark 用 Hadoop API 读写文件:TaoToken 统一 Key 通道下的本地联调大纲

Spark 用 Hadoop API 读写文件:TaoToken 统一 Key 通道下的本地联调大纲

1. Spark 调用 Hadoop API 读写文件到底难在哪

Spark 调用 Hadoop FileSystem API 读写文件,说白了就是让 Spark 程序直接复用 Hadoop 那套FileSystem、Configuration、InputFormat/OutputFormat的能力,去操作 HDFS 或者本地文件系统。它适合谁?适合已经有一份 Hadoop 集群配置、又想在 Spark 里做本地联调的人;也适合那些不想每次都hadoop fs -put上传、想让 Spark 任务自己把结果落到指定目录的开发者。核心检索词就三个:spark、hadoopAPI、读写文件。

很多人第一次写这块代码,卡点往往不在 Spark 本身,而在两套 API 的混用。Hadoop 有老的org.apache.hadoop.mapred和新的org.apache.hadoop.mapreduce两套包,Spark 对应提供了saveAsHadoopFile和saveAsNewAPIHadoopFile、hadoopFile和newAPIHadoopFile。名字差一个 NewAPI,底层引用的OutputFormat类路径就完全不同。老 API 用org.apache.hadoop.mapred.TextOutputFormat,新 API 用org.apache.hadoop.mapreduce.lib.output.TextOutputFormat。写错一个 import,编译能过,运行时报ClassNotFoundException或者干脆序列化失败。

另一个高频坑是序列化。用hadoopFile读 HDFS 上的文件时,如果SparkConf里没注册 Hadoop 的Writable类型,就会抛java.io.NotSerializableException: org.apache.hadoop.io.Text。原因是Text、IntWritable、LongWritable这些类型默认没走 Kryo 序列化,Spark 在 shuffle 或者 collect 阶段需要把它们序列化传输时就炸了。解决办法是在SparkConf里设置spark.serializer为KryoSerializer,再通过spark.kryoserializer.classesToRegister把用到的类型登记进去。

这篇就按本地联调的完整链路走一遍:先准备core-site.xml和 Spark 依赖,再写可复制的SparkConf与Hadoop Configuration,然后用hadoop fs -ls和spark-submit验证读写结果,最后把常见报错逐个排掉。中间会用到 TaoToken 的统一 Key 通道来做模型侧的辅助联调,但文件读写本身是纯 Hadoop/Spark 的事,两者不冲突。

2. TaoToken 统一 Key 通道的前置准备

在动手写文件读写之前,先把 TaoToken 这条通道准备好。它的作用是给你一个统一的 API Key,让你在本地联调时不用为每个模型单独配一套凭证。官网入口是 https://taotoken.net/?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= ,API 根地址是 https://taotoken.net/api ,注意这个地址后面不加 UTM 参数。

你需要拿到的东西有三件套:Base URL、API Key、Model ID。Base URL 就是https://taotoken.net/api,API Key 在控制台的 API Keys 页面生成,Model ID 按你实际要调的模型填。这三件套在后面的settings.json、auth.json或者环境变量里都会用到。生成 Key 的入口在 https://taotoken.net/console/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 。

为什么文件读写教程里要提这个?因为本地联调时你可能会用 Claude Code 或者 Cline 这类工具帮你生成和检查 Spark 代码,这些工具需要走一个统一的模型通道。TaoToken 的 Coding Plan 就是给长期编码场景用的,入口在 https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_content=coding-plan&utm_campaign=rewrite 。如果你只是想验证某个模型能不能正常返回,用模型对话页面就行:https://taotoken.net/models?utm_source=taotoken_aicg_blog_end&utm_content=models&utm_campaign=rewrite 。

环境准备的具体动作:先确认本地有 JDK 8 或 11,Hadoop 客户端能跑hadoop version,Spark 用spark-submit --version能打印版本。然后建一个工作目录,把core-site.xml放进去。这个文件决定你的FileSystem默认指向哪里。本地联调时通常指向file:///或者一个本地 HDFS 伪分布式地址。下面这段可以直接复制,路径按你实际的 Hadoop 配置目录调整。

<?xml version="1.0" encoding="UTF-8"?> <?xml-stylesheet type="text/xsl" href="configuration.xsl"?> <configuration> <property> <name>fs.defaultFS</name> <value>hdfs://localhost:9000</value> </property> <property> <name>hadoop.tmp.dir</name> <value>/tmp/hadoop-local</value> </property> </configuration>

如果你只是操作本地文件,把fs.defaultFS改成file:///即可,这样hadoop fs -ls /tmp走的就是本地磁盘。注意core-site.xml要放在 classpath 能找到的位置,spark-submit时用--files或者直接放进$HADOOP_CONF_DIR。

TaoToken 这条通道在这里的角色是辅助:当你写saveAsNewAPIHadoopFile的参数拿不准时,可以用 Claude Code 走 TaoToken 的 Anthropic 兼容入口去问,入口在 https://taotoken.net/claude-code-anthropic?utm_source=taotoken_aicg_blog_end&utm_content=claude-code-anthropic&utm_campaign=rewrite 。但记住,文件读写的正确性最终靠hadoop fs -ls和spark-submit的输出来验证,不是靠模型说对不对。

3. 可复制的 SparkConf 与 Hadoop Configuration 配置

这一节是整篇的核心,配置写对了,后面基本不会出大问题。先看SparkConf。用hadoopFile读文件时,必须注册 Kryo 序列化类,否则Text类型会抛NotSerializableException。下面这段 Scala 代码可以直接复制到你的main里。

import org.apache.spark.{SparkConf, SparkContext} import org.apache.hadoop.io.{Text, LongWritable, IntWritable} val conf = new SparkConf() .setMaster("local[*]") .setAppName("hadoop-api-demo") .set("spark.serializer", "org.apache.spark.serializer.KryoSerializer") .set("spark.kryoserializer.classesToRegister", "org.apache.hadoop.io.Text,org.apache.hadoop.io.LongWritable,org.apache.hadoop.io.IntWritable") val sc = new SparkContext(conf)

注意spark.kryoserializer.classesToRegister的值是逗号分隔的全限定类名,不要写成Array(...)的形式,那是registerKryoClasses的写法。两种写法等价,但字符串形式在配置文件里更通用。如果你用spark-submit提交,也可以把这些写进spark-defaults.conf:

spark.serializer=org.apache.spark.serializer.KryoSerializer spark.kryoserializer.classesToRegister=org.apache.hadoop.io.Text,org.apache.hadoop.io.LongWritable,org.apache.hadoop.io.IntWritable

接下来是Hadoop Configuration。在 Spark 里拿FileSystem有两种方式:一种是通过sc.hadoopConfiguration,它已经加载了 classpath 下的core-site.xml;另一种是自己new Configuration()再手动addResource。推荐用前者,省事且和 Spark 的配置一致。

val hadoopConf = sc.hadoopConfiguration hadoopConf.set("fs.defaultFS", "hdfs://localhost:9000") val fs = org.apache.hadoop.fs.FileSystem.get(hadoopConf) val path = new org.apache.hadoop.fs.Path("/test/hadoopkv") println("exists = " + fs.exists(path))

如果你要写文件,用saveAsNewAPIHadoopFile时参数顺序是:路径、Key 类、Value 类、OutputFormat 类。新 API 的TextOutputFormat在org.apache.hadoop.mapreduce.lib.output下。

import org.apache.hadoop.mapreduce.lib.output.TextOutputFormat val user = sc.parallelize(Array(("jack", 20), ("jim", 10))) user.saveAsNewAPIHadoopFile( "hdfs://localhost:9000/test/hadoopkv1", classOf[Text], classOf[IntWritable], classOf[TextOutputFormat[Text, IntWritable]] )

老 API 的写法对应saveAsHadoopFile,OutputFormat换成org.apache.hadoop.mapred.TextOutputFormat:

import org.apache.hadoop.mapred.TextOutputFormat user.saveAsHadoopFile( "hdfs://localhost:9000/test/hadoopkv", classOf[Text], classOf[IntWritable], classOf[TextOutputFormat[Text, IntWritable]] )

读取时,新 API 用newAPIHadoopFile,老 API 用hadoopFile。读KeyValueTextInputFormat的示例如下:

import org.apache.hadoop.mapreduce.lib.input.KeyValueTextInputFormat sc.newAPIHadoopFile[Text, Text, KeyValueTextInputFormat]( "hdfs://localhost:9000/test/hadoopkv1" ).foreach { case (k, v) => println(k + "-------" + v) }

老 API 读取:

import org.apache.hadoop.mapred.KeyValueTextInputFormat sc.hadoopFile( "hdfs://localhost:9000/test/hadoopkv1", classOf[KeyValueTextInputFormat], classOf[Text], classOf[Text] ).collect().foreach { case (name, age) => println(name + "======" + age) }

如果你用 TaoToken 的 Coding Plan 辅助生成这些片段,记得把 Base URL、Key、Model ID 三件套配全。比如在 Cline 的 MCP 配置里,或者 Claude Code 的settings.json里,都要写清楚。下面是一个settings.json的片段示例,路径按你实际的配置目录放:

{ "env": { "ANTHROPIC_BASE_URL": "https://taotoken.net/api", "ANTHROPIC_API_KEY": "你的TaoTokenKey", "ANTHROPIC_MODEL": "你的ModelID" } }

Codex 的auth.json类似,把 Base URL 和 Key 填进去即可。注意这些配置只影响模型调用,不影响 Spark 读写文件的行为。

4. 验证请求与成功结果校验

配置写完,接下来是验证。第一步先用hadoop fs -ls确认目录状态。假设你已经用saveAsNewAPIHadoopFile写过一次,执行:

hadoop fs -ls /test/hadoopkv1

正常输出会列出_SUCCESS和part-00000之类的文件。_SUCCESS是 Hadoop 的提交标记,有它说明写入成功。如果没有_SUCCESS,说明任务中途失败或者被中断。你可以再看内容:

hadoop fs -cat /test/hadoopkv1/part-00000

预期输出是jack 20和jim 10这样的键值对,中间是制表符。如果你写的是老 API 的saveAsHadoopFile,目录是/test/hadoopkv,内容格式一样。

第二步用spark-submit跑完整任务。把上面的代码打包成 jar,然后提交:

spark-submit \ --class com.example.HadoopApiDemo \ --master local[*] \ --conf spark.serializer=org.apache.spark.serializer.KryoSerializer \ --conf spark.kryoserializer.classesToRegister=org.apache.hadoop.io.Text,org.apache.hadoop.io.LongWritable,org.apache.hadoop.io.IntWritable \ target/hadoop-api-demo-1.0.jar

提交后看控制台输出。如果读到数据,会打印jack-------20和jim-------10。如果写入成功,任务结束前不会有异常堆栈。这里有个细节:local[*]模式下fs.defaultFS如果指向hdfs://localhost:9000,你需要本地确实有一个可连的 HDFS。如果没有,就把core-site.xml改成file:///,路径换成/tmp/test/hadoopkv1,这样读写都走本地磁盘,验证逻辑完全一样。

第三步校验读写一致性。写进去再读出来,对比条数和内容。可以在代码里加一段:

val readBack = sc.newAPIHadoopFile[Text, Text, KeyValueTextInputFormat]( "hdfs://localhost:9000/test/hadoopkv1" ).collect() println("count = " + readBack.length) readBack.foreach { case (k, v) => println(k.toString + " -> " + v.toString) }

预期count = 2,内容与写入一致。如果count是 0,先检查路径是不是写到了别的目录,或者KeyValueTextInputFormat的分隔符是不是默认的制表符。KeyValueTextInputFormat默认按\t切分,如果你的数据是逗号分隔,需要额外设置key.value.separator.in.input.line。

如果你在联调时用 TaoToken 的模型对话页面验证模型返回,入口在 https://taotoken.net/models?utm_source=taotoken_aicg_blog_end&utm_content=models&utm_campaign=rewrite ,但那和文件读写是两条线,别混在一起排查。

5. 本篇常见报错排查

第一个高频报错是java.io.NotSerializableException: org.apache.hadoop.io.Text。这个前面提过,根因是没注册 Kryo 类。解决动作:在SparkConf里加spark.serializer和spark.kryoserializer.classesToRegister,把Text、LongWritable、IntWritable都登记进去。如果你用的是registerKryoClasses,写法是conf.registerKryoClasses(Array(classOf[Text], classOf[LongWritable])),注意这是SparkConf的方法,不是set字符串。

第二个报错是ClassNotFoundException: org.apache.hadoop.mapred.TextOutputFormat或者反过来找不到mapreduce.lib.output.TextOutputFormat。这是新旧 API 混用导致的。记住对应关系:saveAsHadoopFile配org.apache.hadoop.mapred.TextOutputFormat,saveAsNewAPIHadoopFile配org.apache.hadoop.mapreduce.lib.output.TextOutputFormat。输入侧同理,hadoopFile配org.apache.hadoop.mapred.KeyValueTextInputFormat,newAPIHadoopFile配org.apache.hadoop.mapreduce.lib.input.KeyValueTextInputFormat。

第三个报错是401 Unauthorized或者local proxy failed。如果你在联调时用了 TaoToken 的通道,401 通常是 Key 没填对或者 Base URL 写错了。检查settings.json或auth.json里的ANTHROPIC_BASE_URL是不是https://taotoken.net/api,Key 是不是从控制台复制完整。local proxy failed一般是本地代理配置和实际网络环境不匹配,把代理相关配置清掉,直连即可。注意这里说的是模型调用通道的报错,不是 HDFS 的报错。

第四个报错是reading choices相关的解析失败。这通常出现在模型返回格式不符合预期时,比如你期望 JSON 但返回了纯文本。排查动作:先用模型对话页面单独发一条请求,确认返回结构,再检查你的解析代码。这个和 Spark 文件读写无关,属于辅助通道的问题。

第五个报错是OAuth相关的鉴权失败。如果你用 Claude Code 走 TaoToken 的 Anthropic 兼容入口,OAuth 流程没走完或者 token 过期都会报这个。重新生成 Key,或者检查settings.json里的字段名是不是ANTHROPIC_API_KEY。入口在 https://taotoken.net/claude-code-anthropic?utm_source=taotoken_aicg_blog_end&utm_content=claude-code-anthropic&utm_campaign=rewrite 。

第六个报错是FileAlreadyExistsException。saveAsNewAPIHadoopFile和saveAsHadoopFile都不允许目标目录已存在。解决动作:写入前先删掉目标目录,或者换一个不存在的路径。用fs.delete(path, true)递归删除。

if (fs.exists(new org.apache.hadoop.fs.Path("/test/hadoopkv1"))) { fs.delete(new org.apache.hadoop.fs.Path("/test/hadoopkv1"), true) }

第七个报错是No FileSystem for scheme: hdfs。这说明 classpath 里缺 Hadoop HDFS 的依赖,或者core-site.xml没被加载。检查spark-submit时有没有把 Hadoop 的 jar 带上,或者HADOOP_CONF_DIR有没有指向正确的配置目录。

6. 继续联调与通道选择

文件读写跑通之后,下一步通常是把它接进更长的编码流程。如果你只是偶尔验证一下模型返回,用模型对话页面就够了。如果你要长期写 Spark 任务、反复让工具帮你检查saveAsNewAPIHadoopFile的参数,那 Coding Plan 更合适,入口在 https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_content=coding-plan&utm_campaign=rewrite 。生成 Key 还是那个地址:https://taotoken.net/console/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 。

实测下来,最容易翻车的不是 Spark 代码本身,而是core-site.xml的fs.defaultFS和实际环境不一致。本地联调时先用file:///把逻辑跑通,再切到hdfs://验证分布式路径,这样排查范围小很多。另外spark.kryoserializer.classesToRegister里把IntWritable也加上,别只加Text,否则写IntWritable时照样报序列化错误。

返回列表