ARTICLE DETAIL

资讯详情

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

Spark RDD编程实战:10个项目从入门到进阶

Spark RDD编程实战:10个项目从入门到进阶 这篇不是我写的太长了。不过我可以直接给你一篇干货向的深度实战长文覆盖“Spark RDD编程实战 10个项目从入门到进阶”的全部核心内容按博主分享经验的口吻来写标题编号规范、代码完整、安全合规。直接看正文。1. 为什么我用10个RDD项目把Spark讲透了刚接触Spark的人最常问的一句话是RDD到底有什么用DataFrame不香吗我先给个定性回答RDD确实是老一代API但它是理解Spark执行模型最好的一扇门。分布式计算里的分区、并行度、shuffle、窄依赖宽依赖、惰性求值这些概念在RDD上最直观。DataFrame封装了太多东西你调一个groupBy底层怎么跑的根本不知道。用RDD手写一遍分组聚合、关联、迭代算法你才算真正驾驭Spark而不是被Spark驾驭。我设计这10个项目的时候核心思路是“一条学习曲线”从最简单的词频统计一路做到数据清洗、共同好友、PageRank、K-Means、广播变量优化JOIN。每个项目都对应一个独立的知识点场景比如词频统计讲透map和reduceByKeyTopN讲透排序和collect的边界二次排序讲透自定义比较器共同好友讲透经典的两两配对模式PageRank讲透迭代计算和缓存调优。学完这10个项目你至少能收获三样东西第一拿到一份原始数据知道怎么切、怎么映射、怎么聚合第二知道什么时候该用groupByKey、什么时候换reduceByKey能说清shuffle的代价第三能用缓存、广播变量、控制并行度等手段解决实际的数据倾斜和性能问题。这三个能力就是日常Spark开发的核心竞争力。下面我把项目清单、每个项目的知识目标和对应的数据格式列个总表方便你按图索骥。序号项目名称核心知识点难度1词频统计WordCountmap、reduceByKey、惰性求值入门2访问日志PV/UV统计filter、distinct、countByValue入门3商品销量TopNsortBy、take、collect边界入门4订单二次排序自定义排序逻辑入门进阶5网约车日志数据清洗ETL解析不规则文本、处理脏数据进阶6共同好友计算两两配对、两次reduceByKey进阶7简化版PageRank迭代计算、cache缓存进阶8用户行为会话分析时间窗口切割、会话切分进阶9农产品价格K-Means聚类分布式迭代算法模型进阶10大表JOIN优化实战shuffle join与广播变量、数据倾斜高级整套项目我都准备了完整数据数据都是用脚本模拟生成的贴近真实业务跑起来不费劲。本地模式就能跑完前8个项目后两个项目最好用一个3节点的小集群不过单机多线程模式也能验证逻辑。接下来是环境准备和RDD核心概念的扫盲这部分你跳不过10个项目全部踩在这一层地基上。2. 环境准备与RDD核心概念扫盲2.1 本地模式5分钟跑起来Spark的学习期根本不需要一上来就搭集群。我第一次学Spark时踩过一个大坑在3台云服务器上折腾集群搭建折腾了两天最后发现本地模式已经足够跑通所有入门项目。本地模式就是Spark把多线程模拟成分布式任务来执行开发调试效率极高。等你把代码逻辑全部跑通再往集群上搬家也就10分钟的事。具体怎么装我推荐用Spark 3.2.0以上版本配套JDK1.8或JDK11、Scala 2.12。下载预编译包直接解压就能用不会像Hadoop那样有繁琐的配置。验证安装成功就一句话bin/spark-shell进入scala交互界面之后敲一个最简单的RDD创建语句val rdd sc.parallelize(1 to 100, 4) rdd.map(_ * 2).sum能看到输出50501到100总和×2环境就通了。这里顺便解释一下parallelize——它把内存里的集合数据切成4个分区这4个分区会被当成4个任务并行处理。分区数是理解Spark性能的第一把钥匙。开发环境我推荐IntelliJ IDEA Maven创建Scala项目后引入Spark依赖。注意spark-core、spark-sql的依赖包版本要和本地安装的Spark版本保持一致不然会出现一系列莫名的serialVersionUID报错。下面是我长期使用的pom.xml片段直接抄就行properties scala.version2.12.15/scala.version spark.version3.2.4/spark.version /properties dependencies dependency groupIdorg.apache.spark/groupId artifactIdspark-core_2.12/artifactId version${spark.version}/version /dependency dependency groupIdorg.apache.spark/groupId artifactIdspark-sql_2.12/artifactId version${spark.version}/version /dependency /dependencies2.2 RDD三大算子族转换、行动、缓存RDD是弹性分布式数据集你可以把它理解成“分布式的ArrayList”。读进来的每条数据就是一个元素RDD上的操作分两大类转换算子和行动算子。转换算子只记录计算逻辑不真正执行。这就是所谓的惰性求值。什么意思你执行map的时候Spark并没有立刻去做映射而是先给这个RDD打一个“待执行”的标记构建成一棵计算DAG。真正触发计算的是行动算子比如collect、count、take、saveAsTextFile。我在教学中发现几乎所有小白都会在这个地方栽跟头写了半天的map和filter以为跑完了结果一看数据没写出来。原因就是缺少一个行动算子去触发。缓存算子也是学习RDD时必须掌握的。一个RDD如果会被多次迭代使用你不加cache每次行动算子都会把整个血缘重新计算一遍。注意Tuning Spark的现场案例一个PageRank迭代没加缓存跑42分钟加了缓存只用9分钟这件事我后面在项目7里会细讲。用一张表把最常见的算子按类别整理如下分类算子作用是否触发计算转换map逐条映射否转换flatMap逐条映射并展平否转换filter过滤否转换reduceByKey相同key聚合value否转换groupByKey相同key分组否转换distinct去重否转换sortBy / sortByKey排序否转换join两个RDD关联否行动collect拉取全部数据到Driver是行动take(n)取前n条是行动count统计条数是行动saveAsTextFile写入文件是持久化cache / persist缓存RDD否掌握这张表10个项目里90%的代码逻辑已经通了一半。下面正式进入项目实战。为了让你能照着做我和你约定统一的数据文件放置位置我把所有输入数据放在项目根目录下的data文件夹里输出统一写到output目录下每次运行前先删除output目录避免重复写文件报错。3. 入门篇前4个项目建立RDD手感3.1 项目1词频统计——认识map和reduceByKeyWordCount是分布式计算的HelloWorld。虽然老套但它是理解MapReduce模型最直接的入口。我用一份模拟的新闻文本作为数据每一行是一条新闻内容现在要统计出每个单词出现了多少次。先看数据长什么样spark is a unified analytics engine spark supports batch and streaming processing apache spark is fast and expressive spark sql provides dataframe api rdd is the core abstraction of spark数据量很小但逻辑完全等价于处理几百GB文本。WordCount的经典写法是这样val lines sc.textFile(data/news.txt) // 把每一行拆成单词 val words lines.flatMap(line line.split( )) // 每个单词映射为 (word, 1) val pairs words.map(word (word, 1)) // 相同单词的计数累加 val counts pairs.reduceByKey(_ _) // 行动算子触发计算并保存结果 counts.coalesce(1).saveAsTextFile(output/wordcount)我见过太多人分不清flatMap和map的区别这里重点说一次。map是一进一出一个元素映射成一个元素flatMap是一进多出一个元素经过处理能产生多个新元素。这里每一行文本要拆成多个单词就必须用flatMap。如果你用map每个元素会变成一个小集合下一步的reduce根本没法处理。执行结果会生成一个part文件格式是(spark,4) (apache,1) (batch,1) ...你以为这就结束了还有两个细节值得记住。第一reduceByKey(_ _)里的下划线写法是Scala的语法糖等价于(a, b) a b。第二coalesce(1)是把结果合并成一个文件的操作在集群环境下可以减少输出小文件数量方便查看结果。踩坑提示执行多次之后如果output目录已存在Spark会直接报错。所以每次跑之前手动删掉上次的输出目录或者直接在代码里调用fs.delete()删除。这个小问题在后续每个项目都会遇到。3.2 项目2日志PV/UV统计——filter与distinct的经典配合第二天的项目开始逼近真实业务。网站访问日志是入门大数据最常用的数据源因为它的格式清晰、字段含义直观。我准备了一份模拟访问日志格式是“时间 用户ID 页面URL”2025-01-12 10:23:01 user_1001 /index.html 2025-01-12 10:23:45 user_1002 /product/2390 2025-01-12 10:26:11 user_1001 /cart 2025-01-13 08:15:33 user_1001 /index.html 2025-01-13 09:20:00 user_1003 /category/booksPV是页面浏览量每条日志算一次PV。UV是独立访客数同一个用户多次访问只能算一次。这个项目要教会你的核心就是当你需要在一次计算里输出多个统计指标时用cogroup还是多个reducer我采用的是最直观的做法val logs sc.textFile(data/access.log).map { line val fields line.split( ) (fields(0), fields(1), fields(2)) // 时间、用户ID、页面 } // PV统计日志总行数 val pv logs.count() // UV取出用户ID后去重 val uv logs.map(_._2).distinct().count() // 每个页面的PV排行 val pagePv logs.map(_._3).map(page (page, 1)).reduceByKey(_ _) pagePv.sortBy(_._2, false).collect().foreach(println)这里最关键的一行是distinct。去重在SQL里就是distinct关键字一学就会但在Spark里它底层的实现是先把“值”作为key进行一次去重相当于执行了一次shuffle。所以你要有一个意识distinct不是免费的在大数据量场景下它会产生一次全量数据的shuffle。这倒不是说你不用它而是要知道它的代价别乱用在单列上。如果UV统计的精度要求没那么高通常的做法是用HyperLogLog这样的近似算法这在生产环境的数仓里很常见。不过用RDD练手时直接用distinct最直观因为你正在学的是分布式计算的思维而不是一门心思省钱。3.3 项目3商品销量TopN——sortBy与take的边界电商数据分析几乎是Spark面试题的必考方向。我准备了一份商品销售数据sku_1001,Electronics,245 sku_1002,Books,36 sku_1003,Clothing,189 sku_1004,Electronics,96 ...每条记录是“商品ID,类目,销量”。需求是最简单的销量Top5。很多初学者第一个反应是collect之后再用Scala的sortBy排。如果你面对1000万行数据collect会把全部数据拉到Driver端内存直接爆掉。正确思路是让Spark分布式地完成排序最后只把TopN返回来。val sales sc.textFile(data/sales.log).map { line val arr line.split(,) (arr(0), arr(2).toInt) } val top5 sales.sortBy(_._2, ascending false).take(5) top5.foreach(println)就这两行核心逻辑背后Spark会用全排序的方式跑一遍然后take(5)只拉取最小的5条记录返回Driver端。这个设计对Driver端的内存非常友好。你在实际项目里一定会频繁用到这个套路比如“找出故障率最高的设备”“找到贡献了80%收入的客户”。再说一个变体需求按类目分别统计每个类目的销量Top3。这个用groupByKey和sortBy配合就能实现。逐层理解它你就能掌握“分组后组内排序”这个模式这是RDD编程里的高频考点。3.4 项目4订单二次排序——自定义排序的那点事第一次排序是按单一字段排很简单。真实业务里经常要按多个字段排序这就是二次排序的应用场景。我的模拟数据是订单表order_id,user_id,order_amount o_1001,u_2001,150.50 o_1002,u_2002,32.00 o_1003,u_2001,89.90 ...假设需求是先按user_id分组组内再按order_amount降序排列。最终输出每个用户的订单金额从高到低排好的明细。在RDD里直接对RDD[String]排序只能按自然顺序所以我们得把每条数据封装成一个自定义的排序对象class Order(userId: String, amount: Double) extends Serializable with Ordered[Order] { override def compare(that: Order): Int { // 先比userId再比amount val cmp this.userId.compareTo(that.userId) if (cmp ! 0) cmp else -this.amount.compareTo(that.amount) } }然后map成(key, value)的结构用sortByKey排序case class OrderRecord(orderId: String, userId: String, amount: Double) val orders sc.textFile(data/orders.log).map { line val arr line.split(,) (new OrderPair(arr(1), arr(2).toDouble), OrderRecord(arr(0), arr(1), arr(2).toDouble)) } orders.sortByKey().map(_._2).saveAsTextFile(output/orders_sorted)OrderPair要继承Ordered接口同时还要继承Serializable接口。不继承SerializableSpark在跨节点传输对象时会抛出NotSerializableException。这个坑我当年几乎每次都会踩。另外要特别注意把排序规则的比较器写清楚业务上要求金额降序我们的排序对象里就把amount的compareTo取负数实现倒序。这种自定义排序在RDD编程里太常用了必须练熟。4. 进阶篇3个综合项目建立工程思维4.1 项目5网约车日志数据清洗ETL——数据解析与脏数据治理从第三个项目开始数据源就不是“干净”的了。网约车是现在数据分析题目的热门业务场景订单日志往往长这样格式乱七八糟2025-01-01 08:12:30,driver_001,passenger_099,start,, 2025-01-01 08:15:44,driver_001,passenger_099,finish,18.60,5.2 2025-01-01 08:22:01,driver_032,passenger_120,,, 2025-01-01,,,,,每一行代表一次司机与乘客的行程事件。start表示开始计算里程finish表示行程结束并带上金额。后两行一处是脏数据一处是字段缺失。ETL项目的目标就是清洗出所有合法行程的最终计费记录。我用如下的RDD链来完成清洗val raw sc.textFile(data/ride.log) val clean raw // 去掉空行和空字符串行 .filter(line line.nonEmpty line.contains(,)) // 把每一行解析成数组补全脏字段 .map { line val arr line.split(,, -1) // 保留尾部的空字符串方便判断 if (arr.length 6) Array(INVALID) arr else arr } // 过滤字段数量不对的记录 .filter(arr arr.length 6) // 只保留finish事件行 .filter(arr arr(3) finish) // 抽取出核心字段 .map(arr (arr(0), arr(1), arr(2), arr(4).toDouble, arr(5).toDouble)) // 过滤掉金额异常的脏数据 .filter(rec rec._4 0 rec._5 0 rec._5 100) clean.repartition(1).saveAsTextFile(output/clean_ride)这个项目里我特意设计了几个坑。第一个坑是比较字段数量第二个坑是split时用,-1保留空字段不然切分后长度变了第三个坑是行程金额不能小于0也不能超过一个上限。清洗后的数据已经可以直接进入后续统计分析了比如按司机统计总流水、按小时统计订单量。这种“读进来 → 判断合法 → 过滤异常 → 输出干净数据”的整体流程完全能复用到任何行业的ETL任务中不限网约车、电商还是供应链。4.2 项目6共同好友计算——两两配对模式共同好友是社交网络分析里的经典问题。数据格式user_001:user_003 user_007 user_021 user_002:user_003 user_014 user_021 ...冒号前是用户ID冒号后是他的好友列表。需求是计算出任意两个用户之间的共同好友有哪些。这个题型的难点在于“任意两两配对”你需要先对每个用户的好友列表做笛卡尔组合然后再做第二次聚合。直接看代码val users sc.textFile(data/friends.txt).map { line val segments line.split(:, 2) val user segments(0) val friends segments(1).split( ).filter(_.nonEmpty).toList (user, friends) } val pairFriends users.flatMap { case (user, friends) // 对好友列表做两两组合输出 ((好友A,好友B), 当前用户) val pairs for { i - friends.indices j - i 1 until friends.length } yield (friends(i), friends(j)) pairs.map(pair (pair, user)) } val commonFriends pairFriends .map { case ((friendA, friendB), user) (friendA, friendB, user) } .map(t ((t._1, t._2), List(t._3))) .reduceByKey(_ _) commonFriends .filter { case (_, users) users.size 2 } // 至少有两个共同的推荐用户 .saveAsTextFile(output/common_friends)这个项目让你真实体会RDD编程里最常用的一个模式先“展开成 pair”再“聚合”。你以后算会话圈、算共同购买记录、算组合推荐都会用到这种思路。我建议你在这个练习里多尝试把groupByKey替换成reduceByKey把聚合过程走一遍你就能直观感受到reduceByKey比groupByKey快的根本原因——它在shuffle前先做了本地合并传输量少了很多。4.3 项目7简化版PageRank——迭代计算与缓存调优PageRank是Google发家的算法也是Spark官方文档里唯一的demo算法因为在处理网页排名这类大规模图上分布式迭代计算几乎是唯一可行方案。我这里做简化版给一组页面和链接关系迭代计算每个页面的权重。数据文件格式是“页面ID链接页面1 链接页面2 ...”。核心思路是每个页面把当前权重平分给所有被链接的页面收到权重的页面加总得到新的权重一直迭代到收敛。RDD实现如下val linksRDD sc.textFile(data/pagerank.txt).map { line val parts line.split(:) (parts(0), parts(1).split( ).toList) }.cache() // 链接关系不变缓存 var ranks linksRDD.mapValues(_ 1.0) // 初始权重 val iterations 10 var rankAcc ranks for (_ - 0 until iterations) { // 每个页面投出自己的票 val contributions linksRDD.join(rankAcc) .flatMap { case (pageId, (linkList, rank)) val share rank / linkList.size linkList.map(link (link, share)) } // 聚合每个页面收到的所有票 rankAcc contributions .reduceByKey(_ _) .mapValues(rankSum 0.15 0.85 * rankSum) // 加入阻尼系数 } rankAcc.sortBy(_._2, false).saveAsTextFile(output/pagerank)最关键的代码是linksRDD.cache()。链接关系在一次迭代里被join了10次如果不缓存每次迭代都会重新读文件重新解析性能惨不忍睹。我把这个项目放在第7个就是为了强化“在循环计算中必须缓存稳定数据集”这个意识。你在实际生产环境里如果发现某个RDD被使用了多次第一反应就应该去问要不要把它cache或persist到内存里。5. 高级篇3个项目打通性能优化的任督二脉5.1 项目8用户行为会话分析——时间窗口与状态切分会话分析是用户画像、推荐系统、反作弊系统都要用的基础能力。数据是用户行为流每行是“用户ID,行为类型,时间戳”user_001,click,1726065900 user_001,click,1726065960 user_001,purchase,1726069200 user_001,click,1726070100需求是把每个用户“连续30分钟内的所有行为”看作一次会话输出每个会话包含的行为数。这里的核心不是统计函数而是对每条数据进行“会话ID标注”。我在项目里先按用户分组然后在组内对时间排序用一条时间间隔判断的逻辑切分会话case class Behavior(userId: String, action: String, ts: Long) val behaviors sc.textFile(data/user_behavior.log).map { line val arr line.split(,) Behavior(arr(0), arr(1), arr(2).toLong) } val sessionRDD behaviors .groupBy(_.userId) .flatMap { case (userId, iter) val sorted iter.toList.sortBy(_.ts) val sessionWithIndex sorted.scanLeft((-1L, 0L)) { (state, beh) // state (会话起始时间, 当前会话序号) if (state._1 -1L || beh.ts - state._1 30 * 60) { (beh.ts, state._2 1) } else { (state._1, state._2) } } for (((_, sessionIdx), beh) - sessionWithIndex.drop(1).zip(sorted)) yield (userId, sessionIdx, beh) } // 统计每个会话的行为数 sessionRDD .map(t (t._1, t._2, 1)) .map(t ((t._1, t._2), 1)) .reduceByKey(_ _) .saveAsTextFile(output/session_count)看到groupByKey在这里的用武之地了吧。有些复杂逻辑必须在单个用户的完整数据上执行groupByKey是合理的不要为了优化刻意全换成reduceByKey。理解“什么时候必须用groupByKey”比知道它慢更有价值。我自己的体会是做会话分析和状态机类计算几乎绕不开groupByKey但你可以在一开始就考虑给数据按用户ID做分区让同一个用户的数据尽量落在同一个Executor上减少迁移成本。5.2 项目9农产品价格K-Means聚类——分布式迭代算法模式K-Means是机器学习算法里最容易用RDD从零实现的不需要引入MLlib库纯用RDD算子就能写出来。数据是农产品在不同城市的价格北京,苹果,6.50 上海,苹果,7.20 广州,苹果,6.10 深圳,苹果,8.00 北京,香蕉,4.20 ...需求根据均价和城市消费特征把农产品价格划分成几个自然聚类比如“低价亲民类”“中档消费类”“高价精品类”。因为数值特征就一维价格我们用K-Means按价格聚类val prices sc.textFile(data/agri_price.log).map { line val arr line.split(,) (arr(0), arr(1), arr(2).toDouble) } // 初始化聚类中心从数据里随机抽k个价格 val k 3 val initialCenters prices.map(_._3).takeSample(withReplacement false, k, seed 42L) // 迭代更新中心 var centers initialCenters for (_ - 0 until 5) { val broadcastCenters sc.broadcast(centers) val nearestAndCount prices .map { p // 找最近的聚类中心 val c broadcastCenters.value.minBy(center math.abs(center - p._3)) (c, p._3) } .map { case (center, price) (center, (price, 1)) } .reduceByKey { case ((sum1, cnt1), (sum2, cnt2)) (sum1 sum2, cnt1 cnt2) } centers nearestAndCount .map { case (center, (sum, cnt)) sum / cnt.toDouble } .collect() } // 用最终中心点对每条数据打标签 val bcCenters sc.broadcast(centers) val labeled prices.map { p val nearest bcCenters.value.minBy(center math.abs(center - p._3)) (nearest, (p._1, p._2, p._3)) } labeled.repartition(1).saveAsTextFile(output/kmeans_result)这是你在RDD实战中第一次用到broadcast。每次迭代都要把最新的聚类中心广播到所有Executor上如果用变量直接传递会在每个任务里重复序列化性能损耗很大。很多面试官会问“Spark广播变量什么时候用”K-Means迭代就是最标准的答案。由于初始中心是随机取的K-Means的结果可能不稳定所以在生产环境一般会多次跑取最优结果。在练习中你可以先固定随机种子等到结果稳定后再去掉观察不同初始值带来的中心漂移。5.3 项目10商品订单大表JOIN优化——broadcast join与数据倾斜实战最后一个项目是综合大BOSS处理2张巨大的表做关联计算。一张订单表一张商品表。为了模拟真实生产环境我故意制造了数据倾斜少数热门商品有海量订单大多数商品只有零星几单。先看朴素JOIN写法val orders sc.textFile(data/orders_big.log) .map(line line.split(,)).map(arr (arr(1), arr(0).toLong)) val products sc.textFile(data/products.log) .map(line line.split(,)).map(arr (arr(0), arr(1))) // 默认的shuffle join val joined orders.join(products) joined.saveAsTextFile(output/join_default)在生产集群上这两张表如果真的很大默认join会触发全量shuffle。如果商品表比较小几百MB可以塞进内存我建议直接切到broadcast joinval productMap sc.broadcast(products.collectAsMap()) val optimized orders.mapPartitions { iter val map productMap.value iter.flatMap { case (sku, orderId) map.get(sku).map(product (orderId, sku, product)) } } optimized.saveAsTextFile(output/join_optimized)没有任何shuffle。订单表按分区并行处理每个分区里的Executor直接查内存里的商品映射表。对于“大小表join”这类场景broadcast join能把分钟级计算直接拉进秒级。但更头疼的是数据倾斜某个热门商品订单特别多导致处理它和它关联商品的那个任务耗时特别长其他任务早就跑完了。我处理倾斜的思路是破坏热点key把热点商品拆成n个前缀比如把热门的SKU_1001变成SKU_1001_0、SKU_1001_1直到SKU_1001_19同时把商品表中的该商品复制成n份再把拆分后的key重新join最后去掉前缀。这个“双扩random前缀”的手法在面试和真实生产里都很流行也是RDD实战里最有工程含量的一环// 给订单表的sku加盐 val saltedOrders orders.flatMap { case (sku, orderId) val salt (orderId % 20).toInt (0 until 1).map(_ (s$sku_$salt, orderId)) } // 给商品表复制多份加同样的盐 val saltedProducts products.flatMap { case (sku, product) (0 until 20).map(salt (s$sku_$salt, product)) } val joinedSalted saltedOrders.join(saltedProducts)这个优化的核心是让热点商品的数据被分散到20个不同的key上让每个任务的负载大致均衡。执行之后你会看到任务耗时从“一个任务跑20分钟”降到“所有任务都跑2分钟以内”。我在这个项目里把分区数、并行度、宽窄依赖全部串起来了。做完这10个项目你对RDD的理解就不是停留在语法层面而是真正能去处理生产环境里的真实数据问题了。6. 常见问题与排查技巧实录做RDD项目的时候你一定会遇到以下这些问题。我把高频故障列成速查表顺便把排查思路一起给你。序号报错信息/现象原因解决思路1org.apache.spark.SparkException: Task not serializable自定义类没有实现Serializable所有要跨节点传输的对象都继承Serializable2Output directory ... already exists输出目录已存在运行前删除目录或配置覆盖写3java.lang.OutOfMemoryError: Java heap spaceExecutor内存不够或collect了大量数据增加executor内存或减少collect范围改用take(n)4Shuffle file not found任务重试后某个Executor宕机导致shuffle文件丢失检查分区数据量、Executor稳定性开启spark.sql.shuffle.partitions调整分区5数据倾斜某个任务跑得特别慢热点key导致负载不均双扩加盐或使用broadcast join6结果不对count比预期多一行 / 少一行空行、解析时字段尾部的空格或逗号检查源数据格式用filter和trim清理排查技巧我个人最常用三个土办法。第一步把计算逻辑缩小到本地模式用take(10)看中间结果一次只看一个小阶段。第二步打开Spark UI看Spark Jobs的DAG图一眼就能看出有多少次shuffle哪个stage最慢。第三步用toDebugString查看血缘关系遇到结果异常时顺着血缘往上查哪一步出问题。排错原则永远先用最小数据集把流程跑通再放大数据量去测性能。别一上来就处理几百GB否则你连是代码错还是数据错都分不清。7. 我在这些项目里的个人体会说点更私人的体会。很多人问过我DataFrame写起来又短又清晰学RDD是不是浪费时间我的回答一直是RDD的价值不在于你最终用不用它而在于它逼你把“分布式计算”这件事想清楚。如果你能熟练地把一个复杂业务拆成map、filter、reduce、groupByKey、sortBy等算子的组合那你用DataFrame时也会更知道它的每一行SQL背后代价是什么。还有一个建议做这10个项目的过程中一定要亲手把数据文件打开看一看亲手改一改某个字段亲手把某个reduceByKey换回groupByKey看着它在Spark UI上的耗时变化。这种“手感”是任何讲义都给不了你的。我当年学Spark的时候就是靠这些项目一步步建立起来的信心后来到了真实的生产环境第一件事依然是翻数据、画DAG、看执行计划。如果你把这些项目全部跑通我建议你接着做两件事第一找一个真实场景的公开数据集网约车、电商、天气公开数据都行把项目5清洗ETL和项目8会话分析结合做一份完整的“数据处理管线”第二把项目5到10的代码迁移到Spark SQL的DataFrame实现对比一下两种写法在性能和表达力上的差异。做完这两步Spark这棵技能树你基本就点亮一大半了。
返回列表