BDAS:EC2 + PySpark,以及 DataCamp 那 14 分
这是四次迭代里唯一有环境风险的一次:SSH 密钥、安全组、实例类型、Spark 版本、Java 版本,任何一处不通都会让你在第 10 周晚上抓头发。所以这一章的第一条建议在第 4 章就说过,这里再说一次:第 8 周实验课当天就把环境跑通,跑一个 count() 截图存证。剩下的两周你只需要翻译代码。
环境:先跑通,再写代码
课程会提供 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 |
| 5 | pyspark 或 SparkSession 起来 | 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 的转换是惰性的。filter、select、withColumn、join 全都只是在记账,什么都没算。只有动作(count、show、collect、write)才会触发真正的执行。
这解释了新手最困惑的一个现象:十行转换代码秒回,一个 .count() 卡三分钟。那三分钟是前面十行的账一起还的。
实操后果:不要在探索时反复 count()。如果一份中间结果要用多次,.cache() 它(并在用完 .unpersist())。
差异二:数据分在很多块上,跨块移动很贵
DataFrame 被切成若干分区,每个分区由一个 task 处理。操作分两类:
- 窄依赖(
filter、map、select):每个分区自己算完就行,不需要跨机器传数据。 - 宽依赖(
groupBy、join、distinct、orderBy):同一个 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.3 | 60.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,未发现明显倾斜」也算完成任务——重点是你知道要看。
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 迭代分,是整学期投入产出比最高的学习时间。
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 的使用边界,以及一张过完就能交的自检清单。