简介:大数据分析课程综合实验包,涵盖MapReduce词频统计(wordCount)、PageRank链接分析、关联规则挖掘(Apriori)、k-means聚类和推荐系统五个子实验,适合高校大数据、数据科学方向学生或基础算法入门者参考。整套实验包含完整的任务书说明文档、Python实现源码、数据集及运行结果文件,其中wordCount实验以包含百万级单词的9个源文件模拟分布式节点,要求实现9个map节点与3个reduce节点的多线程词频统计,并设计了combine与shuffle环节的思考题。资源压缩包共56个文件,核心内容为py算法脚本、csv/txt数据文件与docx实验说明,另含README与归一化数据等预处理文档,总大小约115MB,按lab1至lab5分目录组织,检索和复现都比较方便。包内源码经过实训流程验证,关键输出均已保留供对照,可帮助学习者快速理清五个经典大数据算法的核心逻辑。目前已有501人学习/下载,可作为实验报告撰写和实践操作的参考蓝本。 如果你正在做大数据分析相关的课程实验,大概率会碰到这套经典组合:WordCount、PageRank、关系挖掘、K-means、推荐系统算法。五个实验看起来是五个独立任务,实际上是一条精心设计的成长路径,从最简单的分布式词频统计,一路走到推荐系统这样贴近真实业务的场景。

这篇博文把我自己把这五个实验从零跑通、反复调参、最后整理成报告的全过程写下来,重点讲讲每一步该关注什么、代码里哪些细节容易翻车、以及怎么验证你的实验结果不是“恰好碰对了”。无论你是正在赶实验报告的学生,还是想系统入门大数据算法链路的学习者,这份复盘应该都能帮你省掉不少试错时间。

1. 实验选型里的递进逻辑:为什么偏偏是这五个

1.1 一条从“会调API”到“会做算法取舍”的爬坡路线

很多同学会把五个实验当成五个独立任务逐个做,做完就忘。我建议反过来:先把这五个实验当成一个整体去看,你才会知道老师设计这份实验清单时到底想让你掌握什么。

  • WordCount :入门标配。用最简单的词频统计让你理解分布式计算框架的Map/Reduce模型,搞清楚数据是怎么被切分、打散、合并的。
  • PageRank :从“一次计算”升级到“迭代计算”。同一个数据集要反复算很多轮,每一轮的输出是下一轮的输入,这就涉及收敛判断、悬挂节点处理这些WordCount里完全没有的问题。
  • 关系挖掘 :从“统计频率”进入“发现规则”。不再是算一个数,而是从大量事务数据里挖出“买了A的人还爱买B”这种隐含关联。
  • K-means :无监督学习的代表。没有标签,全靠数据自身的距离结构把样本分成K堆,是后面所有聚类问题的基础。
  • 推荐系统 :把前面学的所有能力缝合起来。要处理用户-物品矩阵、算相似度、做预测、还要设计评估指标,最接近真实互联网场景。

这五步分别对应了:计算框架入门 → 迭代算法设计 → 规则发现 → 无监督聚类 → 完整业务链路。难度是逐步抬升的,而且每一步都在复用前面的能力。

1.2 实验之间能复用的技术资产

我不建议你每个实验都从零开始写,很多代码是可以反复用的。

WordCount里学到的RDD操作(map、flatMap、reduceByKey)在PageRank和K-means里会反复出现;PageRank里的迭代更新和收敛判断,放进K-means里同样成立;关系挖掘里对“支持度-置信度”的评估思路,和推荐系统里对准确率、召回率的评估思路也是同构的。

把这些共同的思维模式提炼出来,你会发现五个实验真正的核心不是某个算法公式,而是三种通用能力: 把数据切成键值对的能力、迭代更新的能力、设计评估指标的能力 。

2. 实验环境与数据准备:最容易被低估的前置环节

2.1 Hadoop还是Spark:别只按课程要求选

不少学校还在用Hadoop MapReduce做这套实验,但如果你有自主选择权,我更推荐Spark。原因不是Hadoop过时了,而是作为学习工具,Spark的调试反馈速度要快太多。MapReduce每个Job都要落盘,一个简单的迭代可能几十秒甚至几分钟就没了,而Spark默认基于内存,小数据集秒级出结果。

