卷 VI · 交付CH 23工位 23/24

BDAS:EC2 + PySpark,以及 DataCamp 那 14 分

工位 Iteration 4 23% + DataCamp 14% 37% 合计 评分员在找:你真的在分布式环境里跑了(分区数、shuffle、执行时间),以及你诚实说明了「为什么这个规模需要/不需要 Spark」。

这是四次迭代里唯一有环境风险的一次:SSH 密钥、安全组、实例类型、Spark 版本、Java 版本,任何一处不通都会让你在第 10 周晚上抓头发。所以这一章的第一条建议在第 4 章就说过,这里再说一次:第 8 周实验课当天就把环境跑通,跑一个 count() 截图存证。剩下的两周你只需要翻译代码。

EC2 / SSH惰性求值分区与 shuffle数据倾斜DataCamp

环境:先跑通,再写代码

课程会提供 AMI 或指引,具体镜像名和实例类型按课程实验课为准。但无论细节怎么变,跑通的顺序是固定的,而每一步都要留一张截图

#动作验证方式(截图)常见卡点
1启动 EC2 实例(用课程指定 AMI)实例状态 running + 公有 IP区域选错(AMI 只在特定区域可用)
2安全组开 22(SSH)与 Jupyter 端口安全组入站规则截图只对自己的 IP 开放,不要 0.0.0.0/0
3密钥权限 chmod 400 key.pem → SSH 连上终端提示符Windows 上权限问题 → 用 WSL 或 PowerShell 的 icacls
4启动 Jupyter,本地端口转发浏览器里的 Notebook 页面token 找不到 → 看终端输出的 URL
5pysparkSparkSession 起来Spark UI(4040 端口)截图Java 版本不匹配
6读数据 → count()行数输出 + Spark UI 的 Jobs 页S3 路径权限 / 本地路径写错
# 本地 → EC2,并把 Jupyter 的 8888 转发到本地
chmod 400 infosys722.pem
ssh -i infosys722.pem -L 8888:localhost:8888 ubuntu@<公有IP>

# 实例上
jupyter notebook --no-browser --port=8888
# 复制输出里带 token 的 URL 到本地浏览器

# 别忘了:不用的时候 Stop 实例(不是 Terminate,否则磁盘上的东西没了)
# Stop 之后公有 IP 会变,下次连接要重新查
⚠ 三件会让你损失时间或钱的事

1. 忘记关实例。按小时计费,一个周末就能烧掉课程额度。每次做完立刻 Stop,并在 AWS 控制台设一个预算告警。

2. Terminate 了实例。Stop 保留磁盘,Terminate 会删掉(除非改过设置)。所有代码放 GitHub,数据放 S3 或本地也留一份——课程要求用 GitHub,正好利用它。

3. 所有工作只在实例上。把 Notebook 定期 git push写一份自己的 README(第 4 章说过):实例类型、AMI 名、SSH 命令、启动 Jupyter 的命令、数据路径。第 10 周的你会感谢第 8 周的你。

三处思维差:pandas 用户最容易栽的地方

差异一:什么都没跑,直到你要结果

Spark 的转换是惰性的。filterselectwithColumnjoin 全都只是在记账,什么都没算。只有动作(countshowcollectwrite)才会触发真正的执行。

这解释了新手最困惑的一个现象:十行转换代码秒回,一个 .count() 卡三分钟。那三分钟是前面十行的账一起还的。

实操后果:不要在探索时反复 count()如果一份中间结果要用多次,.cache() 它(并在用完 .unpersist())。

差异二:数据分在很多块上,跨块移动很贵

DataFrame 被切成若干分区,每个分区由一个 task 处理。操作分两类:

  • 窄依赖filtermapselect):每个分区自己算完就行,不需要跨机器传数据
  • 宽依赖groupByjoindistinctorderBy):同一个 key 的数据必须先聚到一起,要把数据写到磁盘、通过网络传给别的机器再读回来——这叫 shuffle,是 Spark 里最贵的操作,而且它会把执行计划切成两个 stage。

Demo 的默认设置(1 200 万行、8 个分区、一次 groupByKey)给出:2 个 stage、1 次 shuffle、504 万行经过 shuffle、估算 2.7 分钟。切到「join 再 group」会看到 3 个 stage、2 次 shuffle、1 008 万行过网络。

差异三:倾斜——一个热键就能毁掉整个作业

把 Demo 里的「数据倾斜」滑块从 0 拖到 0.3(意思是 30% 的数据落进同一个 key):

