ARTICLE DETAIL

资讯详情

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

Spark学习01——创建RDD的所有方法:从内存集合到外部存储的完整实践

Spark学习01——创建RDD的所有方法:从内存集合到外部存储的完整实践 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_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_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://读 HDFSs3n://或s3a://读对象存储。这里第一个大坑就是路径格式。在 Windows 上写file:///D:/data/people.txt是三个斜杠Linux 上是file:///home/user/people.txt。少一个斜杠或者用反斜杠就会报IllegalArgumentException: Wrong FS或者找不到文件。textFile还支持目录、压缩文件、通配符。比如textFile(/my/directory)读整个目录textFile(/my/directory/*.txt)读所有 txttextFile(/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/apiKey 在控制台生成Model ID 按你选的模型填。控制台地址是 https://taotoken.net/console?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content API Keys 管理在 https://taotoken.net/api-keys?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_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可以用来限定行范围数据量大时建议加上。把四种方式的返回类型列个表对照一目了然方法返回类型适用场景textFileRDD[String]按行处理文本wholeTextFilesRDD[(String, String)]按文件处理小文件sequenceFileRDD[(K, V)]读写 SequenceFilenewAPIHadoopRDDRDD[(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 配套。如果你用的是 TaoTokenBase URL 是https://taotoken.net/apiKey 在 https://taotoken.net/api-keys?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_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.xHBase 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_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content 。需要验证模型是否通可以用模型对话页面 https://taotoken.net/chat?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content 发一条测试消息。长期编码任务用 Coding Plan https://taotoken.net/coding-plan?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content 。Claude Code 相关接入参考 https://taotoken.net/claude-code-anthropic?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_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.parallelism8 \ SparkLearnExample.jar注意集群模式下setMaster(local[*])要去掉或者用--master覆盖。路径也要从file://改成hdfs://。最后说个实用技巧创建 RDD 后先别急着写业务逻辑先count()和getNumPartitions确认数据量和分区数符合预期。这一步花几秒钟能省掉后面几小时的排查。数据量小就少分区数据量大就多分区但别超过集群核数的 4 倍否则调度开销会吃掉收益。这套流程我在本地和集群都跑过最容易出问题的永远是路径和依赖版本。把这两个盯住六种创建方式基本一次过。
返回列表