我对这两个方案的使用感受是这样的:

对比维度 Hadoop MapReduce Spark
API上手难度 Java代码量大,类接口多 Python/Java/Scala都行,PySpark最友好
迭代计算效率 每轮都写磁盘,慢 内存计算,快很多
调试成本 日志多、堆栈长 本地模式可直接查UI,定位快
课程匹配度 如果课程指定了要用,没法换 很多课程也接受Spark,作为MapReduce的进阶实现

如果你最终决定用Spark,建议版本组合用Spark 3.x + Python 3.9/3.10 + JDK8/11,这套组合目前最稳。安装时注意JDK版本别太新,JDK17在某些老版本Spark上会有模块访问报错。

提示:先用本地模式把逻辑跑通,再考虑伪分布式或集群。绝大多数实验的测试数据量小到根本不需要集群,本地模式能让你把精力放在算法上,而不是集群运维上。

2.2 五个实验的数据集怎么准备

数据集的挑选直接影响实验工作量。我的建议是:

  • WordCount :随便找一本英文电子书转成txt,或者直接用Project Gutenberg上的公开文本。注意必须是纯文本格式,编码统一成UTF-8。
  • PageRank :别一开始就用真实爬虫数据,先用一个6到8个节点的小图验证算法正确性,再换大图测性能。小图结构类似 fromNode toNode ,每行一条有向边。
  • 关系挖掘 :用经典的超市购物篮数据,网上可以找到Groceries数据集,一行为一个事务,每个商品用单词或编号表示。
  • K-means :自己用随机数生成二维高斯分布的若干簇,这种数据画出来后肉眼看得很清楚,验证聚类效果极其直观。也可以从UCI等公开数据集下载,但Iris这类数据更适合做分类实验。
  • 推荐系统 :最常用的是MovieLens数据集,一个小版本约10万条评分记录,包含userId、movieId、rating、timestamp四个字段。数据量适中,特征清晰。

数据集准备好了,最好统一放到一个 data/ 目录下,并在代码里用相对路径或配置文件引用,不要在代码里写死绝对路径,否则换台机器跑就直接崩。

3. WordCount 与 PageRank:从一次计算到迭代计算的思维跳转

3.1 WordCount实验要做对的几个设计点

WordCount代码本身不难,核心逻辑基本就是:把每行文本按分隔符切词,输出 (word, 1) ,然后按单词累加。但如果你想拿高分,或者想真正理解分布式框架,有几个点值得刻意设计。

第一个是 Combiner 。Map端输出后如果直接全部传给Reduce,会产生大量网络传输。Combiner本质是在Map端先做一次局部合并,比如一个节点上出现100次“the”,先合成 (the, 100) 再传出去。要注意Combiner必须满足交换律和结合律,词频统计天然满足,所以适合做。如果你用Spark, reduceByKey 本身自带本地合并效果,比 groupByKey 聪明得多。

第二个是 数据倾斜 。真实文本里“the”、“a”这类高频词可能出现几十万次,单个Reduce任务会拖慢整体速度。实验里可以重新设计Partitioner,让高频词分散到不同Reduce任务,或者加一个随机key前缀再分两步聚合。这是面试常考点,实验报告里写出来会很加分。

第三个是 自定义输出格式和Counter 。你可以用Counter统计总词数、过滤掉的空行数,这能让实验结果更可信,而不是只输出一个WordCount结果文件就完事。

一段最简PySpark版本的WordCount大概是长这个样子的:

def tokenize(line):
    for word in re.split(r'[^\w]+', line.strip().lower()):
        if word:
            yield (word, 1)

word_counts = (
    sc.textFile("data/input.txt")
      .flatMap(tokenize)
      .reduceByKey(lambda a, b: a + b)
      .sortBy(lambda x: x[1], ascending=False)
)
word_counts.saveAsTextFile("output/wordcount")

3.2 PageRank迭代中的三个工程细节

PageRank公式本身不算复杂,核心迭代是 PR(A) = (1-d) + d * sum(PR(T)/C(T)) ,其中每个页面把它的PR值平均分给所有出链页面。真正让新手翻车的是三个工程细节。

第一个是 阻尼因子d的取值 。通常取0.85,它的含义是模拟用户浏览网页时有15%的概率会随机跳到任意一个页面,这样能避免排名在环路上死循环。如果d=1,整个迭代可能会发散或者陷入特定环路。

