ARTICLE DETAIL

资讯详情

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

Spark RDD编程实战:10个项目全流程解析与性能调优

Spark RDD编程实战:10个项目全流程解析与性能调优 1. 项目整体设计10个项目的选型逻辑与学习路径1.1 为什么选这10个项目到底在练什么能力先说实话网上Spark RDD的教程一抓一大把但大多数都停在WordCount这个层面讲一讲flatMap、reduceByKey就结束了你照着敲完还是不知道自己能做什么。我整理这套Spark RDD编程实战项目时最核心的思路就是按“一个RDD程序需要掌握的能力”来选题——文本读取、JSON解析、日志分析、多表关联、TopN、分区控制、共享变量、数据倾斜优化最后用两个完整业务场景把前面所有知识点串成一条线。为什么要这样设计因为工作中真正遇到的RDD作业几乎都是这些能力的组合。比如给你一份网约车订单日志你需要先解析非结构化文本再过滤脏数据然后做维度聚合最后按条件排序导出。这一套流程里如果只知道map和reduceByKey应付不了真实数据里的各种意外情况。我把10个项目按“入门—进阶—调优”三层递进每个项目只比上一个多一两个新知识点保证小白不会一看就放弃又能感受到难度在稳步提升。选数据时我也做了些考虑。日志数据、订单数据、农产品价格数据、成绩数据这些都是实际业务里最常见的格式比凭空造的a、b、c字段更容易理解需求。你别小看这些选择学习一件东西最好的方式就是让它贴近真实场景。很多初学者看到“用户表、订单表、金额统计”就天然能理解要做什么这比抽象的key-value数据有感觉得多。1.2 为什么现在还要学RDD它到底解决什么问题我知道有人会问Spark SQL和DataFrame不香吗确实我日常写分析任务也优先用DataFrame但RDD依然是理解Spark执行引擎最直接的一条路径。RDD可以看作整个Spark的底层抽象基石分区、依赖、Shuffle、惰性求值这些概念只有落到RDD层面你才能真正看清。我从面试和实际辅导的经验来看很多人简历里写着“熟悉Spark”但一问reduceByKey和groupByKey的区别就卡壳一谈数据倾斜就只会说“加盐”。这些恰恰是RDD最能讲清楚的问题。比如reduceByKey为什么比groupByKey好因为reduceByKey在map端就做预聚合Shuffle的数据量小很多。这个结论在DataFrame里看不到但RDD代码一跑打开Spark UI看Shuffle字节数一下就明白了。RDD在两类场景下依然有不可替代的优势。第一种是处理完全非结构化的数据比如自定义格式的日志、规则复杂的文本DataFrame的schema反而碍事。第二种是高度定制化的清洗逻辑比如一个字段要根据前后文判断才能决定怎么解析RDD的mapPartitions能给你最大的灵活性。这10个项目就跑通了这两类场景同时把Spark底层的执行原理埋在里面属于那种“现在练了后面一定会用上”的投资。2. 环境准备与数据构造别让环境拖慢学习进度2.1 本地开发环境怎么搭最省事这10个项目全部用local模式就能跑不需要一上来就搭集群。学习阶段集群只有坏处提交任务要等资源、看日志要翻YARN页面、改代码要打包上传。本地模式下Spark会起多个线程模拟executorShuffle、分区、缓存这些机制照样会触发足够理解原理。等你把项目跑完再去看集群部署和提交参数会轻松很多。版本上我选的是Spark 3.3.2 Scala 2.12 JDK8这个组合经过大量生产验证非常稳。如果你用Maven或SBT核心依赖只需要一条libraryDependencies org.apache.spark %% spark-core % 3.3.2如果连构建工具都不想配置直接下Spark发行版把jars目录添加到IDEA的Libraries里也能跑。但这里有个坑我必须提前说IDEA里跑local模式经常会看到一堆关于log4j的警告严重时会直接报错。解决办法是在Run Configuration的VM options里加上一行-Dlog4j2.loggerContextFactoryorg.apache.logging.log4j.simple.SimpleLoggerContextFactory这个配置我当年查了快两小时才搞定。加上之后日志干净排查自己的print输出也方便。2.2 配套数据是怎么来的以及为什么它比真实数据更好用这套项目的“附完整数据”不是随便生成几千行随机文本而是专门设计了数据特征。比如日志数据access.log我生成了500行字段包含IP、时间、请求URL、状态码、响应字节数但故意在中间混了几行IP为空、状态码为“abc”的坏数据。订单数据orders.json大概300条字段有订单号、用户ID、商品品类、金额、下单时间偶尔缺金额字段。成绩数据scores.csv则做了一定的数据倾斜——某个科目的数据量比其他科目大很多方便后面的倾斜处理项目有素材可用。这里顺便给你一个生成测试数据的模板你可以按需扩展。比如生成随机访问日志的Python脚本import random, time ips [f192.168.1.{i} for i in range(1, 20)] for _ in range(100): ip random.choice(ips) ts int(time.time()) url random.choice([/login, /cart, /order, /pay]) print(f{ip}\t{ts}\t{url}\t200\t{random.randint(100, 5000)})学习阶段的测试数据不需要很大但一定要包含“脏数据”和“倾斜特征”。完全没有脏数据你练不了解析和过滤数据完全均匀你体会不到partition倾斜带来的性能差异。我在每个数据文件里都埋了这些点跑项目的时候留意一下收获会更大。3. 从零开始入门项目的完整拆解3.1 项目1WordCount——用一条数据处理链路理解RDDWordCount虽然被写烂了但它确实覆盖了RDD最核心的四个动作读取、逐条转换、聚合、输出。直接上代码我加了比较详细的注释import org.apache.spark.{SparkConf, SparkContext} object WordCount { def main(args: Array[String]): Unit { val conf new SparkConf().setAppName(WordCount).setMaster(local[*]) val sc new SparkContext(conf) val lines sc.textFile(data/words.txt) .filter(_.trim.nonEmpty) // 第一层过滤空行 val counts lines .flatMap(_.split(\\s)) // 按空白字符切分注意多个空格 .map(word (word, 1)) // 转成key-value对 .reduceByKey(_ _) // 按word聚合计数 .sortBy(_._2, ascending false) // 按出现次数降序 counts.collect().foreach(println) sc.stop() } }第一个要注意的点是切分规则。用\\s而不是 是因为真实文本里经常有多个空格和制表符混在一起一个空格切分会出现空字符串影响统计结果。第二个点是filter放在textFile之后能提前丢弃空行减少后面的无效处理。代码跑完以后我建议你把filter那行注释掉再跑一次对比结果你会对“脏数据是怎么影响结果的”有直观感受。跑WordCount时顺手做两件事。第一打开Spark UI的localhost:4040页面看一眼这个作业的DAG图观察哪些操作是宽依赖stage会切分哪些是窄依赖。第二把输入文件的行数扩大10倍看作业变化。我见过太多人只会跑结果从没看过执行计划这样学RDD等于白学。3.2 项目2日志PV/UV统计——行动算子与缓存的关系PV/UV是日志分析最基础的需求。PV是总访问次数直接count就行UV是独立访客数需要按IP去重再统计。代码不长val accessRdd sc.textFile(data/access.log) .map(line line.split(\t)) .filter(cols cols.length 5 cols(0).trim.nonEmpty) val pv accessRdd.count() val uv accessRdd.map(cols cols(0)).distinct().count()这段代码背后有个容易被忽视的性能问题。accessRdd本身是从磁盘读取并解析出来的一个RDDcount和distinct是两个不同的action如果这两个动作之间没有缓存每次action都会从头读取、解析一遍数据。本地测试数据量小没事换成上GB的日志多跑几个action会明显变慢。解决办法很简单在第一个action执行前加上缓存accessRdd.cache() // 或者 accessRdd.persist(StorageLevel.MEMORY_AND_DISK)我建议用persist(StorageLevel.MEMORY_AND_DISK)而不是cache()因为cache只把数据放在内存数据量大时可能导致executor频繁GC而MEMORY_AND_DISK会在内存放不下时溢写到磁盘更稳。你可以把这段代码跑两次一次加persist一次不加然后对比Spark UI的“Input”列读取字节数直观体验缓存带来的收益。3.3 项目3读取JSON并完成字段清洗——非结构化数据入门JSON在真实项目里太常见了。Spark的RDD没有直接读JSON的API常规做法是textFile读整行再逐条解析。我用的是json4s库SBT依赖加一行libraryDependencies org.json4s %% json4s-jackson % 4.0.6解析和清洗的核心逻辑如下import org.json4s._ import org.json4s.jackson.JsonMethods._ implicit val formats: DefaultFormats.type DefaultFormats val rawRdd sc.textFile(data/orders.json) val parsedRdd rawRdd.map { line try { val json parse(line) val orderId compact(render(json \ orderId)) val userId (json \ userId).extract[Int] val amount json \ amount match { case JDouble(v) v case JInt(v) v.toDouble case _ 0.0 } (orderId, userId, amount) } catch { case _: Exception (, 0, 0.0) // 解析失败用占位值 } }.filter { case (orderId, _, amount) orderId.nonEmpty amount 0.0 }这项目有两条经验值得你记下来。第一extract遇到字段类型不匹配会直接抛异常所以解析必须套try catch脏数据换成占位值之后再用filter清洗掉比在map里写if判断干净得多。第二JSON里的金额字段可能来自不同业务系统有时是整数有时是浮点数不能直接extract[Double]我用match表达式做了类型适配。这类问题在DataFrame里有schema推断兜底但在RDD里全靠自己处理写一次就能提升对数据的敏感度。我还故意在订单数据里放了一条JSON格式完全错误的记录以及一条缺金额字段的记录这两个case都是为这个项目准备的。4. 能力进阶聚合、关联与自定义操作4.1 项目4分组统计与TopN——groupByKey的替代方案需求很常见成绩表按科目分组统计每个科目的前三名。最直觉的写法是groupByKey后逐个排序但我给两个方案第二个你以后在真实作业里会非常受用。case class Score(name: String, subject: String, score: Int) val rdd sc.textFile(data/scores.txt).map { line val cols line.split(,) Score(cols(0), cols(1), cols(2).toInt) } // 方案一groupByKey后排序直观但把所有值都集中到一个迭代器 val top3ByGroup rdd .map(s (s.subject, s.score)) .groupByKey() .mapValues(iter iter.toList.sortBy(-_).take(3)) // 方案二用TreeSet维护最大值只保留Top3 import scala.collection.mutable val top3ByTreeSet rdd.aggregateByKey(mutable.TreeSet.empty[Int])({ (set, v) if (set.size 3) set v else if (v set.head) (set v).last else set }, (s1, s2) (s1 s2).last) val top3Optimized top3ByTreeSet.mapValues(_.toList.sortBy(-_))方案二看起来绕但原理非常简单用一个容量为3的TreeSet保持当前最大三个值来一个数就比一下插不插进去。这样不管某个科目有多少条记录内存里始终只保留3个分数。groupByKey会把这个科目所有分数全部放到一个迭代器里如果参与排序的数据量是百万级别内存开销差距就变得非常大。在实际项目中我还遇到过更极端的情况某个组有上亿条数据groupByKey直接把executor搞OOM。所以你现在可能觉得aggregateByKey难读我建议把这段代码多敲几遍配合Spark UI看stage的shuffle字节数你会发现方案二的shuffle数据量明显更小。能看懂这行代码说明你对聚合的理解已经不只是表层了。4.2 项目5订单与用户关联——RDD Join的常见陷阱这项目模拟电商场景有一份用户表、一份订单表需要算每个用户的消费总额。先看基础版本val userRdd sc.textFile(data/users.txt) .map(parseUser) .map(u (u.id, u)) val orderRdd sc.textFile(data/orders.txt) .map(parseOrder) .map(o (o.userId, o.amount)) val result userRdd.join(orderRdd) .map { case (userId, (user, amount)) (user.name, amount) } .reduceByKey(_ _)这里有个坑我当年亲自踩过以为join完就结束了结果输出的是每个用户的订单明细不是总额。原因是当一个用户有多条订单时join会把用户和每一条订单分别组合一次产生多行记录。所以join之后必须接一步reduceByKey把金额聚合成总额才算完成需求。进阶优化的思路是如果用户表很小、订单表很大join会把小表复制到每个分区去匹配大表走一遍完整的Shuffle。更合理的方式是先把小表collect成Map再用广播变量带到每个executor在map端直接查表val userMap sc.broadcast(userRdd.collectAsMap()) orderRdd.mapPartitions { iter val users userMap.value iter.flatMap { case (userId, amount) users.get(userId).map(name (name, amount)) } }.reduceByKey(_ _)广播变量加mapPartitions是我处理“大表关联小表”最常用的手段。它能完全避免Shuffle作业跑起来速度提升非常明显。这里有一个容易忽略的细节collectAsMap把数据拉到Driver端之后driver内存要能放得下这个小表所以只适合真正意义上的小表如果你广播一个几个GB的表内存直接爆掉。4.3 项目6自定义分区器——解决热点分区问题分区器平时用得少但一旦遇到数据分布不均它就是撬动性能的杠杆。默认的HashPartitioner按key.hashCode取模某些key的hash值如果恰好集中数据就会倾斜到同一分区。自定义Partitioner可以按业务规则控制数据落点。import org.apache.spark.Partitioner class RangePartitioner(partitions: Int) extends Partitioner { override def numPartitions: Int partitions override def getPartition(key: Any): Int key match { case s: String if s.startsWith(A) 0 case s: String if s.startsWith(B) 1 % numPartitions case _ 2 % numPartitions } } val partitionedRdd rdd.partitionBy(new RangePartitioner(3))这个示例偏教学真实场景里你会按更有业务含义的规则来分区。比如用户ID是数字可以根据ID段划分日志按日期字符串划分到对应分区下游按分区读取时非常方便。分区的时候有个组合操作值得一提repartitionAndSortWithinPartitions可以一次完成分区加排序比先repartition再sort少走一次Shuffle数据量大时效率差距很明显。关于分区数量我的建议是跟executor核心数挂钩一般控制在核心数的2到3倍比较合适。分区太少并行度上不去一个task处理的数据量太大分区太多task调度开销和序列化成本反而拉低性能。新人最容易犯的毛病就是无脑设置1000个分区觉得并行度越高越好实际经常适得其反。5. 进阶必做性能调优与数据倾斜实战5.1 项目7共享变量与累加器——统计与广播的正确姿势这项目的目标很简单搞懂累加器怎么统计脏数据广播变量怎么避免重复传输。先讲累加器。很多人会在foreach里直接写var count 0; count 1然后在Driver端打印count发现总是0。原因是foreach里的闭包会被序列化发送到executor执行Driver端的count变量只是个副本executor改的是自己的副本Driver当然看不到。正确做法是用Spark自带的累加器val badLineCounter sc.longAccumulator(badLines) rawRdd.foreach { line if (line.startsWith(#) || line.trim.isEmpty) { badLineCounter.add(1L) } } println(sbad lines: ${badLineCounter.value})累加器是唯一允许在executor端更新、在Driver端读取结果的变量它在设计上就是只写不读的executor之间也不能看到彼此。这个特性让它天然适合做全局计数比如统计清洗过程中丢弃了多少条非法记录。我在做数据质量监控时经常定义好几个累加器分别统计缺失字段数、解析失败数、金额异常数作业跑完直接打印一份数据质量报告很有用。广播变量在4.2里已经演示过这里再补一个注意事项广播变量的数据量不要太大。如果要广播几十GB的RDD那等于每个executor都存一份完整副本内存直接爆炸。它是用来优化“小表广播”的适用范围一般是几MB到几十MB的量级。判断标准就一条广播传出去的数据量必须远远小于每次task通过网络拉取的重复数据量否则不划算。5.2 项目8二次排序与多字段排序——进阶排序技巧排序需求在业务里很常见但RDD没有直接提供“二次排序”算子。比如我想按用户分组组内按金额降序、再按下单时间升序排列就需要把多个字段组合成排序键。case class Order(userId: String, amount: Double, ts: Long) val sortedRdd orderRdd.sortBy(o (o.userId, -o.amount, o.ts))Scala的元组排序是字典序第一位相同再比第二位。把金额取负数放进元组第二位就能实现金额降序时间放在第三位就是最后的排序维度。这个方法简单直接适合单级或多级排序场景。如果需要按分区内部排序输出强烈推荐repartitionAndSortWithinPartitions。它会先在分区阶段对数据排序再让下游读取。和先repartition再sort相比少了一次Shuffle数据量越大优势越明显。我做用户分群分析时经常用它按用户ID重分区同时保证每个分区内按时间有序下游处理连续操作时特别方便。还有一个容易忽略的排序细节如果排序字段中大多是重复值而你是想取TopN优先用rdd.sortBy(...).take(n)因为take可以用堆来避免全排序全量sort在数据量很大时代价高得多。很多新人不知道这个区别在几十亿数据上跑sortBy然后take(10)等于白白浪费一次全量排序。5.3 项目9数据倾斜处理——从定位到解决数据倾斜是Spark作业最常见的问题也是我项目里花篇幅最多的一章。定位方式简单跑作业时打开Spark UI进入某个Stage页面如果某个task的运行时间远长于其他task而且输入数据量明显偏大基本就是倾斜了。处理倾斜先看是不是“聚合类倾斜”。如果是某些key的数据量异常大可以试试两阶段聚合也就是先把倾斜key加随机前缀打散到不同分区做一轮局部聚合再去掉前缀做一轮全局聚合import scala.util.Random val skewedRdd orderRdd.map { case (cat, amount) val prefix if (cat 爆款品类) Random.nextInt(10) else 0 ((prefix, cat), amount) }.reduceByKey(_ _) .map { case ((prefix, cat), sum) (cat, sum) } .reduceByKey(_ _)思路很好理解加盐让原本集中的key分散到10个分区各自算一遍这步把压力分摊去掉盐之后每个key已经变成一个小聚合结果全局聚合的压力就小很多。这个方案对reduceByKey和aggregateByKey的倾斜都有效。但要注意倾斜如果发生在join场景这个方法就不适用了。因为join不能随便加盐加了盐会把原本匹配的key拆开导致数据错乱。处理join倾斜要复杂一些通常的思路是先把热点key找出来热点key的数据走广播map端join非热点数据走普通join最后union到一起。我在项目9的完整代码里给了这个实现核心思想就是“大表热点逻辑拆分小表热点广播匹配”。数据倾斜还有一个容易忽略的来源partitionBy之后写入分区如果某个分区数据量特别大同样会造成末尾task长时间运行。这种问题用自定义分区器就能解决回头看4.3的方案就行。要记住倾斜不只有聚合一种形态先定位到具体Stage再选对应的解法。6. 综合实战两个贴近业务场景的RDD程序6.1 综合项目A网约车订单数据清洗这个项目是很多公司招聘时喜欢考的题目给一批网约车订单原始数据字段非常不规范要求输出一份干净数据。我设计的清洗流程分四步解析、过滤、补全、输出。def parseLine(line: String): Option[OrderRecord] { val cols line.split(,) if (cols.length 7) None else Some(OrderRecord( cols(0), cols(1), cols(2).toDouble, cols(3).toLong, cols(4), cols(5), cols(6))) } val cleanRdd rawRdd .flatMap(parseLine) .filter(r r.amount 0.0 r.amount 5000.0) .filter(r r.isValidTime)这里有两个核心经验。第一用flatMapOption而不是mapfilter因为Option天然把解析失败变成空结果flatMap会自动丢掉None代码逻辑干净很多不用额外写try catch。这在RDD处理脏数据时是非常好用的模式。第二业务过滤条件要结合领域知识。网约车订单金额大于0且小于5000超过这个区间基本可以判断为异常数据或测试数据时间字段解析失败说明记录本身有问题直接淘汰。真实场景里不同业务的“合理范围”不一样你要懂一点业务才能定好这个边界。数据里我故意放了几条重复订单号所以还需要去重。这时候需要注意完全一样的重复好处理rdd.map(_.orderId).distinct()就行但真实业务里重复数据往往不完全相同比如订单时间差了十几秒这时候要先定义“什么算重复”再决定保留哪条。RDD阶段写这种精细逻辑比DataFrame更顺手你有充分的控制权。6.2 综合项目B农产品价格数据分析第二个综合项目从“分析农产品批发价格”出发需求是按品类计算月均价再算环比涨幅输出涨幅最高的品类。这项目把前面的解析、分组、排序、TopN全串起来了。case class Price(date: String, category: String, price: Double) val monthPrice priceRdd .map(p ((p.date.substring(0, 7), p.category), p.price)) .aggregateByKey((0.0, 0))( { case ((sum, cnt), price) (sum price, cnt 1) }, { case ((sum1, cnt1), (sum2, cnt2)) (sum1 sum2, cnt1 cnt2) } ) .mapValues { case (sum, cnt) sum / cnt } val momChange monthPrice .map { case ((month, cat), avg) (cat, (month, avg)) } .groupByKey() .flatMap { case (cat, values) val sortedValues values.toList.sortBy(_._1) sortedValues.sliding(2).map { case List((m1, p1), (m2, p2)) (cat, m2, (p2 - p1) / p1) } }这里我特意没有用groupByKey来算月均价而是用aggregateByKey保存(累计值, 条数)。原因和之前TreeSet的例子一样groupByKey会把一个品类下所有价格都拉到迭代器里而aggregateByKey在map端就能预聚合数据量大了之后效率差异非常明显。这个意识一定要建立起来。环比涨幅算完以后你还可以做一步趋势延伸输出按月份排序的CSV文件用任意BI工具画一张趋势图。RDD阶段不需要接数据库先把“算得出来、算得对”完成。等到后面学Spark SQL你会发现同样的需求写得更短但底层的理解已经足够了。综合项目最大的价值是逼你把前面的知识揉在一起用。跑通这两个项目你已经具备了独立写一个中等复杂度RDD作业的能力。剩下的就是多看、多写、多调优。7. 踩坑笔记运行RDD作业时的高频问题7.1 内存问题OOM和任务迟迟不结束RDD程序最容易挂在内存问题上。我把高频场景整理成一张排查表方便你以后快速对照现象常见原因解决方向Executor OOM单分区数据量过大、缓存策略不当增加分区数、改用MEMORY_AND_DISK或序列化存储Driver OOMcollect()拉取的数据量太大用take(n)抽样观察结果、分批处理Task运行时间极长数据倾斜或分区数过小两阶段聚合、自定义分区器GC时间占比高对象创建过多、堆内存紧张开启off-heap、减少大对象缓存这里面最容易被忽略的是collect()。本地小数据跑没问题生产环境一换大表collect会把全量数据拉到Driver端内存直接打满作业卡死。我建议新手在任何场景下都先用take(100)看结果确认逻辑没问题再考虑全量输出。这一条能帮你躲过一大半运行期事故。7.2 序列化与依赖问题RDD程序里最常见的异常是org.apache.spark.SparkException: Task not serializable。原因通常是你在map等转换操作里引用了Driver端的某个普通对象而这个对象没有实现Serializable。排查办法很简单看报错里提到哪个类把它改成实现Serializable或者换成更轻量的数据结构。Scala的case class默认实现了Serializable所以我会尽量用case class封装中间数据。另一个高频问题是运行环境缺少依赖报NoClassDefFoundError。用IDEA本地跑经常遇到因为SparkContext启动时jars没带全。解决办法是运行时把Spark安装目录的jars加入ClassPath或者用assembly插件打成fat jar再提交。如果你在公司集群提交作业别忘记用--jars把json4s这类第三方库带上去。7.3 输出路径已存在和分区数异常的坑写文件系统时Spark要求输出目录不存在。同一段代码跑第二次会直接报FileAlreadyExistsException。处理办法是在作业开头统一清理旧输出目录再执行主逻辑。本地文件系统可以用你熟悉的命令在脚本里我习惯用一个函数封装删除逻辑避免重复写。还有一个很隐蔽的问题如果你最后用了coalesce(1)结果文件会合并成一个所有数据压到同一个task上数据量大时极易卡死。如果只是想去掉大量小文件我建议用coalesce配合适当数量而不是直接降到1。输出文件数量其实影响下游读取的并行度这个细节在生产环境里经常踩雷做数据仓库的同学应该深有体会。这10个项目跑完我最深的体会是RDD编程不是比谁算子背得熟而是比谁能准确判断数据在分布式环境里是怎么流动的。你在本地跑通所有代码只是第一步真正有价值的是把每个例子的参数改一改、把Spark UI打开看看亲眼看一看Shuffle的字节数和task耗时。我自己也是写到第4个、第5个项目时才突然想明白reduceByKey和Shuffle的关系。建议你用配套数据跑完后再找一份真实业务数据试一遍遇到倾斜和异常记录就回头看这篇的排查表。代码和测试数据都在对应章节里把路径换成你自己的就行祝一次跑通。
返回列表