01大数据算法与技术概述
Part A |课程纲要
A1. 课程定位 & 学习目标
-
做什么:理解并行计算基础,用 Python(
joblib、multiprocessing)并行化数据处理和机器学习任务(练习 KMeans)。 -
为什么:串行计算在大数据上慢,并行可以充分利用多核 CPU,加速训练或处理。
-
怎么做:
-
学会把任务拆分成可独立处理的子任务。
-
使用
joblib/multiprocessing并行化这些子任务。 -
观察加速效果,理解调度和 IO 开销。
-
-
例子:KMeans 的
n_init并行、多批样本 assignment 并行。 -
Tips:并行并不总是越多越快,要考虑 CPU 核数、内存和 IO 瓶颈。
A2. 评分构成与门槛
-
做什么:了解评分标准,确保报告和演示得分最大化。
-
为什么:报告、演示、作业、合作每块都有分,要有针对性。
-
怎么做:
-
报告 50 分:方法(并行点)、实验(基线 vs 并行)、分析(瓶颈、误差)、附录(代码)
-
演示 30 分:脚本清晰、跑通 demo、回答问题
-
作业/平时 10 分:提交、讨论
-
合作/沟通 10 分:Git 提交记录、分工表
-
-
例子:报告中写“用 joblib 并行 KMeans n_init,CPU 4 核,加速比 2.5×”。
A3. 时间线(可直接抄进日历)
-
做什么:规划 8 周学习任务。
-
为什么:保证每周有产出,按时完成实验和报告。
-
怎么做:
-
Week1:环境搭建(Python、joblib、scikit-learn、Jupyter)
-
Week2:选择数据、定义问题、提交计划
-
Week3:串行实现 KMeans,记录基线时间
-
Week4:joblib 并行化 n_init / 参数搜索
-
Week5:数据分片 assignment 并行
-
Week6:性能优化(线程数、缓存、结果稳定性)
-
Week7:报告草稿 + PPT
-
Week8:提交报告 + 现场演示
-
-
例子:Week3 拆成“实现 / 测试 / 记录结果”3 个任务块,写进日历。
A4. 官方建议的数据源 & 题材
-
做什么:选合适数据集做实验。
-
为什么:并行化任务需要足够大数据量才能看出加速效果。
-
怎么做:
-
数据源:Kaggle / UCI ML / OpenML
-
题材:
-
文本:新闻语料、日志 → TF-IDF / 文档聚类
-
电商:订单 → 商品共现
-
图像:批量特征提取
-
仿真:Monte Carlo → 天然并行
-
-
-
Tips:优先选“每个样本独立处理”的任务,容易并行化。
A5. 学术诚信 & 沟通
-
做什么:确保作业和报告符合规范。
-
为什么:学术不诚信会扣分,沟通记录是评分证据。
-
怎么做:
-
不抄袭、注明引用
-
保留 Git 提交记录、issue 日志、分工表
-
提供
requirements.txt或environment.yml
-
-
例子:报告附录写明 Python 版本、库版本、运行命令。
Part B |课堂内容
B1. 什么是“大数据”
-
做什么:理解大数据的特征。
-
为什么:决定任务拆分和并行策略。
-
怎么做:
-
3V:Volume(量大)、Velocity(速度快)、Variety(类型多)
-
一句话理解:不是更大的 Excel,而是需要改变处理范式,拆任务并行处理。
-
-
例子:1000 万条日志 → 拆成 1000 个文件,每个核处理一个文件。
B2. 数据按结构分四类
-
做什么:分类数据类型,决定处理方式。
-
为什么:不同结构适合不同并行策略。
-
怎么做:
-
结构化:表格 / CSV / 数据库 → 可直接 SQL 聚合
-
半结构化:JSON / XML → 解析字段
-
日志 / 时间序列:按时间追加 → 做窗口/聚合
-
非结构化:文本、图像 → 先特征化/向量化
-
-
例子:银行流水表 → 结构化;论坛帖子 → 非结构化。
B3. 数据存放三层
-
做什么:理解数据存放位置。
-
为什么:并行处理需知道 I/O 源。
-
怎么做:
-
本地 Sandbox:快速迭代
-
中间层 Cluster / Object Storage:并行处理
-
生产 / 服务层:数据库 / 流平台
-
-
例子:开发阶段在 Jupyter,稳定后迁移到 S3 + Spark 并行。
B4. 生态与角色
-
做什么:区分 BI 和 Data Science,理解团队角色。
-
为什么:明确职责,划分并行任务。
-
怎么做:
-
BI:报表可视化,结构化数据
-
DS:建模、预测
-
四类角色:Data Engineer / Data Scientist / ML Engineer / Analyst
-
-
例子:Data Engineer 做 ETL 并行,DS 做 KMeans 并行实验。
B5. 工程取舍三连
-
做什么:理解 ACID / BASE / CAP。
-
为什么:决定数据一致性与可用性策略。
-
怎么做:
-
银行 → ACID
-
社交 → BASE
-
网络分区 → CAP 折中
-
-
Tips:大规模分布式任务通常可容忍最终一致。
B6. 数据并行(Data Parallelism)
-
做什么:掌握数据切分 + 并行计算模式。
-
为什么:核心并行思想。
-
怎么做:
-
切分数据 → map(每核处理)
-
汇总结果 → reduce(聚合)
-
-
例子(词频统计):
map(doc_id, text):
for word in tokenize(text):
emit(word, 1)
reduce(word, counts_list):
emit(word, sum(counts_list))
B7. 我能做的项目并行点
-
做什么:明确实验并行点。
-
怎么做:
-
数据清理 / Tokenize
-
特征工程
-
模型参数搜索
-
KMeans n_init
-
Assignment 分片
-
批量预测
-
超参数验证 / CV
-
-
例子:KMeans n_init=10 → 并行跑 10 个随机初始化。
B8. 核心心得
-
做什么:总结并行注意事项。
-
为什么:避免踩坑。
-
怎么做:
-
粒度合适
-
IO vs CPU 辨别
-
并行度不要超过硬件极限
-
固定随机种子确保复现
-
Part C |练习与答案
练习 1|按结构分类数据
-
题目:服务器日志、CT、银行表、论坛帖子
-
答案:
-
服务器日志 → 日志/流
-
CT → 半结构化 / 二进制 + 元数据
- 银行表 → 结构化
-
论坛帖子 → 非结构化 + 半结构化元数据
-
练习 2|商品共现 MapReduce
-
方法一(pairs):
map(order_id, items_list):
for each pair (i, j):
emit((i,j), 1)
reduce(pair(i,j), counts):
emit((i,j), sum(counts))
-
方法二(stripes):
map(order_id, items_list):
for item in items_list:
stripe = {}
for other in items_list:
if other != item:
stripe[other] += 1
emit(item, stripe)
reduce(item, stripes_list):
merged = {}
for s in stripes_list:
for k,v in s.items():
merged[k] += v
emit(item, merged)
练习 3|KMeans 并行化方案
-
一句话:并行 n_init、或样本 assignment 分片、或 mini-batch 批处理。
Part D |项目行动路线图
-
阶段 0(准备):Python、joblib、scikit-learn 环境 + 数据下载
-
阶段 1(基线):串行实现 KMeans + 时间记录
-
阶段 2(并行化 1):并行 n_init + 加速比
-
阶段 3(并行化 2):Assignment 分片 + 聚合
-
阶段 4(性能调优):调整 n_jobs、线程/进程模式
-
阶段 5(结果与稳定性):复现测试、随机种子
-
阶段 6(报告 + 演示):PPT + demo
Part E |口袋速记
-
并行 vs 分布式:多核单机 vs 多台机器
-
Joblib 示例:
Parallel(n_jobs=4)(delayed(f)(x) for x in xs)
-
Multiprocessing:
Pool.map(func, iterable)
-
KMeans 并行点:n_init, assignment, mini-batch
-
并行陷阱:IO-bound 不一定线性加速,过高并行度会慢
Part F |待解决问题
-
IO 瓶颈优化? → 缓存、批量读、减少 shuffle
-
云上 vs 本地成本?
-
高维稀疏 KMeans 优化?
-
joblib 并行迁移到 Spark / Dask?
-
并发调试与日志策略?
附录:常用代码片段
-
Joblib 并行:
from joblib import Parallel, delayed
import time, math
def cal_sqrt(i):
time.sleep(0.5)
return math.sqrt(i**2)
num = 10
res = Parallel(n_jobs=4)(delayed(cal_sqrt)(i) for i in range(num))
print(res)
-
Multiprocessing Pool:
from multiprocessing import Pool
import math
def cal_sqrt(i):
return math.sqrt(i**2)
if __name__ == '__main__':
with Pool(4) as p:
results = p.map(cal_sqrt, range(10))
print(results)
-
KMeans 并行 n_init:
from sklearn.cluster import KMeans
from joblib import Parallel, delayed
import numpy as np
X = np.random.randn(10000, 10)
def run_kmeans(seed):
km = KMeans(n_clusters=5, n_init=1, random_state=seed)
km.fit(X)
return km.inertia_
results = Parallel(n_jobs=4)(delayed(run_kmeans)(s) for s in range(8))
print(results)
更多推荐


所有评论(0)