第二个是 悬挂节点的处理 。如果一个网页没有任何出链,它的PR值会“吞掉”整个网络的质量。处理方法是把悬挂节点传出的PR平均分配给图中所有节点,或者干脆把这些节点单独揪出来处理。

第三个是 收敛判断 。不要写死迭代次数,而是用两轮迭代之间所有节点PR值变化量的绝对值之和作为判断依据,当它小于某个阈值(例如1e-6)时停止。实际调参时你会发现,稀疏大图和密集小图的收敛速度差异非常大。

PageRank在Spark里的迭代思路是维护两个RDD:一个是 links (节点到出链列表的映射),另一个是 ranks (节点到当前PR值的映射),每轮通过join把两者合并,再计算贡献值。下面是一段核心结构:

for i in range(max_iter):
    contributions = links.join(ranks).flatMap(
        lambda url_links_rank: [
            (url, rank / len(links_list))
            for url in links_list
        ]
    )
    new_ranks = contributions.reduceByKey(add).mapValues(
        lambda score: 0.15 + 0.85 * score
    )
    delta = ranks.join(new_ranks).map(
        lambda url_r1_r2: abs(r1 - r2)
    ).sum()
    ranks = new_ranks
    if delta < tol:
        break

3.3 这两个实验我踩过的真实问题

第一坑是 迭代RDD的血统过长 。用Spark做PageRank如果不做Checkpoint,几十次迭代之后,RDD的血统(Lineage)会非常长,一旦某个节点需要重算,会从最初源头一直重跑,导致栈溢出或性能骤降。解决办法是每隔几轮调用一次 rdd.checkpoint() ,把中间结果落到可靠存储里。

第二坑是 小图数据格式里的空格和换行符 。很多图数据有前导空格或空行,解析时如果直接 split(" ") ,会发现多出空字符串节点。我建议解析边时统一用 line.strip().split() ,并且过滤掉 ==2 的数据行。

第三坑是 算法做对了但结果排序和预期不一致 。PageRank算完之后,如果两个页面的PR值在小数点后6位完全一样,排序会不稳定。实验报告里最好对最终PR值做一次 sortByKey(False) ,并把精度保留到合理位数再输出,这样结果更稳定也更美观。

4. 关系挖掘与 K-means:数据挖掘里的两个经典陷阱区

4.1 关系挖掘搞清楚支持度、置信度、提升度就够了

关系挖掘实验最常见的实现是Apriori算法,它的核心逻辑是:先找到所有满足最小支持度的频繁项集,再用频繁项集生成置信度满足阈值的关联规则。这个过程听起来简单,但Apriori有个非常关键的性质—— 如果一个项集不频繁,它的所有超集也一定不频繁 。这个性质是剪枝的基础,也是Apriori效率的来源。

以超市购物篮为例,数据长这样:

bread, milk, eggs
bread, milk, diapers
milk, eggs

最小支持度设为2/3的话,单项集 bread 出现2次、 milk 出现3次、 eggs 出现2次、 diapers 出现1次。这样 diapers 直接剪枝掉,不用再考虑它和任何商品的组合。这个步骤在MapReduce或Spark中实现时,关键在于频繁项集的逐层统计——先算单个商品频次,过滤后生成候选2项集,再统计、再过滤,直到候选集为空。

Apriori在大数据集上性能很差,因为它需要反复扫描数据集并生成海量候选集。如果你用Spark,可以直接用MLlib里的 FPGrowth 做对比实验,它会用FP树压缩数据集,比Apriori快好几个量级。我个人建议实验里“手写Apriori理解原理,再调FPGrowth看差距”,这会让实验报告非常有层次。

评估规则时有三个标准:

  • 支持度 :规则在所有事务中出现的占比,衡量规则是否有足够的数据支撑。
  • 置信度 :在包含前件的所有事务中,同时包含后件的比例,衡量规则的可靠性。
  • 提升度 :规则的实际置信度和后件独立出现概率的比值。提升度大于1才说明前件对后件有正向影响。

很多新手只看置信度,结果挖出一堆“买牛奶就会买牛奶”这种废话规则。实验报告里一定要把提升度也列出来,这是展示你理解到位的关键。