倾斜最大分区 / 平均分区估算耗时
0(均匀)1.19×2.7 分钟
0.360.8×77.4 分钟

为什么这么夸张:一个 stage 的耗时等于它最慢那个 task 的耗时。shuffle 之后默认有 200 个分区,如果 30% 的数据全落在其中一个分区里,那一个 task 要处理的量是平均的 60 倍——其余 199 个分区早就跑完了,你还在等它。

怎么发现:Spark UI 的 Stages 页里看 task 耗时分布,如果 Max 远大于 Median,就是倾斜。这张截图放进 BDAS 报告,是「我真的在分布式环境里跑过」最有力的证据。

怎么缓解:给热键加随机盐再两阶段聚合、把小表 broadcast 掉、或者干脆在 groupBy 前先过滤掉那个异常 key。报告里能写出「我们检查了 task 耗时分布,Max/Median = X,未发现明显倾斜」也算完成任务——重点是你知道要看。

◆ 关于 reduceByKey 与 groupByKey:本 demo 诚实的地方

Demo 里把 groupByKey 换成 reduceByKey,shuffle 次数和记账行数是一样的——因为它是一个按行数记账的模型。

但在真 Spark 上,reduceByKey 通常快好几倍,原因是它会先在每个分区内做局部聚合(map-side combine),只把聚合后的少量结果发出去;而 groupByKey 把原始记录全部发过网络。差别不在「传了多少行」,在传了多少字节

这也是 DataFrame API 比 RDD API 更值得用的理由:Catalyst 优化器会自动帮你做这类优化(谓词下推、列裁剪、局部聚合)。所以 BDAS 那次用 DataFrame/Spark SQL,不要用 RDD——除非你想在报告里专门讨论两者差异(那会是个不错的加分点)。

PySpark 的九步骨架

from pyspark.sql import SparkSession, functions as F
from pyspark.ml import Pipeline
from pyspark.ml.feature import Imputer, StringIndexer, OneHotEncoder, VectorAssembler, StandardScaler
from pyspark.ml.classification import LogisticRegression, DecisionTreeClassifier
from pyspark.ml.evaluation import BinaryClassificationEvaluator, MulticlassClassificationEvaluator

spark = (SparkSession.builder.appName('INFOSYS722-BDAS')
         .config('spark.sql.shuffle.partitions', '200')   # 记下这个数,报告要写
         .getOrCreate())

## 2.1 读数据
df = spark.read.csv('s3a://.../u5mr_raw.csv', header=True, inferSchema=True)
print('行数', df.count(), '分区数', df.rdd.getNumPartitions())   # 两个数都写进报告

## 2.2 / 2.4 描述与质量
df.describe().show()
df.select([F.count(F.when(F.col(c).isNull(), c)).alias(c) for c in df.columns]).show()
df.groupBy('region').count().orderBy('count', ascending=False).show()   # 暴露大小写不一致

## 3.2 清洗
df = (df.withColumn('region', F.initcap(F.trim(F.lower(F.col('region')))))
        .withColumn('vaccine', F.when(F.col('vaccine') == 999, None).otherwise(F.col('vaccine')))
        .withColumn('gdp_pc', F.when(F.col('gdp_pc') > 100000, F.col('gdp_pc') / 1000)
                               .otherwise(F.col('gdp_pc')))
        .dropDuplicates(['region', 'year']))          # 业务键去重

## 3.3 构造
df = df.withColumn('log_gdp', F.log10('gdp_pc'))

## 7.1 划分(注意:randomSplit 不分层,不平衡数据要手动分层)
train, test = df.randomSplit([0.7, 0.3], seed=42)

## 3.2 / 3.5 / 6:整条 Pipeline(所有统计量只在 train 上 fit)
stages = [
    Imputer(strategy='median', inputCols=num_cols, outputCols=num_cols),
    StringIndexer(inputCol='region', outputCol='region_idx', handleInvalid='keep'),
    OneHotEncoder(inputCols=['region_idx'], outputCols=['region_ohe']),
    VectorAssembler(inputCols=num_cols + ['region_ohe'], outputCol='raw_features'),
    StandardScaler(inputCol='raw_features', outputCol='features', withMean=True, withStd=True),
    LogisticRegression(labelCol='label', featuresCol='features', maxIter=100),
]
model = Pipeline(stages=stages).fit(train)

