1. 先搞清楚 RDD 到底是什么,为什么创建方式值得单独学
RDD 是 Spark 里最基础的数据抽象,全称弹性分布式数据集。你可以把它理解成一份「只读、可分片、能重算」的数据清单:它不直接存数据,而是记录「这份数据从哪来、怎么算出来」。真正跑任务时,Spark 才按这份清单把数据切成若干分区,分发到不同 executor 上并行处理。
对刚入门的人来说,创建 RDD 是绕不过去的第一道坎。原因很直接:后面所有的 map、filter、reduceByKey 都建立在「你手上已经有一个 RDD」这个前提上。而创建方式又分好几类,从内存集合、本地文件、HDFS、对象存储,一直到 HBase 这类外部系统,每种方式的参数含义、分区行为、返回类型都不一样。如果这一步没跑通,后面学算子就是空中楼阁。
这篇聚焦的场景很明确:Spark 入门者想在自己机器上把「创建 RDD 的所有方法」一次性跑通,同时希望本地环境和集群环境用同一套 Key/API 通道来管理配置,避免到处散落密钥。我会把 parallelize、makeRDD、textFile、wholeTextFiles、sequenceFile、newAPIHadoopRDD 全部过一遍,每个都给可复制代码和验证动作,重点讲清楚分区数、路径格式、返回类型这三个最容易踩坑的地方。
适合谁看:写过一点 Scala 或 Java、装过 Spark、但一遇到textFile路径报错或者wholeTextFiles返回类型对不上就卡住的人。如果你还没装 Spark,建议先把本地环境跑起来再回来对照,因为下面的验证步骤都依赖能实际提交任务。
先说一个贯穿全文的配置思路。入门阶段最烦的是环境变量、密钥、集群地址散落在各个文件里,换个环境就要改一遍。我习惯把这类连接信息统一走一个 API 通道来管理,本地和集群读同一份配置,减少「本地能跑集群报错」的排查成本。TaoToken 就是干这个的:它提供一个统一的 Key 和 API 入口,把模型调用、编码计划、控制台这些能力收敛到一处。官网入口是 https://taotoken.net/?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= ,API 地址是 https://taotoken.net/api 。注意 API 地址不带 UTM 参数,配置里直接写这个就行。
需要说明的是,TaoToken 管的是「连接通道和密钥」,不替代你的 Spark 运行环境,也不替代编辑器。Spark 本身还是跑在你自己的 JVM 和集群上。它的价值在于:当你后面要接大模型做代码辅助、或者用 Coding Plan 跑长期编码任务时,不用每个工具单独配一套密钥。下面进入正题,先看内存集合创建。
2. 从内存集合创建:parallelize 与 makeRDD 的区别和分区检查
从内存集合创建 RDD 是最快上手的方式,适合做实验和单元测试。核心就两个方法:parallelize和makeRDD。很多人以为它俩完全一样,其实有个细节值得说清楚。
parallelize的签名是parallelize[T](seq: Seq[T], numSlices: Int = defaultParallelism)。它把本地集合切成分区,默认分区数取决于你的运行模式:local 模式下通常是 CPU 核数,集群模式下由spark.default.parallelism决定。
makeRDD有两个重载。第一个重载和parallelize完全一致,就是换个名字;第二个重载接收Seq[(T, Seq[String])],也就是每个元素可以带「位置信息」,告诉 Spark 这份数据更希望被调度到哪些节点上。这个在需要数据本地性的场景才有意义,入门阶段基本用不到,但面试常问。
先看一段可复制的代码,把两种方式都跑一遍,并打印分区数:
import org.apache.spark.{SparkConf, SparkContext} object CreateRDDMemory { def main(args: Array[String]): Unit = { val conf = new SparkConf() .setAppName("CreateRDDMemory") .setMaster("local[*]") val sc = new SparkContext(conf) // 方式一:parallelize val rdd1 = sc.parallelize(List("zhangsan", "lisi", "wangwu"), 2) println("rdd1 分区数 = " + rdd1.getNumPartitions) rdd1.foreach(println) // 方式二:makeRDD 第一种重载,等价于 parallelize val rdd2 = sc.makeRDD(List("zhangsan", "lisi", "wangwu"), 2) println("rdd2 分区数 = " + rdd2.getNumPartitions) rdd2.foreach(println) // 方式三:makeRDD 第二种重载,带位置信息 val seqWithLoc = Seq( ("zhangsan", Seq("node1")), ("lisi", Seq("node2")) ) val rdd3 = sc.makeRDD(seqWithLoc) println("rdd3 分区数 = " + rdd3.getNumPartitions) rdd3.foreach(println) sc.stop() } }跑完之后你会看到类似输出:rdd1 分区数 = 2,然后三行名字。这里有个检查动作很关键:把numSlices从 2 改成 5,再跑一次,观察分区数变化。如果集合只有 3 个元素却要 5 个分区,Spark 会创建 5 个分区,其中两个是空的。这个现象在后续mapPartitions里会体现出来,空分区也会走一遍函数,容易埋坑。
实测下来,入门阶段建议显式传分区数,不要依赖默认值。因为默认值在不同运行模式下不一样,本地跑得好好的,打包到集群可能分区数暴涨或暴跌,影响并行度和 shuffle 行为。
还有一个常见误区:以为parallelize会把数据复制到所有节点。实际上它是在 driver 端持有集合,然后按分区切分后分发。如果集合特别大(比如几百万条),driver 内存会吃紧,这时候应该改用外部存储方式。内存集合创建只适合小数据量实验。
分区数怎么定?一个经验值是「分区数 = 集群总核数的 2 到 4 倍」。本地local[*]下,defaultParallelism就是核数。你可以用sc.defaultParallelism打印出来看看,再决定传多少。
到这里内存方式就通了。接下来进入更常用的外部存储方式,先从最典型的textFile开始,它也是报错最多的地方。
3. 从文件系统创建:textFile 路径格式与分区参数配置
textFile是日常用得最多的创建方式,支持本地文件系统、HDFS、S3 等。它的签名是textFile(path: String, minPartitions: Int = defaultMinPartitions),返回RDD[String],每一行是一个元素。
路径前缀决定读哪里:file://读本地,hdfs://读 HDFS,s3n://或s3a://读对象存储。这里第一个大坑就是路径格式。在 Windows 上写file:///D:/data/people.txt是三个斜杠,Linux 上是file:///home/user/people.txt。少一个斜杠或者用反斜杠,就会报IllegalArgumentException: Wrong FS或者找不到文件。
textFile还支持目录、压缩文件、通配符。比如textFile("/my/directory")读整个目录,textFile("/my/directory/*.txt")读所有 txt,textFile("/my/directory/*.gz")读压缩文件。压缩文件在读取时会自动解压,但注意压缩格式要能被 Hadoop 的编解码器识别,.gz和.bz2一般没问题。
第二个参数minPartitions控制最小分区数。默认情况下,Spark 为文件的每个块创建一个分区,HDFS 块默认 128MB。你可以传更大的值请求更多分区,但不能比块数少。比如一个 256MB 的文件在 HDFS 上是 2 个块,你传minPartitions = 1,实际还是 2 个分区;传 4,就会切成 4 个分区。
下面这段配置可以直接复制,注意把路径换成你自己的:
import org.apache.spark.{SparkConf, SparkContext} object CreateRDDTextFile { def main(args: Array[String]): Unit = { val conf = new SparkConf() .setAppName("CreateRDDTextFile") .setMaster("local[*]") val sc = new SparkContext(conf) // 本地文件,注意 file:/// 三个斜杠 val rdd = sc.textFile("file:///D:/github/SparkLearnExample/examples/src/main/resources/people.txt", 2) println("textFile 分区数 = " + rdd.getNumPartitions) rdd.foreach(println) // 读整个目录 val rddDir = sc.textFile("file:///D:/github/SparkLearnExample/examples/src/main/resources") println("目录分区数 = " + rddDir.getNumPartitions) sc.stop() } }跑通后你会看到people.txt的每一行被打印出来。检查动作有两个:第一,把minPartitions从 2 改成 1 再改成 4,观察getNumPartitions的变化,理解「不能少于块数」这条规则;第二,故意把路径写成file://D:/...(两个斜杠),看报错信息长什么样,记住这个错,以后一眼能认出来。
如果你用 HDFS,路径写成hdfs://namenode:8020/user/data/people.txt。端口号按你集群实际配置来,常见是 8020 或 9000。S3 的话用s3a://bucket/path,需要额外配置 access key 和 secret key,这部分建议走统一的密钥管理,别硬编码在代码里。
关于密钥管理,这里可以顺带说下 TaoToken 的接入方式。如果你后面要用大模型辅助写 Spark 代码,或者用 Coding Plan 跑长期任务,可以在配置里统一填 Base URL、Key、Model ID 三件套。Base URL 填https://taotoken.net/api,Key 在控制台生成,Model ID 按你选的模型填。控制台地址是 https://taotoken.net/console?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= ,API Keys 管理在 https://taotoken.net/api-keys?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= 。这样本地和集群读同一份配置,不用改代码。
回到 textFile。还有一个细节:它返回的是RDD[String],每行一个字符串。如果你需要行号,得自己用zipWithIndex,但注意这会触发一次额外作业。如果文件里有表头,记得用filter去掉或者用first单独取。
textFile讲完,接下来是它的「兄弟」wholeTextFiles,返回类型完全不同,这是新手最容易搞混的地方。
4. wholeTextFiles、sequenceFile 与 newAPIHadoopRDD 的返回类型对照
wholeTextFiles和textFile最大的区别在返回类型。textFile返回RDD[String],每行一个元素;wholeTextFiles返回RDD[(String, String)],键是文件路径,值是整个文件的全部内容。也就是说,它把每个小文件整体读成一个字符串,适合文件不大但需要按文件处理的场景。
签名是wholeTextFiles(path: String, minPartitions: Int = defaultMinPartitions): RDD[(String, String)]。注意它读的是目录时,会把目录下每个文件作为一个元素。如果文件很大,单个字符串可能撑爆内存,所以它只适合小文件。
val rdd4 = sc.wholeTextFiles("file:///D:/github/SparkLearnExample/examples/src/main/resources") rdd4.foreach { case (path, content) => println("文件路径: " + path) println("内容长度: " + content.length) }跑完你会看到每个文件的完整路径和内容长度。检查动作:对比textFile读同一目录的输出,你会发现textFile是把所有文件的行混在一起,而wholeTextFiles是按文件分组。这个区别在做「按文件聚合」时非常关键。
sequenceFile读的是 Hadoop 的 SequenceFile 格式,返回RDD[(K, V)],K 和 V 必须是 Hadoop Writable 接口的子类。签名是sequenceFile[K, V](path: String, keyClass: Class[K], valueClass: Class[V], minPartitions: Int)。常见用法是sc.sequenceFile("hdfs://xxx", classOf[Text], classOf[Text])。注意这里要传 class 对象,不是字符串。
newAPIHadoopRDD是最通用的方式,能读任意 Hadoop 输入格式,包括 HBase。签名是newAPIHadoopRDD[K, V, F <: InputFormat[K, V]](conf: Configuration, fClass: Class[F], kClass: Class[K], vClass: Class[V]): RDD[(K, V)]。读 HBase 的典型写法:
import org.apache.hadoop.conf.Configuration import org.apache.hadoop.hbase.HBaseConfiguration import org.apache.hadoop.hbase.client.Result import org.apache.hadoop.hbase.io.ImmutableBytesWritable import org.apache.hadoop.hbase.mapreduce.TableInputFormat val hconf: Configuration = HBaseConfiguration.create() hconf.set("hbase.zookeeper.quorum", "xxx:2181,xxx:2181,xxx:2181") hconf.set(TableInputFormat.INPUT_TABLE, "xxx") hconf.set(TableInputFormat.SCAN_COLUMNS, "fields:phone_no fields:contacts_list") val pairRdd = sc.newAPIHadoopRDD( hconf, classOf[TableInputFormat], classOf[ImmutableBytesWritable], classOf[Result] ) pairRdd.foreach(println)这里返回的是RDD[(ImmutableBytesWritable, Result)],键是行键,值是 HBase 的 Result 对象。注意SCAN_COLUMNS里列族和列用冒号分隔,多列用空格分隔。注释掉的SCAN_ROW_START和SCAN_ROW_STOP可以用来限定行范围,数据量大时建议加上。
把四种方式的返回类型列个表对照,一目了然:
| 方法 | 返回类型 | 适用场景 |
|---|---|---|
| textFile | RDD[String] | 按行处理文本 |
| wholeTextFiles | RDD[(String, String)] | 按文件处理小文件 |
| sequenceFile | RDD[(K, V)] | 读写 SequenceFile |
| newAPIHadoopRDD | RDD[(K, V)] | 任意 Hadoop 输入格式,如 HBase |
检查动作:对每个方法都打印getNumPartitions和count(),确认数据真的读进来了。count()会触发实际计算,如果路径错或者格式不对,这里就会报错。
这里要提醒一句:newAPIHadoopRDD读 HBase 时,如果 ZooKeeper 地址写错,会卡住很久然后报连接超时。建议先用hbase shell确认集群能连上,再写 Spark 代码。另外,别在生产库上直接跑全表扫描,加行范围或者列过滤。
5. 常见报错排查:401、路径错误、分区数异常与 OAuth 问题
这一节把入门阶段最常撞的报错集中过一遍,每个都给现象、原因、解决动作。
第一个是401 Unauthorized。这个通常出现在你调用外部 API 时,比如用大模型辅助编码、或者访问需要鉴权的存储。现象是请求直接被拒,日志里有 401。原因一般是 Key 没填、填错、或者过期。解决动作:去控制台重新生成 Key,确认 Base URL 和 Key 配套。如果你用的是 TaoToken,Base URL 是https://taotoken.net/api,Key 在 https://taotoken.net/api-keys?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= 生成。注意 Base URL 不要带多余路径,Key 不要有多余空格。
第二个是local proxy failed。这个报错通常和网络代理配置有关。现象是连接超时或者代理拒绝。解决动作:检查你的环境变量里有没有残留的代理设置,比如http_proxy、https_proxy。如果有,确认代理地址是否可达;如果不需要代理,直接清掉这些变量。在 Spark 里,还要检查spark-submit有没有传--conf spark.hadoop.*.proxy之类的参数。
第三个是reading choices相关报错,完整信息类似Exception in thread "main" java.lang.NoSuchMethodError或者reading choices字段解析失败。这个多出现在依赖版本冲突时,比如 Spark 版本和 Hadoop 版本不匹配,或者 HBase 客户端版本和集群不一致。解决动作:用mvn dependency:tree看依赖树,把冲突的包排除掉,统一版本。Spark 3.x 建议配 Hadoop 3.x,HBase 2.x。
第四个是 OAuth 相关报错。如果你用某些云服务或者需要 OAuth 认证的接口,可能会遇到 token 过期或者 scope 不对。现象是 403 或者invalid_token。解决动作:重新走一遍授权流程,确认 scope 包含你需要的权限。如果是长期任务,建议用能自动刷新的凭证方式,别用短期 token。
第五个是路径错误,前面提过,这里给完整报错样例:java.lang.IllegalArgumentException: Wrong FS: file:///D:/data, expected: hdfs://namenode:8020。原因是你在 HDFS 模式下用了file://路径。解决动作:确认运行模式,本地模式用file://,集群模式用hdfs://。或者用sc.hadoopConfiguration.set("fs.defaultFS", "hdfs://namenode:8020")统一设置。
第六个是分区数异常。现象是getNumPartitions返回 0 或者特别大。返回 0 通常是空目录或者过滤后没数据;特别大可能是minPartitions传太大,或者读了很多小文件。解决动作:先count()确认有数据,再根据数据量调整分区数。小文件多的话,考虑用wholeTextFiles或者先合并。
第七个是ClassNotFoundException,比如找不到TableInputFormat。原因是 HBase 相关依赖没打进包。解决动作:在pom.xml里加hbase-client和hbase-mapreduce依赖,打包时用mvn package带上依赖,或者用--jars传。
把这些报错整理成排查顺序:先看路径和 FS 是否匹配,再看依赖版本,再看鉴权和网络,最后看分区和数据量。大部分问题在前两步就能定位。
如果你在排查过程中需要查文档,接入文档在 https://taotoken.net/doc?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= 。需要验证模型是否通,可以用模型对话页面 https://taotoken.net/chat?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= 发一条测试消息。长期编码任务用 Coding Plan https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= 。Claude Code 相关接入参考 https://taotoken.net/claude-code-anthropic?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= 。
6. 把六种创建方式串成一个可运行工程
最后给一个完整的工程结构,把前面所有方式串起来,你可以直接复制到 IDE 里跑。目录结构建议这样:
SparkLearnExample/ src/main/scala/ CreateRDDFun.scala src/main/resources/ people.txtpeople.txt内容随便写几行,比如:
zhangsan 20 lisi 25 wangwu 30CreateRDDFun.scala完整代码:
import org.apache.hadoop.conf.Configuration import org.apache.hadoop.hbase.HBaseConfiguration import org.apache.hadoop.hbase.client.Result import org.apache.hadoop.hbase.io.ImmutableBytesWritable import org.apache.hadoop.hbase.mapreduce.TableInputFormat import org.apache.hadoop.io.Text import org.apache.spark.{SparkConf, SparkContext} object CreateRDDFun { def main(args: Array[String]): Unit = { val sparkConf = new SparkConf() .setAppName("CreateRDDFun examples") .setMaster("local[*]") val sc = new SparkContext(sparkConf) create01(sc) create02(sc) create03(sc) create04(sc) // create05 和 create06 需要 HDFS/HBase 环境,按需打开 // create05(sc) // create06(sc) sc.stop() } private def create01(sc: SparkContext): Unit = { val rdd = sc.parallelize(List("zhangsan", "lisi", "wangwu"), 2) println("create01 分区数 = " + rdd.getNumPartitions) rdd.foreach(println) } private def create02(sc: SparkContext): Unit = { val rdd2 = sc.makeRDD(List("zhangsan", "lisi", "wangwu"), 2) println("create02 分区数 = " + rdd2.getNumPartitions) rdd2.foreach(println) } private def create03(sc: SparkContext): Unit = { val rdd3 = sc.textFile("file:///D:/github/SparkLearnExample/examples/src/main/resources/people.txt", 2) println("create03 分区数 = " + rdd3.getNumPartitions) rdd3.foreach(println) } private def create04(sc: SparkContext): Unit = { val rdd4 = sc.wholeTextFiles("file:///D:/github/SparkLearnExample/examples/src/main/resources") println("create04 分区数 = " + rdd4.getNumPartitions) rdd4.foreach { case (path, content) => println("路径: " + path + " 长度: " + content.length) } } private def create05(sc: SparkContext): Unit = { val rdd5 = sc.sequenceFile("hdfs://xxx", classOf[Text], classOf[Text]) rdd5.foreach(println) } private def create06(sc: SparkContext): Unit = { val hconf: Configuration = HBaseConfiguration.create() hconf.set("hbase.zookeeper.quorum", "xxx:2181,xxx:2181,xxx:2181") hconf.set(TableInputFormat.INPUT_TABLE, "xxx") hconf.set(TableInputFormat.SCAN_COLUMNS, "fields:phone_no fields:contacts_list") val pairRdd = sc.newAPIHadoopRDD( hconf, classOf[TableInputFormat], classOf[ImmutableBytesWritable], classOf[Result] ) pairRdd.foreach(println) } }跑之前检查三件事:第一,people.txt路径和代码里一致;第二,local[*]模式不需要额外配置;第三,create05和create06需要真实集群,本地跑先注释掉。
跑通后你会看到每个方法的输出和分区数。建议把create03的minPartitions改成 1、4、8 各跑一次,观察分区数变化,这是理解分区最直观的方式。
如果你要把这套代码打包到集群,spark-submit命令大概长这样:
spark-submit \ --class CreateRDDFun \ --master yarn \ --deploy-mode client \ --conf spark.default.parallelism=8 \ SparkLearnExample.jar注意集群模式下setMaster("local[*]")要去掉,或者用--master覆盖。路径也要从file://改成hdfs://。
最后说个实用技巧:创建 RDD 后先别急着写业务逻辑,先count()和getNumPartitions确认数据量和分区数符合预期。这一步花几秒钟,能省掉后面几小时的排查。数据量小就少分区,数据量大就多分区,但别超过集群核数的 4 倍,否则调度开销会吃掉收益。
这套流程我在本地和集群都跑过,最容易出问题的永远是路径和依赖版本。把这两个盯住,六种创建方式基本一次过。