4.2 K-means实验的三大调参点

K-means是五个实验里代码负担最小的一个,但也是参数影响最大的一个。我见过太多同学把K设置成3,随机初始化一次就出结果,然后发现每次跑出来的簇都不一样。这背后的原因是初始化点对算法结果影响极大。

第一个调参点是 K值的选择 。课程里最常用的是肘部法则:分别跑K=2、3、4、5,计算每个K下所有样本到所属中心点距离的平方和(SSE),然后画出曲线,找到“肘部”位置。SSE下降变缓的那个K就是推荐值。当然也有更严谨的轮廓系数,但对本科实验来说肘部法则已经够了。

第二个调参点是 初始中心点的选取 。随机初始化容易收敛到局部最优,最好使用K-means++策略:第一个中心点随机选,后续每个中心点都尽量选离已有中心点远的样本点。Spark MLlib里的 KMeans 默认已经实现了K-means++,但如果你手写K-means,务必把初始化这一步写清楚。

第三个调参点是 标准化 。如果特征维度之间的量纲差异很大(比如一列是0~100,另一列是0~1),距离计算会被大数值特征主导,聚类结果基本失效。所以K-means之前必须做Min-Max标准化或Z-score标准化。这一步经常被忽略,但恰恰是实际业务中决定聚类效果的关键。

K-means在Spark里的实现核心是迭代地用 reduceByKey 统计每个簇的样本总和与样本数量,再计算新的中心点。下面是一个简化实现片段:

centroids = initial_centroids
for i in range(max_iter):
    points_with_cluster = points.map(lambda p: (nearest_centroid(p, centroids), p))
    partials = points_with_cluster.mapValues(lambda p: (p, 1.0))
    sums = partials.reduceByKey(lambda a, b: (a[0] + b[0], a[1] + b[1]))
    new_centroids = sums.mapValues(lambda s: s[0] / s[1]).collectAsMap()
    centroids = [new_centroids[k] for k in sorted(new_centroids)]

4.3 这两类算法的评估边界

关系挖掘和K-means都属于“结果需要人工解释”的算法,所以特别容易得到看似合理但实际没有意义的结论。关联规则里高置信度不一定代表有因果,可能是后件本身出现概率就很高;K-means对凸形簇效果好,对长条形或环形分布的数据效果会非常差。

写实验报告时建议主动讨论这些局限性,这比堆一堆数字更能体现你的分析能力。

5. 推荐系统实验:从相似度公式到离线评估的全链路

5.1 该手写UserCF还是ItemCF

推荐系统实验最经典的实现是协同过滤,分为基于用户的(UserCF)和基于物品的(ItemCF)两种。课程实验里一般建议用户自己实现一种,再用另一种做对比。

UserCF的思路是:找到和当前用户口味最相似的一群用户,把那些用户喜欢但当前用户没见过的物品推荐过来。ItemCF的思路则是:找到和目标物品最相似的物品,只要用户喜欢过某个物品,就推荐和它类似的物品。

实际互联网里ItemCF更常用,因为物品相似度可以离线算好,在线推荐时只需要查表,响应速度更快。而用户相似度随着用户行为不断变化,实时计算代价高。但UserCF在新闻推荐、社区推荐这种物品生命周期短的场景里反而更合适。

手写协同过滤时建议至少包含以下模块:

  • 数据加载模块:把MovieLens数据读成用户-物品评分字典。
  • 相似度计算模块:计算用户间或物品间的相似度矩阵。
  • 预测模块:根据相似用户或相似物品做加权评分预测。
  • 评估模块:把数据集按时间或随机划分为训练集和测试集,计算误差指标。

5.2 相似度计算的细节和评分预测

相似度计算是协同过滤的胜负手。常见的有三种:Jaccard相似度只看交集大小,适合布尔数据;余弦相似度看向量夹角,适合评分数据;皮尔逊相关系数对用户的评分习惯做了去均值处理,能消除“有人总打高分,有人总打低分”的影响。

以UserCF为例,预测用户u对物品i的评分时,通常先找出与u最相似的K个用户,然后用这些用户对i的评分做加权平均,权重就是相似度值。如果某些用户没评过i,就直接跳过。最后要做归一化,防止评分被放大或缩小。公式细节不难,但很容易写错,我建议先用一个3用户×3物品的小矩阵手算一遍期望结果,再跑代码验证。