## 7.2 / 8.1 执行与评估
pred = model.transform(test)
auc = BinaryClassificationEvaluator(labelCol='label', metricName='areaUnderROC').evaluate(pred)
f1  = MulticlassClassificationEvaluator(labelCol='label', metricName='f1').evaluate(pred)
print(f'AUC={auc:.3f}  F1={f1:.3f}')
pred.groupBy('label', 'prediction').count().show()      # 混淆矩阵

## 8.2 增益表(第 19 章的分箱,用窗口函数)
from pyspark.sql.window import Window
w = Window.orderBy(F.desc(F.element_at('probability', 2)))
(pred.withColumn('decile', F.ntile(10).over(w))
     .groupBy('decile')
     .agg(F.count('*').alias('n'), F.sum('label').alias('hits'))
     .orderBy('decile').show())

## 9.3 保存整条 Pipeline(含预处理参数)
model.write().overwrite().save('s3a://.../model_v1')

三个必须写进报告的数字:getNumPartitions() 的结果、spark.sql.shuffle.partitions 的设置、以及从 Spark UI 抄下来的执行时间与 shuffle 读写量。这三个数字是「你真的用了大数据工具」的唯一硬证据。

✎ 如果你的数据其实很小(多数人是这样)

不要假装它大。按第 3 章那段模板诚实写「当前规模不需要 Spark,采用它是为了验证可扩展性」,然后做一个真正有说服力的动作:把数据人为放大若干倍,测量执行时间随规模的变化。

import time
for k in [1, 10, 100]:
    big = df
    for _ in range(k - 1):
        big = big.union(df)          # 复制放大
    t0 = time.time()
    big.groupBy('region').agg(F.avg('u5mr')).collect()
    print(f'{k}×({big.count():,} 行)耗时 {time.time() - t0:.1f} 秒')

三个数据点画一条线,配一句「耗时随规模近似线性增长,说明方案在 100 倍规模下仍可行」。这就是 LO2 里「understand the key components of the computing environment for Big Data」的直接证据,比任何空话都强,而且只要五分钟。

DataCamp:14 分,只需要花时间

25 小时以上,占 14%,第 7 周(9 月 14 日)截止。用 uni 邮箱登课程指定的 DataCamp 环境。按对这门课的用处排序,建议这样刷:

优先级方向为什么时数
★★★PySpark / Big Data with PySpark直接就是 Iteration 4 要用的东西。先刷这个,第 9 周会省下几小时8–10 h
★★★Supervised Learning with scikit-learn直接是 Iteration 3 的主线;Pipeline 与交叉验证那两章尤其对得上第 13、17 章6–8 h
★★pandas 数据清洗对应 Step 3。如果你 pandas 已经很熟,可以跳4 h
★★Tableau对应 Iteration 3 的可视化交付4 h
Power BI顺带覆盖不计分的 MSAS,一份时间两处用3 h

两条纪律:从第 1 周开始,每周 2.5 小时——它是唯一一项「攒不起来」的作业,因为时数按实际学习时间记;② 优先刷 PySpark,因为它同时是 14% 的 DataCamp 分和 23% 的 BDAS 迭代分,是整学期投入产出比最高的学习时间。

▣ 本章交付物 —— Iteration 4 交什么

1. 报告(Step 1–8):Step 1–2、5 复用;3–4、6–8 用 PySpark 重做;加一节「三条产线结果对比」(ISAS / OSAS / BDAS 的指标差异及原因——划分方式、默认参数、编码差异、Imputer 无 KNN 等)。

2. 环境证据:EC2 实例截图、SSH 连接截图、Spark UI 截图(Jobs 页 + Stages 页的 task 耗时分布)

3. 三个数字:初始分区数、spark.sql.shuffle.partitions、执行时间与 shuffle 读写量。

4. 可扩展性证据:1×/10×/100× 的耗时表 + 一句结论。

5. GitHub 仓库:Notebook、README(含环境搭建步骤)、数据获取脚本或说明。报告里给出仓库链接。

6. Pipeline 保存的证据model.write().save())——对应 Step 9.1/9.3 的部署与维护。

这一章的一句话

Spark 的三处思维差是惰性求值(动作之前什么都没跑)、分区与 shuffle(宽依赖切 stage 且要过网络)、以及倾斜(一个热键让作业从 2.7 分钟变 77 分钟,因为 stage 的耗时等于最慢那个 task);而如果你的数据其实不大,就诚实写出来,然后用 1×/10×/100× 的耗时表证明可扩展性——那比假装数据很大得分高。

最后一章:那 17% 的研究论文、第 12 周的数据伦理、生成式 AI 的使用边界,以及一张过完就能交的自检清单。