ARTICLE DETAIL

资讯详情

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

Spark大数据课后习题全解:从RDD、SQL到Streaming与调优实战

Spark大数据课后习题全解:从RDD、SQL到Streaming与调优实战 第一次系统地把《Spark大数据技术与应用》的课后习题从头刷到尾是带两个刚转方向的同学的时候。他们卡住的地方出奇一致书上的概念看得懂例题跟着敲也能跑一到课后习题就不知道从哪下手。更麻烦的是写完了也不知道对不对——没有标准答案、没有判分、更说不清为什么这么写。这三个不知道其实指向同一件事习题真正考的不是语法记忆而是在信息不完备的情况下把一段业务描述翻译成分布式计算逻辑的能力。Spark、大数据这套东西落到工程里核心就是这件事。这篇内容把这套翻译过程摊开讲题眼怎么找、环境怎么搭才不出错、RDD 与 Spark SQL 的题分别用什么套路解、Streaming 和综合案例题怎么串成链路、调优与内存类题目怎么自己给自己判分以及最后怎么把一堆零散习题拼成能写进简历的完整项目。适合正在啃这本教材的在校生也适合从 Java、Python、数据库方向转过来、需要尽快把 Spark 用起来的人。基础好不好都能看代码部分我会把每一步的意图讲清楚。1. 先看清题目在考什么再动手敲第一行代码大部分人做课后习题的顺序是错的翻到题、打开 IDE、开始写。正确顺序是先花十分钟把这一章的题目通读一遍判断它们落在哪个能力层上。教材的习题看着零散其实分类很清楚不同类型题目的验证方式完全不同用同一种方式对待必然低效。1.1 用章节分布反推一张知识地图《Spark大数据技术与应用》这类教材的章节编排基本遵循同一条主线Spark 概述与运行架构、开发环境与运行模式、RDD 编程基础、RDD 进阶键值对算子与持久化、Spark SQL、Spark Streaming、MLlib、GraphX、性能调优。习题的分布也跟着走但权重大头通常在 RDD 编程和 Spark SQL 这两块一般能占到全部题目的六成以上。这不是偶然——这两块是日常写 Spark 作业代码时出场率最高的部分。我给这三类题型做了一张对照表做之前先给题目归个类能省掉大量无用功题型典型问法真正考察的能力怎么判断自己做对了概念与架构题简述 DAG 与 Stage 的划分依据Driver 和 Executor 各自负责什么对执行模型的理解程度能用自己的话讲给不懂的人听懂不靠背诵算子与代码题用 RDD 实现某统计写出某指标的 SQL把业务描述翻译成算子链结果可复现换一份数据仍然对调优与配置题数据倾斜怎么处理内存参数如何设置定位问题设计对照实验能给出改动前后的可量化对比这里有个很实用的判断标准凡是你能在五分钟内靠回忆答出来的题都不算掌握只能算记住了。真正的掌握体现在换一个问法你还能答。比如题目问的是统计每个单词出现的次数你把代码背下来了但如果我问统计每个单词在多少个不同的文件里出现过需要把 value 去重再计数你能不能立刻改出来改不出来就说明前面的题白做了。1.2 概念题的坑全在话术正确但理解错误概念题是课后习题里最容易被轻视的一类。很多人翻两页书找一段话抄上去就算完事。问题在于教材上的表述是浓缩过的你抄的时候并不理解它省略了什么。举一个几乎每本 Spark 教材都会出现的经典问题reduceByKey和groupByKey有什么区别标准答案通常是reduceByKey会在 map 端做预聚合网络传输量更小性能更好。这句话没错但如果你只记住这句面试官接着问那你知道 map 端预聚合具体发生在哪个阶段吗、为什么groupByKey不能做预聚合你就答不上来了。完整理解应该包含这几层reduceByKey的聚合函数满足结合律所以同一条记录在同一个分区内可以先做一次局部合并这一步发生在 shuffle write 之前叫 map 端 combinegroupByKey只做分组不做聚合分区内的所有 value 必须原样攒进一个迭代器再整体通过网络发给下游所以它在网络开销和内存占用两方面都吃亏。真正需要说清楚的是代价。假设有 1 亿条记录、100 万个 key用groupByKey时每个分区要维护的 value 列表加起来就是全量数据量一旦某个分区的数据超过可用内存就会 spill 到磁盘。而reduceByKey在 map 端先把它压到接近 key 的数量级通常能减少一到两个数量级的数据量。这个数量级差异才是它被反复强调的原因不是更快这三个字。我自己的习惯是做概念题时强制自己多问三层它做了什么、什么时候做、不做会怎样。三层都能答才算过关。1.3 直接抄答案会在真项目里付出代价我见过太多同学的习题本写得很漂亮代码能跑出结果但被问一句你这个结果为什么是对的就沉默了。原因很简单他们从网上抄来的代码只保证了对某个特定输入输出正确而工程里没有特定输入这回事。举个具体的例子。很多版本的词频统计习题示例数据里不包含空行和标点符号所以一段简单的flatMap(_.split( )).map((_, 1)).reduceByKey(_ _)就能出正确答案。可一旦数据换成真实的日志文件里面混着空行、连续空格、制表符、大小写不一致的单词这段代码的输出就全是噪声空字符串会被当成一个高频 keySpark和spark会被算成两个词。所以做代码题时我的建议是自己给数据加脏手动往输入文件里插几行空行、加几个多余空格、把个别单词改成大写跑一遍看结果还对不对。这一步花不了五分钟但它能逼着你把清洗逻辑写进代码里而不是依赖数据的干净程度。2. 环境与数据准备习题跑不起来八成都卡在这一步课后习题的技术难度往往不高真正消耗时间的是环境。教材配套的版本和现在能下载到的版本之间隔着好几年直接照书上的命令敲大概率会遇到各种莫名其妙的报错。这一节把常见的坑集中说清楚。2.1 本地模式够用别一上来就搭集群一个很常见的误区刚学到第三章就去折腾分布式集群搭了两三天Spark 一行代码没写。做课后习题阶段单机本地模式足以覆盖九成以上的题目。本地模式同样会走完整的 DAG 划分、Stage 切分、shuffle 流程只是 Executor 线程和 Driver 在同一个 JVM 里你还能用调试器单步跟踪反而是学习执行模型最好的环境。启动方式很简单进到 Spark 解压目录后# Scala / Java 交互式环境使用本机全部 CPU 核心 ./bin/spark-shell --master local[*] # Python 交互式环境 ./bin/pyspark --master local[*] # 提交打包好的作业 ./bin/spark-submit \ --master local[4] \ --class com.example.WordCount \ --executor-memory 2g \ ./target/wordcount.jar input/ output/local[*]里的星号表示用满本机所有核心local[4]表示固定四个线程。这个数字不是随便写的它决定了并行 task 的数量上限。如果你本机是 8 核用local[1]跑你会看到 Spark UI 上永远只有一个 task 在跑误以为程序写得慢其实是并行度被自己锁死了。什么时候才需要伪分布式或 on YARN两种情况一是题目明确要求验证资源调度、Executor 动态分配、多应用共存二是数据规模大到单机内存放不下。除此之外本地模式跑不动的原因通常是代码写法问题不是集群问题。2.2 版本三角JDK、Scala、Spark 必须对齐这是新手最容易踩、报错信息又最不友好的坑。Scala 是编译型语言Spark 的二进制包是用特定版本的 Scala 编译的如果你的开发环境用另一个 Scala 版本编译代码运行时会出现NoSuchMethodError、NoSuchMethodError或者更隐蔽的java.lang.NoClassDefFoundError——代码编译期完全正常一提交就崩。对照关系大致是这样Spark 2.x 系列主要对应 Scala 2.11 和 JDK 8Spark 3.0 到 3.2 一般对应 Scala 2.12、JDK 8 或 11Spark 3.3 之后开始提供 Scala 2.13 的构建版本JDK 11、17 都可以用。教材如果写于 Spark 2.x 时期书上的 API 和参数名会和新版本有出入比如spark.sql.crossJoin.enabled这类参数在新版本里被调整过。我的建议很直接如果想严格跟书就装书对应年代的 JDK 8 Scala 2.11 Spark 2.4别折腾如果本来就想用新版本那就接受书上的部分示例需要微调但核心概念和绝大部分算子 API 是向后兼容的不影响学习。用 Maven 管理项目时依赖要写清楚 Scala 的二进制版本后缀dependency groupIdorg.apache.spark/groupId artifactIdspark-core_2.12/artifactId version3.2.4/version /dependency dependency groupIdorg.apache.spark/groupId artifactIdspark-sql_2.12/artifactId version3.2.4/version /dependency注意_2.12这个后缀它必须和你本机scala -version输出的前两位一致。如果用了 IDEA 的 Scala 插件还要检查 Project Structure 里指定的 Scala SDK 版本这里不对编译能过、运行必崩。2.3 数据自己造比到处找数据集靠谱教材习题用到的数据一般是几类网站访问日志、电商订单、学生成绩、超市销售流水。这些数据在网上能找到一些公开版本但字段格式往往和书上对不上改数据的时间比写代码还长。更快的做法是写个生成脚本自己造。用 Python 几十行就能生成几万到几十万条结构可控的数据import random from datetime import datetime, timedelta users [fu{str(i).zfill(5)} for i in range(1, 2001)] cities [北京, 上海, 广州, 深圳, 杭州, 成都] categories [手机, 电脑, 家电, 图书, 服饰, 食品] with open(orders.csv, w, encodingutf-8) as f: f.write(order_id,user_id,city,category,amount,order_time\n) base datetime(2024, 1, 1) for i in range(1, 200001): uid random.choice(users) # 故意让部分用户下单更频繁制造出天然的数据倾斜 if random.random() 0.05: uid u00001 t base timedelta(minutesrandom.randint(0, 60 * 24 * 90)) f.write(f{i},{uid},{random.choice(cities)},{random.choice(categories)}, f{round(random.uniform(10, 5000), 2)},{t.strftime(%Y-%m-%d %H:%M:%S)}\n)这段脚本里我特意加了一行让 5% 的概率把用户改成同一个 ID。这不是随手写的它会在后续做分组统计时制造出一个典型的数据倾斜场景。等你做到调优那一章的习题这份数据就直接变成了现成的实验材料不用再另外造。提示生成数据时保留一些脏特征——空值、格式不统一的时间字符串、异常大的金额。习题里的数据越干净你学到的越少。3. RDD 编程题的解题套路从题目描述翻译到算子链RDD 编程是课后习题的重头戏也是最容易出写得出来但写得不优雅问题的地方。我总结出一个四步翻译法几乎能应对所有 RDD 类题目。3.1 把一段中文题目拆成四个问题拿到题目后不要急着想算子先回答这四个问题输入是什么形态一行一条纯文本还是多字段的结构化数据还是已经是键值对这决定了第一步用textFile、map拆字段还是parallelize。要按什么分组分组键就是最终reduceByKey、groupByKey或aggregateByKey里的那个 key。很多题目的难点其实是分组键怎么构造比如统计每个城市每个品类的销售额key 就是一个二元组(city, category)。每一步的中间形态是什么这是最关键的一步。每写一个算子强迫自己说出现在这个 RDD 的每条记录长什么样。说不出来说明你在瞎写。最终要什么形态是要一个数字reduce、count、一个集合collect、take、还是一份落盘文件saveAsTextFile举个具体例子。题目统计每个城市销售额最高的三个品类及其金额。拆解过程是输入是 CSV 格式的订单明细第一步按逗号拆字段并取出城市、品类、金额三项中间形态是(city, category, amount)第二步把分组键构造成(city, category)值构造成(amount, 1)第三步用reduceByKey聚合出((city, category), (totalAmount, count))第四步需要按城市分组再排序这一步用groupByKey拿到city - List[(category, totalAmount)]对每个城市的列表排序取前三。写出来大概是这样val lines sc.textFile(orders.csv) // 跳过表头拆字段过滤掉字段数不对的脏行 val base lines.filter(!_.startsWith(order_id)) .map(_.split(,)) .filter(_.length 6) .map(a ((a(2), a(3)), a(4).toDouble)) val cityCategorySum base.reduceByKey(_ _) .map { case ((city, cat), amt) (city, (cat, amt)) } val top3 cityCategorySum.groupByKey() .mapValues(iter iter.toList.sortBy(-_._2).take(3)) top3.collect().foreach(println)注意这里groupByKey是可以接受的因为分组之后每个城市的品类数量有限不会撑爆内存。但如果分组后的 value 列表可能非常大就必须换写法。很多人学到groupByKey性能差之后一刀切地不用它反而写出了更别扭的代码。判断标准是看分组后每个 key 的 value 数量级而不是看算子名字。3.2 词频统计的六种变体练透一类题的通用解词频统计几乎是所有 Spark 教材的第一道代码题看起来太简单其实它可以派生出大量变体每一种都在考察不同的东西。我整理了常见的几种变体关键改动考察点基础词频无flatMapmapreduceByKey的基本组合忽略大小写与标点清洗前置数据预处理意识按词频降序取前 NsortBy(-_._2).take(N)全局排序与局部 top-N 的区别统计每个词出现在多少个文件里value 去重再计数distinct或Set的使用过滤停用词广播停用词表广播变量的使用时机按首字母分组输出二次分组多阶段聚合的思路其中统计每个词出现在多少个文件里这一题特别值得做因为它能暴露一个常见错误。很多人写成reduceByKey(_ _)那统计的是出现总次数不是文件数。正确思路是先构造(word, fileName)的去重对再按 word 计数val fileRDD sc.wholeTextFiles(input/*.txt) val wordPerFile fileRDD.flatMap { case (path, content) content.split(\\s).filter(_.nonEmpty) .map(w (w.toLowerCase, path)) }.distinct() // 同一个词在同一文件里只算一次 val wordFileCount wordPerFile.map { case (w, _) (w, 1) }.reduceByKey(_ _)distinct()这一步会触发 shuffle代价不低但这是准确性的必要条件。如果你在做题时只盯着能不能出结果很容易漏掉这一步然后得到一个看起来合理但完全错误的数字。3.3 宽窄依赖和 shuffle为什么你的作业越跑越慢做到进阶章节题目会开始问哪些算子会产生 shuffle。这个问题背后是一整套执行机制值得花时间吃透因为它直接决定你后面能不能看懂 Spark UI。窄依赖指父 RDD 的每个分区最多被一个子分区使用比如map、filter、flatMap、union数据不需要跨节点移动可以在同一个 task 里流水线执行。宽依赖指父 RDD 的多个分区数据要打散重组到多个子分区比如groupByKey、reduceByKey、join、sortByKey、repartition每个宽依赖都会切断流水线形成一个新的 Stage 边界。这就解释了一个现象同样的一段逻辑加一次repartition或者groupByKey运行时间可能翻好几倍。因为 shuffle 要把数据写到磁盘、通过网络传输、再读回来过程中还涉及序列化和反序列化。Spark UI 的 Stage 页面里那些耗时特别长、task 数量特别多的 Stage基本都是 shuffle 造成的。有个细节很多人忽略reduceByKey虽然也是宽依赖但它比groupByKey快原因前面讲过是 map 端预聚合。所以在需要先分组再聚合的场景下优先用reduceByKey、aggregateByKey、combineByKey这类带预聚合能力的算子而不是先groupByKey再map。3.4 缓存和持久化什么时候该用什么时候是浪费RDD 持久化是必考知识点但习题往往只问有哪些存储级别不问你什么时候该用。实际写代码时判断标准只有一个这个 RDD 会不会被多个行动算子重复使用。如果一份中间结果只被算一次加cache()毫无意义反而多占内存。如果同一个 RDD 会被反复访问比如在一个循环里对同一份数据做多次查询那就必须缓存否则每次行动算子都会从源头重新计算整条血缘链。val cleaned sc.textFile(access.log) .map(parseLine) .filter(_.isValid) .cache() // 后面会被统计多次值得缓存 val pv cleaned.count() val uv cleaned.map(_.userId).distinct().count() val top cleaned.map(x (x.url, 1)).reduceByKey(_ _).take(10)存储级别也要选对。默认的MEMORY_ONLY数据放不下时会丢弃分区下次用到重新计算MEMORY_AND_DISK会溢写到磁盘适合数据量偏大又需要反复访问的情况。序列化存储MEMORY_ONLY_SER省内存但费 CPU做反序列化有开销。我的经验是先用默认级别跑看 Spark UI 的 Storage 页签里 cache 的实际占用和命中率再决定要不要换成序列化。注意cache()是惰性的它只是打了一个标记真正发生缓存是在第一次遇到行动算子的时候。所以在cache()之后立刻用collect()把整个数据集拉回本地是双重错误——既把数据全量缓存了又把 Driver 内存打爆了。4. Spark SQL 与 DataFrame 类题目很多人在这里假装会了从 RDD 转到 Spark SQL很多人会有一种虚假的进步感SQL 写起来确实短代码量掉了一大截跑得也快。但习题里的 SQL 题目恰恰是最能区分会写 SQL和会写 Spark SQL的地方。4.1 从 RDD 到 DataFrame 的思维切换写 RDD 时你的思维是数据在节点之间怎么流动写 DataFrame 时思维要换成我要什么结果。这个切换看起来是语法差异实际是抽象层次的跃升。核心差别在于执行优化。DataFrame 的转换逻辑会被 Catalyst 优化器处理成逻辑计划再经过规则优化和代价模型生成物理计划最后才编译成 RDD 操作。这意味着你写的 SQL 和实际执行的算子可能完全不是一回事谓词下推会把过滤条件尽可能推到数据源附近列裁剪会只读取用到的列这些优化你在 RDD 层面得自己手动实现。举个直观的例子。用 RDD 读一个 CSVmap(_.split(,)).filter(_(2) 北京)整个文件的所有字段都被读进来再过滤。用 DataFrame 写成spark.read.csv(...).filter($city 北京).select(order_id, amount)如果底层是 Parquet 这类列式存储实际只会读取city、order_id、amount三列的对应数据块。数据量差几十倍的时候性能差异是量级上的。所以做 SQL 类题目时我建议强制自己回答一个问题我这段逻辑Catalyst 能帮我优化吗如果答案是不能比如逻辑藏在 UDF 里那就要考虑换个写法。4.2 窗口函数题一套模板打天下教材里 Spark SQL 的习题经常出现取每个用户最近一笔订单计算每个品类销售额的环比给订单按金额排名这类需求。这三类看着不同其实全部是窗口函数的固定套路。套路是三步定义分区PARTITION BY、定义排序ORDER BY、选择窗口帧默认是分区内从头到当前行或者整分区的聚合。取每个用户最近一笔订单SELECT user_id, order_id, amount, order_time FROM ( SELECT *, ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY order_time DESC) AS rn FROM orders ) t WHERE rn 1;累计与环比SELECT user_id, order_time, amount, SUM(amount) OVER (PARTITION BY user_id ORDER BY order_time ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS cum_amount, LAG(amount, 1) OVER (PARTITION BY user_id ORDER BY order_time) AS prev_amount FROM orders;这里有两个高频错误。第一ROW_NUMBER、RANK、DENSE_RANK混用。如果同一个用户存在两条时间完全相同的订单ROW_NUMBER只会留一条RANK会让两条都排第一都留DENSE_RANK也会都留。题目没说清楚时你要主动判断业务上该保留一条还是多条而不是随便挑一个。第二窗口帧忘了写。不写ROWS BETWEEN时默认帧是RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW它按值而不是按行括范围如果排序字段有重复值计算结果和你想的不一样。需要按行累计时必须显式写ROWS。4.3 数据源读写题格式陷阱比逻辑陷阱多读取 CSV 写入 Parquet读取 JSON 做嵌套字段展开按分区写入数据这类题目逻辑本身不难坑全在格式细节上。列几个我踩过的Schema 推断把数字变成字符串。spark.read.csv默认开启inferSchema它会扫描数据抽样推断类型但遇到混杂内容时就放弃了全按字符串处理。结果你写amount 100时比较的是字符串9 100会返回 true。稳妥做法是显式定义 schemafrom pyspark.sql.types import StructType, StringType, DoubleType, TimestampType schema StructType() \ .add(order_id, StringType()) \ .add(user_id, StringType()) \ .add(amount, DoubleType()) \ .add(order_time, TimestampType()) df spark.read.option(header, True).schema(schema).csv(orders.csv)时间格式不统一。数据里同时存在2024-01-01 10:00:00和2024/01/01 10:00直接转timestamp会得到一堆 null。用to_timestamp指定格式或者分情况处理别指望自动识别万能。写入时的覆盖行为。mode(overwrite)在动态分区写入的场景下会把你指定的整个分区目录清空而不只是覆盖这次涉及的分区——如果你的输出路径下有别的数据就一起没了。需要只覆盖对应分区时要配合spark.sql.sources.partitionOverwriteModedynamic使用。这个坑在没有版本控制的环境里杀伤力极大。小文件问题。一个 200 个分区的 DataFrame 写出去就是 200 个文件如果每个文件只有几十 KB下游读取时会因为大量小文件而变慢。习题里通常不考这个但只要你开始处理真实数据就必须处理。常见的做法是先repartition或coalesce控制文件数再写出。coalesce不触发 shuffle适合减少分区repartition会触发全量 shuffle适合增大分区或需要均匀分布。5. Streaming 与综合案例题把零散知识点串成一条链路到了教材后半段题目会从实现某个功能变成完成一个完整流程。这类综合题是最接近真实工作的也是最容易做偏的。5.1 微批处理的时间语义先搞清楚再写代码Spark Streaming 的题目绕不开几个概念批处理间隔、窗口、检查点、水位线。定义不难背难的是搞明白它们之间的相互关系。批处理间隔batch interval决定了多久切一次数据。设成 10 秒那每 10 秒会生成一个微批处理上一批累积的数据。这里有个反直觉的点批处理间隔不是处理耗时的上限。如果你的处理逻辑要跑 30 秒那批次会排队堆积你会看到调度延迟持续上升最终拖垮整个作业。所以间隔的设置依据是处理一批数据的平均耗时一般让它留出 50% 以上的余量。事件时间和处理时间的区别也常被考。事件时间是数据里自带的发生时间处理时间是数据到达系统的时间。网络延迟、上游故障都会让两者产生偏差。结构化流里的水位线watermark就是用来处理这种偏差的设一个 10 分钟的水位线意味着系统愿意等待迟到的数据最多 10 分钟超过这个时间的迟到数据会被丢弃。不设水位线状态会无限增长内存迟早撑不住。# 结构化流按事件时间开窗10 分钟水位线容忍迟到数据 result (spark.readStream .format(kafka) .option(kafka.bootstrap.servers, localhost:9092) .option(subscribe, orders) .load() .selectExpr(CAST(value AS STRING) AS json_str) .select(from_json(json_str, schema).alias(d)) .select(d.*) .withWatermark(order_time, 10 minutes) .groupBy(window(order_time, 5 minutes), category) .agg(sum(amount).alias(total))) query (result.writeStream .outputMode(update) .option(checkpointLocation, /tmp/ck/orders) .format(console) .start()) query.awaitTermination()注意outputMode的选择。append只输出新结果update输出有变化的结果complete输出全量结果。带聚合的查询用append时只有在水位线推进、窗口真正关闭后才会输出你会感觉数据怎么半天不出来这不是 bug是语义决定的。5.2 综合题的拆解顺序从需求到输出中间不要跳步综合案例题一般会描述一个业务场景比如统计某电商平台各品类的销售情况。题目给的描述往往只有两三行但要做的东西很多。我的拆解顺序固定为五步第一步明确指标口径。同一个销售额是按下单时间算还是按支付时间算要不要扣除退款含不含运费这些题目通常不写你要自己选一个合理的口径并在注释里写清楚。这一步决定了后面所有的代码。第二步确定数据模型。把输入数据的字段列出来标注每个字段的类型、是否可能为空、格式是否统一。第三步清洗。过滤掉金额为负或为零的记录、处理时间字段的格式、拆分或合并字段、去重。第四步生成明细宽表。这一步是把清洗后的数据整合成一张字段完整、粒度统一的表通常写成 Parquet 落盘方便后续反复查询。第五步在宽表上聚合指标输出结果。这五步看着啰嗦但它是把一道大题拆成五道小题的唯一办法。我在做综合题时会把每一步的输出都落盘并抽样看几行确认形态对了再往下走。一路写到底再运行是效率最低的做法因为你不知道自己是在哪一步错的。5.3 从复购率这类指标题看口径的重要性教材里有一类题目会要求计算某类业务指标比如复购率。这类题目的难点从来不是 SQL是口径。复购率的常见定义是在观察周期内购买两次及以上的用户数除以周期内有购买行为的用户数。听起来很清楚但落到代码里马上会遇到几个问题观察周期怎么定自然月、滚动 30 天同一用户同一天下两单算不算复购跨品类的购买算不算先给一个按自然月统计的实现WITH monthly AS ( SELECT user_id, DATE_FORMAT(order_time, yyyy-MM) AS mon, COUNT(DISTINCT order_id) AS order_cnt FROM orders WHERE amount 0 GROUP BY user_id, DATE_FORMAT(order_time, yyyy-MM) ) SELECT mon, COUNT(CASE WHEN order_cnt 2 THEN 1 END) AS repeat_users, COUNT(*) AS active_users, ROUND(COUNT(CASE WHEN order_cnt 2 THEN 1 END) / COUNT(*), 4) AS repeat_rate FROM monthly GROUP BY mon ORDER BY mon;这段代码里有几个值得注意的选择。用COUNT(DISTINCT order_id)而不是COUNT(*)是为了防止一条订单被拆成多行比如订单包含多个商品时被重复计数。过滤amount 0是为了排除退款或异常数据——如果不加退款用户会被算成活跃用户指标就被稀释了。我在实际项目里还遇到过一个更隐蔽的问题大促期间用户会在同一天集中下单用自然月统计会让大促月份的复购率虚高。解决办法是改用滚动窗口统计或者把复购定义成跨月复购即两个不同月份都有购买。哪种更合理取决于业务方想看什么但无论选哪种都要在产出结果时把口径写清楚否则同一份数据不同的人算出来的数字不一样最后扯皮的一定是你。6. 调优与内存类题目最难自己判分的一块调优题是课后习题里最特殊的一类没有标准答案也难以验证。你写了一个参数配置怎么知道它是不是对的这一节讲的就是怎么把调优题变成可验证的实验。6.1 内存模型的几个数字得能算出量级Spark 1.6 之后引入了统一内存管理理解它只需要记住几个数字之间的关系。每个 Executor 的 JVM 内存里有一部分是预留的大约 300MB用于内部对象和元数据剩下的部分才是统一内存池这个池子由spark.memory.fraction控制默认 0.6也就是可用内存的 60%。统一内存池内部又分执行内存和存储内存两块初始比例由spark.memory.storageFraction决定默认 0.5也就是各占一半。关键在于这个比例不是硬边界。存储内存不够时可以借用执行内存的空闲部分执行内存紧张时甚至可以驱逐存储在内存里的数据前提是没有被持久化到磁盘。所以调参的逻辑不是把某一块调大而是根据作业特征是偏计算比如大量 join、聚合还是偏缓存比如反复查询同一份数据来调整。举个实际算例。一台 16 核 64GB 的机器要跑一个 Spark 作业怎么分配假设留 8GB 给系统和其他进程剩 56GB 给 Spark。如果spark.executor.cores设成 4可以起 4 个 Executor每个约 14GB减去 300MB 预留统一内存池大约 8.2GB。这个量级心里要有数因为当你说我给每个 Executor 16GB时你应该能算出实际可用于计算的内存是多少。顺便说一句 Executor 核心数的选择。设得太少单个 Executor 内并行度不足设得太多单个 Executor 的 JVM 内存压力大GC 时间会明显上升。经验值是 3 到 5 之间4 是最常见的起点。这背后还有个 HDFS 客户端吞吐的考虑——每个核心对应一个并发读写流核心数过多时反而会因为并发流太多导致 I/O 效率下降。6.2 数据倾斜识别、定位、处理三步走数据倾斜是调优题里出现频率最高的主题因为它太常见了。前面第 2.3 节那份造的数据如果你拿去按用户聚合就会看到明显的倾斜——u00001这个用户占了 5% 的记录处理它的那个 task 会比其他 task 慢几个数量级。识别倾斜的方法很直接打开 Spark UI进到有问题的 Stage看 task 的耗时分布。如果大部分 task 在 2 秒内结束少数几个跑了 5 分钟那就是典型的倾斜。另一个信号是 shuffle read size 的分布个别 task 读到几百 MB 而其他只有几 MB。处理手段按优先级排列第一先看能不能过滤。如果热点 key 是空值或测试数据直接在这一步过滤掉成本最低。第二小表广播。如果是 join 引发的倾斜且其中一张表足够小默认阈值 10MB由spark.sql.autoBroadcastJoinThreshold控制让 Spark 把整张小表广播到每个 Executor把 shuffle join 变成 map join倾斜自然消失。第三加盐打散。如果大表 join 大表且必须处理热点 key就得用加盐。思路是给热点 key 拼上一个随机后缀比如u00001变成u00001_0到u00001_9把一条记录拆成十条分散到不同分区同时把小表对应的记录也复制十份保证 join 能匹配上。这个做法会增加数据量所以只对热点 key 加盐不要全量加盐否则得不偿失。from pyspark.sql import functions as F SALT 10 # 只对出现次数超过阈值的 key 加盐 hot_keys df.groupBy(user_id).count().filter(F.col(count) 100000) \ .select(user_id).rdd.flatMap(lambda r: [r[0]]).collect() hot_bc spark.sparkContext.broadcast(set(hot_keys)) F.udf(string) def add_salt(uid): if uid in hot_bc.value: return f{uid}_{random.randint(0, SALT - 1)} return uid salted df.withColumn(salt_key, add_salt(user_id))第四让优化器帮忙。Spark 3.x 的自适应查询执行AQE可以自动处理一部分倾斜把spark.sql.adaptive.enabled和spark.sql.adaptive.skewJoin.enabled打开它会根据运行时的 shuffle 统计信息自动拆分过大的分区。这是成本最低的尝试应该放在手动加盐之前。6.3 调优题怎么自己给自己判分这是我最想说的一点没有对照组调优就是玄学。做调优题的完整流程应该包含四步构造基线、设计单变量实验、采集指标、解释差异。基线就是用默认参数跑一遍记录总耗时、各 Stage 耗时、shuffle 读写量、GC 时间。然后每次只改一个参数其他保持不动再跑一遍对比。指标从哪来Spark UI 是最主要的来源。Jobs 页面看整体耗时Stages 页面看每个阶段的时间构成包括 shuffle write、shuffle read、GC、task 反序列化等细分项Environment 页面确认参数是否真的生效——这一步很多人漏掉因为参数名写错了不会报错只会静默忽略。举个我自己做过的例子。有一次处理一份 20GB 的数据默认配置下跑了 24 分钟。改了三处把spark.sql.shuffle.partitions从默认的 200 调到 800把 join 的小表改成广播开了 AQE。再跑一次是 9 分钟。这个结果里改动三处是没法归因的。后来我逐个回退测了一遍发现其中两处贡献很小主要收益来自广播 join。如果当时只记录调完了快了 15 分钟那我就什么也没学到。所以做这类题时我的建议是把每次实验的参数、耗时、关键指标都记到一张表里。这张表本身就是你最有价值的笔记面试时能拿出来讲比背一堆参数名有说服力得多。7. 把习题集变成项目集从做完到做好刷完一本教材的习题你手里会有一堆散落的代码片段。这些东西如果不整合过两个月就忘光了也很难在求职或汇报时拿出来展示。最后这一节讲怎么改造。7.1 单个题目怎么扩成端到端链路一道统计各城市销售额的题往上加三层就是一条完整链路第一层是数据接入。原来代码里读的是本地 CSV改成从 Hive 表或对象存储读取加上数据源配置和错误处理。这一步会让你接触到连接配置、权限、读取失败重试这些真实问题。第二层是调度与增量。原来是一次性跑全量改成按天增量处理加上调度脚本。这时候要处理的核心问题是数据的幂等性——重跑一次同一天的作业结果会不会翻倍这就逼着你理解overwrite和append、分区覆盖的区别前面第 4.3 节讲的那些坑会立刻变成你的实际问题。第三层是质量校验。在输出结果之前加几个检查主键有没有重复、金额有没有为负、行数有没有剧烈波动、关键字段的空值率是多少。校验失败就中断并告警。这一层看起来和 Spark 无关但它是把能跑的脚本变成能交付的作业的分界线。做完这三层你手上就不是一道题了而是一个可以用三分钟讲清楚的项目。讲的时候重点不是我用了什么算子而是我遇到了什么问题、怎么定位、为什么这么解决。7.2 面试官会顺着哪几条线追问做过一遍习题的人在面试里最容易被追问的还是那几条主线。我列一下做习题时可以带着这些问题反向检查自己的理解Stage 是怎么划分的宽依赖出现在哪里如果同一个 RDD 被多次使用会发生什么shuffle 在物理上是怎么实现的数据写到哪里为什么 shuffle 是性能瓶颈内存不够时会发生什么spill 是怎么触发的怎么判断作业有没有频繁 GC数据倾斜怎么发现的如果加盐之后还是慢下一步查什么小文件问题怎么来的怎么解决coalesce和repartition该用哪个这几条线其实都能在教材的习题里找到对应题目只是题目比较孤立。你在做的时候如果能把它们串起来——比如因为宽依赖产生了 shuffle因为 shuffle 的数据分布不均导致了倾斜因为倾斜所以要考虑加盐加盐会增加数据量所以要配合调整分区数——那你就真正理解了这一章。7.3 一份我用了很久的复盘模板每做完一道有难度的题我会填这么几个字段题目要求、我一开始的思路、实际卡住的地方、最终方案、有没有更优解、如果数据量放大十倍会发生什么。最后一项特别有用它能逼着你跳出能跑就行的心态。举个填过的例子。题目要求对日志做用户行为路径分析我一开始用groupByKey把每个用户的行为收集起来再排序数据量小的时候完全没问题。填到放大十倍这一项时我发现热门用户的行为可能有几百万条一个 key 的 value 列表就能把 Executor 撑爆。后来改成用窗口函数在 SQL 层面处理让 Spark 自己做分段问题就消失了。这个模板的价值不在于记录而在于它强制你在做完之后多想一步。习题做完就忘是因为缺少这最后一步的加工而那些你花时间复盘过的题会在很久之后遇到类似场景时突然浮现出来。最后分享一个小技巧把整本书的习题按基础算子、SQL、流处理、调优四类各挑三道重新用最新版本的 Spark 写一遍把代码整理成一个仓库每道题配一段简短的说明注释。这件事花不了太多时间但它把学过变成了手里有东西这两者在后面的路上差别很大。
返回列表