5.3 离线评估不能只看“跑通了”

很多同学做完推荐系统实验就只输出一堆Top-N推荐列表,然后截图结束。但推荐系统是一个链路,真正的价值在于评估算法好坏。

离线评估时第一步是划分数据集。我强烈建议按时间切分,用用户前80%的行为做训练,后20%做测试,这样才能模拟真实场景。如果随机划分,训练集和测试集高度重合,指标会虚高。

第二步是选指标。预测评分用RMSE和MAE,衡量预测值和真实值的误差;Top-N推荐用准确率、召回率、F1,衡量推荐列表里有多少是用户真正喜欢的。

第三步是考虑覆盖率和多样性。如果算法只推荐热门item,准确率不会差,但用户不会有任何惊喜感。覆盖率指推荐的物品占总物品的比例,多样性指每个用户的推荐列表之间的差异程度。这两个指标在实验报告里体现出来,会立刻拉开和其他报告的差距。

如果你用的是Spark MLlib,直接调 ALS 做矩阵分解也可以,和手写UserCF/ItemCF形成对比组。ALS能有效缓解稀疏问题,但需要调 rank (潜在因子数)、 regParam (正则化系数)和 alpha (置信度参数),这个调参过程本身也是很好的实验内容。

6. 让实验结果可信的验证套路与调试清单

6.1 五个实验分别怎么验证结果对不对

我把每个实验的验证方法整理成了一套固定的动作,做完一个实验就对着清单检查一遍:

实验 快速验证方法
WordCount 先用Python的collections.Counter统计同一份数据的结果,与MapReduce/Spark输出做diff。
PageRank 构造一个5节点以内的小图,手算两到三迭代的期望值,与程序前几轮输出对比。
关系挖掘 用经典Apriori测试数据集(如Groceries的子集)跑一遍,和已知频繁项集结果比对。
K-means 用随机生成的二维高斯簇数据,画散点图检查聚类结果是否和肉眼判断一致。
推荐系统 用一个人工构造的极小用户-物品评分矩阵,手推预测评分,再和程序输出对比。

这套验证逻辑的核心是 先用小数据验证正确性,再用大数据验证性能和扩展性 。直接拿真实大数据调试算法,出了问题根本分不清是逻辑错还是数据怪。

6.2 调试阶段的通用排查顺序

如果你发现实验结果不对,请按下面的顺序排查,不要一上来就觉得算法写错了:

先检查输入数据的解析是否干净。我用Spark时最常踩的坑是文本编码不一致、行尾有 \r 、CSV字段带引号导致解析多出引号字符。这一步排查通常能解决50%以上的“结果不对”。

再检查键值对阶段是否有数据丢失。比如PageRank里有些节点出现在目标列但没出现在源列,这些节点初始PR值就是0,会导致分布异常。推荐系统的评分数据也可能有重复记录,要提前去重或用 rating 字段做聚合。

最后才检查算法逻辑。调试时尽量加日志,打印关键中间结果。Spark可以用 take(5) 、 count() 、 collect() 等方法快速检查每个RDD的内容。千万别每次都跑到最后才看输出,那样排查效率太低。

6.3 让实验报告更有分量的几个建议

实验报告不只是贴代码和截图,我建议每个实验都留下这些痕迹:不同参数下结果的对比表、数据分析时的直观观察、对算法局限性的讨论。

以K-means为例,不要只说“K=3时效果最好”,而是放一张肘部法则曲线图,一段SSE的变化数据,解释为什么K=3是合适的。推荐系统实验里,把ALS不同rank下的RMSE列成一张表,再把最佳参数下Top-N推荐的准确率和覆盖率一起展示。这些内容比起一个光秃秃的运行结果截图,有说服力得多。

最后一个小建议:所有实验代码统一用Git管理,一个实验一个分支。跑实验的时候把关键参数、数据集规模、运行时间都记录在README里。我自己的习惯是写完代码先记录baseline,再逐项优化,这样实验报告根本不用临时编数据——你每一步的调参记录就是最扎实的内容。

本文还有配套的精品资源,点击获取 menu-r.4af5f7ec.gif

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