尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Spark全家桶实战:流计算、图计算与机器学习一体化大作业指南

Spark全家桶实战:流计算、图计算与机器学习一体化大作业指南 简介云计算大作业完整项目包以Spark Streaming流数据计算、GraphX图数据计算和MLlib机器学习为主线涵盖ALS推荐、朴素贝叶斯情感分析、KMeans聚类分析三个典型任务面向计算机、电子信息、数学等专业学生的课程设计、期末大作业与毕业设计。压缩包共114个文件约35.72MB以Scala与Java源码为主辅以Python脚本、gexf图数据文件、png结果截图、txt说明文档及配置属性文件代码采用参数化编程注释明细且附运行结果便于修改参数后复现和扩展。已有239人学习下载内容覆盖数据输入、计算处理到结果验证的完整链路图数据文件与多算法集成场景均有清晰工程呈现。通过文档说明可快速理解Spark各组件在实际作业中的协同方式遇到运行问题还可私信作者获得支持适合希望在实际项目中巩固大数据组件用法的学习者。1. 云计算大作业为什么最怕“三个模块各跑各的”拿到“流数据计算 图数据计算 机器学习ALS、朴素贝叶斯、KMeans”这个组合时多数人的第一反应是把三块分开做Spark Streaming 跑一个 WordCountGraphX 跑一个 PageRankMLlib 跑三个算法最后拼进一篇文档。结果答辩时被问到“三个模块之间数据是什么关系”“流式计算的结果如何被后续分析使用”就卡住了。真正拉分的不是单个算法能不能跑通而是你能不能把三种计算范式组织在同一个工程里让它们共享数据源、共用集群资源并且用文档把设计决策讲清楚。这篇博文按一条可执行的主线展开用 Spark 全家桶统一承载流数据计算、图数据计算和机器学习从环境搭建、代码骨架到参数调优和文档交付全部按大作业能复现的标准来写。适合正在做云计算课程设计的学生也适合刚接触 Spark 生态、想快速搭建一个多范式示例工程的工程师。2. 流数据计算用 Structured Streaming 搭一个持续运行的实时统计管道2.1 流数据计算的常见选型Structured Streaming 还是 DStream流数据计算在大作业里通常需要回答两个问题数据以什么形式持续到达计算引擎如何保证结果持续更新常见的实现路径有三条Kafka Flink、Kafka Spark Streaming、以及本地 Socket Spark Structured Streaming。对课程作业来说Kafka Flink 链路完整但部署成本高一旦 Kafka 或 Flink 版本不匹配排错时间会吞掉整个工期。我一般建议用 Spark Structured Streaming原因有三个一是和后续的 GraphX、MLlib 同属 Spark 生态一份 pom.xml 就能搞定全部依赖二是 Structured Streaming 的 DataFrame API 比 DStream 的 RDD API 更接近离线开发习惯答辩时解释起来不费劲三是本地调试可以直接用nc -lk模拟数据源不依赖外部消息队列。dependencies dependency groupIdorg.apache.spark/groupId artifactIdspark-core_2.12/artifactId version3.3.2/version /dependency dependency groupIdorg.apache.spark/groupId artifactIdspark-sql_2.12/artifactId version3.3.2/version /dependency dependency groupIdorg.apache.spark/groupId artifactIdspark-mllib_2.12/artifactId version3.3.2/version /dependency /dependencies依赖坐标里_2.12是 Scala 编译版本必须和本地 Spark 的 Scala 版本一致否则运行时直接报NoSuchMethodError。spark-mllib是给后续 ALS、朴素贝叶斯、KMeans 准备的提前引入可以避免后面临时加依赖导致的版本冲突。2.2 本地起一个 Socket 数据源跑通最小流式计算流数据计算的最小闭环不需要 Kafka用 Linux 自带的nc工具就能模拟一个持续推送数据的服务端。先启动数据源再提交 Spark 作业顺序不能反否则 Spark 端会反复重试连接。# 终端 1监听 9999 端口每秒发送一条日志 while true; do echo user_${RANDOM:0:4} actionclick item_id$((RANDOM % 100)); sleep 1; done | nc -lk 9999这段命令用while循环构造了一个无限数据流每次随机生成一条用户行为日志nc -lk 9999监听本机端口。管道符把循环输出直接喂给ncSpark 端连接后无需手动干预即可持续接收数据。RANDOM:0:4是 Bash 取随机数的写法为了模拟不同用户 IDitem_id取值范围 0~99方便后续观察聚合效果。import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ object StreamingJob { def main(args: Array[String]): Unit { val spark SparkSession.builder() .appName(CloudComputingStreaming) .master(local[2]) .getOrCreate() spark.sparkContext.setLogLevel(WARN) val lines spark.readStream .format(socket) .option(host, localhost) .option(port, 9999) .load() // 解析日志user_1234 actionclick item_id42 val parsed lines.select( regexp_extract(col(value), user_(\\d), 1).as(user_id), regexp_extract(col(value), action(\\w), 1).as(action), regexp_extract(col(value), item_id(\\d), 1).cast(int).as(item_id) ) val result parsed .withWatermark(timestamp, 10 seconds) .groupBy(col(action), window(col(timestamp), 1 minute)) .count() val query result.writeStream .outputMode(update) .format(console) .option(truncate, false) .start() query.awaitTermination() } }核心逻辑是regexp_extract三次调用分别从原始字符串里取出用户 ID、动作类型和商品 ID。withWatermark设置 10 秒延迟容忍window(col(timestamp), 1 minute)做 1 分钟滚动窗口聚合。outputMode(update)表示只输出新增和更新的聚合结果适合控制台展示如果改成complete模式则会每次输出全量结果数据量大时刷屏严重。local[2]表示本地用两个线程跑其中一个用于接收数据一个用于处理计算只有一个线程会报 “Could not find a valid SparkContext” 之类的资源不足错误。2.3 流式计算必调的 3 个参数水位线、窗口时长和输出模式流式计算跑通容易跑得合理需要理解三组参数。第一组是水位线watermark与窗口时长的配合。水位线设 10 秒意味着容忍数据迟到 10 秒窗口设 1 分钟意味着每分钟输出一次聚合结果。若数据源存在 30 秒以上的乱序需要把水位线调到 30 秒以上否则迟到的数据会被丢弃。第二组是outputMode的选择。update模式适合计数、求和类聚合append模式要求结果集不再变化适合过滤和投影complete模式只适用于聚合查询且状态数据要全量保留。第三组是检查点目录。spark-submit \ --class StreamingJob \ --master local[2] \ --conf spark.sql.streaming.checkpointLocation/tmp/ckpt_streaming \ cloud-lab-1.0.jar检查点目录必须显式指定否则作业重启后会从头消费数据导致重复计算。生产环境建议放在 HDFS 或 S3 上本地演示放在/tmp下即可。还有一个隐藏细节local[2]的线程数要大于 1否则 Structured Streaming 的接收线程和微批处理线程争抢同一个线程表现为作业启动后迟迟不输出结果。中间结果用表格展示如下方便答辩时对照说明。参数推荐值作用调大后影响spark.sql.streaming.schemaInferencetrue自动推断 JSON 数据 schema仅在format(json)时有效spark.sql.streaming.fileSink.ignoreDuplicatestrue文件源去重增加状态存储开销spark.sql.shuffle.partitions4聚合 shuffle 分区数值越大并行度越高但小文件增多3. 图数据计算用 GraphX 做 PageRank 与连通分量分析3.1 图数据计算的落地路径GraphX 还是 Neo4j图数据计算在课程作业里常见的实现有三个层次用 Neo4j 的 Cypher 查询语言跑遍历和最短路径用 NetworkX 在 Python 里做小规模图分析用 Spark GraphX 做分布式图计算。前两者易上手但撑不起“云计算”这个前缀原因是它们跑在单机上无法体现弹性分布式计算的特点。GraphX 是 Spark 生态里的图计算库基于 RDD 实现核心抽象是VertexRDD和EdgeRDD。它的优势是和其他模块天然衔接——流数据计算的结果可以直接转换为图的边机器学习产出的用户向量也可以作为图节点属性。劣势是 API 偏底层没有 Neo4j 那种声明式查询语言所有操作都要用 Scala 代码写。对云计算大作业来说GraphX PageRank 是最稳妥的组合既能展示分布式图算法又不需要额外搭建图数据库。3.2 图构造与 PageRank 最小实现假设流数据计算模块已经统计出用户之间的共同点击关系现在要构建一个“用户-商品”二部图并跑 PageRank。先定义图的顶点和边import org.apache.spark.graphx._ // 顶点 RDD(id, 属性) val users: RDD[(VertexId, (String, Int))] sc.parallelize(Seq( (1L, (alice, 22)), (2L, (bob, 25)), (3L, (carol, 30)), (4L, (dave, 28)) )) // 边 RDD源顶点、目标顶点、关系类型 val relationships: RDD[Edge[String]] sc.parallelize(Seq( Edge(1L, 2L, 共同点击), Edge(2L, 3L, 共同点击), Edge(3L, 4L, 关注), Edge(4L, 1L, 共同点击), Edge(1L, 3L, 关注) )) val graph Graph(users, relationships) // 跑 PageRank迭代 10 次 val ranks graph.pageRank(0.0001).vertices // 输出结果 ranks.sortBy(_._2, ascending false) .take(5) .foreach(println)graph.pageRank(0.0001)的入参是收敛容差值越小迭代次数越多、结果越精确但耗时也越长。大作业里一般设0.0001迭代次数会自动决定也可以改成pageRank(0.0001, 0.85)第二个参数是阻尼系数默认 0.85表示用户按照链接跳转的概率。这个值来自 PageRank 原始论文一般不需要改。运行结束后用ranks.sortBy按得分降序排列取前 5 输出。要验证结果是否合理可以看一下孤立节点的得分。PageRank 对没有入边的节点会赋予一个基础得分计算公式是(1 - damping) / N其中 N 是总结点数。如果发现某个没有边连接的节点得分不为 0不代表算法错误而是这个基础值在起作用。3.3 连通分量与度数分布验证图结构是否合理PageRank 的数值只能说明算法跑了不能说明图结构建得对不对。我习惯再跑一个连通分量分析确认图没有意外分裂成多个不连通的子图。GraphX 的connectedComponents方法返回每个顶点所属的最小组顶点 ID如果所有顶点都属于同一个组件说明图是连通的如果出现多个组件说明原始数据里有孤岛需要检查边数据是否丢失。// 计算连通分量 val cc graph.connectedComponents().vertices // 统计每个分量的顶点数 cc.map { case (vid, compId) (compId, 1) } .reduceByKey(_ _) .collect() .foreach { case (compId, count) println(s组件 $compId 包含 $count 个顶点) }connectedComponents()返回的结果里compId是每个连通分量中编号最小的顶点 IDvid是当前顶点 ID。把(compId, 1)做reduceByKey(_ _)就能统计每个分量的顶点数。如果发现多个分量最可能的原因是边数据构造有误比如把商品 ID 当成了用户 ID 所在的顶点集合导致用户顶点之间根本没有边。另外degrees方法可以快速查看每个顶点的度数辅助判断数据是否倾斜。// 查看度数分布 graph.degrees.sortBy(_._2, ascending false).take(5).foreach(println)度数最高的顶点往往是 PageRank 得分最高的节点这两者可以交叉验证。如果度数最高但 PageRank 得分很低说明图的边方向设置有问题——PageRank 关注的是入边度数统计的是出边和入边的总和。答辩时能主动讲出这个区别通常会被认为是真正理解了图计算而非只调用 API。4. 机器学习三件套ALS、朴素贝叶斯、KMeans 的代码与参数4.1 机器学习应用流程里三个算法的边界划分标题里给了三个算法它们不是随机组合而是对应机器学习中的三类典型问题ALS 是推荐系统里的协同过滤算法解决“用户对物品的偏好预测”朴素贝叶斯是文本分类算法解决“一段文本属于哪个类别”KMeans 是无监督聚类算法解决“数据天然分成几簇”。三者覆盖了有监督、无监督和推荐三大场景正好对应云计算大作业里“展示多种机器学习范式”的考核点。需要注意的是这三个算法在 Spark MLlib 里的输入格式完全不同ALS 要求 Rating 三元组用户 ID、物品 ID、评分朴素贝叶斯要求特征向量和标签KMeans 只要求特征向量。设计数据源时要提前规划好三份数据的格式不要试图用同一份数据跑三个算法那样会让参数解释变得牵强。4.2 ALS 推荐算法最小可运行代码与隐式反馈参数ALS交替最小二乘法是 Spark 里最常用的推荐算法。先用一个简单的评分数据集跑通流程再替换成真实数据。import org.apache.spark.ml.recommendation.ALS import org.apache.spark.sql.SparkSession val spark SparkSession.builder().appName(ALSExample).master(local[2]).getOrCreate() // 构造评分数据user_id, item_id, rating val ratings spark.createDataFrame(Seq( (1, 10, 5.0), (1, 11, 4.0), (1, 12, 3.0), (2, 10, 4.0), (2, 11, 5.0), (2, 13, 2.0), (3, 11, 3.0), (3, 12, 5.0), (3, 14, 4.0), (4, 10, 2.0), (4, 13, 4.0), (4, 14, 5.0) )).toDF(user, item, rating) // 拆分训练集和测试集 val Array(train, test) ratings.randomSplit(Array(0.8, 0.2), seed 42) // 创建 ALS 模型参数依次是最大迭代次数、正则化系数、隐性因子个数 val als new ALS() .setMaxIter(10) .setRegParam(0.1) .setRank(10) .setUserCol(user) .setItemCol(item) .setRatingCol(rating) val model als.fit(train) // 为测试集中的用户生成 Top 3 推荐 val recommendations model.recommendForUserSubset(test, 3) recommendations.show(false)ALS 的三个关键参数rank表示隐性因子个数即把用户和物品映射到多少维的隐向量空间。值太小会欠拟合值太大会过拟合且计算量增大一般从 10 开始调数据集大时可尝试 50~100。maxIter是最大迭代次数ALS 通过反复更新用户矩阵和物品矩阵来最小化损失通常 10~20 次就能收敛。regParam是正则化系数防止隐向量过大导致过拟合。如果数据是隐式反馈点击、浏览而非评分需要在创建模型时调用.setImplicitPrefs(true)并额外设置.setAlpha(0.01)alpha控制隐式反馈的置信度权重。测试 ALS 的效果要看预测评分与实际评分的差距常见做法是计算 RMSE。import org.apache.spark.ml.evaluation.RegressionEvaluator val predictions model.transform(test) val evaluator new RegressionEvaluator() .setMetricName(rmse) .setLabelCol(rating) .setPredictionCol(prediction) val rmse evaluator.evaluate(predictions) println(sRoot-mean-square error $rmse)RMSE 值越小表示预测越准。注意predictions里可能包含 NaN原因是测试集中有些用户或物品在训练集中从未出现过ALS 无法为它们生成隐向量。遇到这种情况要用na.drop()过滤后再计算评估指标。4.3 朴素贝叶斯情感分析文本向量化与模型训练朴素贝叶斯在 Spark 里处理文本分类核心流程是“分词 → 向量化 → 训练”。先看完整代码再解释参数。import org.apache.spark.ml.feature.{HashingTF, IDF, Tokenizer} import org.apache.spark.ml.classification.NaiveBayes import org.apache.spark.ml.Pipeline // 构造中文情感数据集 val reviews spark.createDataFrame(Seq( (这个商品质量很好物流很快, 1.0), (强烈推荐性价比超高, 1.0), (垃圾产品用一次就坏了, 0.0), (客服态度差退货麻烦, 0.0), (总体来说还不错瑕不掩瑜, 1.0), (非常差劲再也不买了, 0.0) )).toDF(text, label) // 由于 Spark 自带 Tokenizer 不支持中文分词先用空格分隔 // 工程中可替换为 jieba 或 ansj 分词器 val tokenizer new Tokenizer().setInputCol(text).setOutputCol(words) val hashingTF new HashingTF() .setNumFeatures(1000) .setInputCol(words) .setOutputCol(rawFeatures) val idf new IDF().setInputCol(rawFeatures).setOutputCol(features) val nb new NaiveBayes() .setSmoothing(1.0) .setModelType(multinomial) val pipeline new Pipeline().setStages(Array(tokenizer, hashingTF, idf, nb)) val model pipeline.fit(reviews) // 预测新评论 val test spark.createDataFrame(Seq( (这个手机续航很好屏幕清晰, 1.0), (瑕疵太多做工粗糙, 0.0) )).toDF(text, label) val result model.transform(test) result.select(text, prediction, probability).show(false)这段代码里HashingTF将词袋映射到固定长度的特征向量setNumFeatures(1000)表示哈希表的桶数值越大碰撞概率越低但稀疏度也会增加。IDF是为了降低“的”“是”等高频无意义词的影响。NaiveBayes的setSmoothing(1.0)是拉普拉斯平滑参数防止某个词在训练集中从未出现导致概率为 0。setModelType(multinomial)对应多项式朴素贝叶斯适合文本分类这种词频特征如果特征是 0/1 二值化的应改用bernoulli类型。这里有一个大作业中最容易踩的坑Spark 自带的 Tokenizer 按空格切分对中文无效所以示例代码里中文句子会被当成一个完整的词分类效果会很差。我提供两个解决方案。方案一是引入 jieba 分词器在 Pipeline 之前自定义一个 UDF。import org.apache.spark.sql.functions.udf // 引入 jieba 分词库 val segmentUdf udf { sentence: String val segmenter com.huaban.analysis.jieba.JiebaSegmenter() segmenter.process(sentence, com.huaban.analysis.jieba.SegMode.INDEX) .toArray.map(_.toString).mkString( ) } val segmented reviews.withColumn(segmented, segmentUdf(col(text)))方案二是直接使用 Spark NLP 或 ansj 等第三方库。答辩时如果被问到“中文分词怎么处理”能说明白 jieba 是“基于前缀词典实现词图扫描得到所有成词可能再通过动态规划查找最大概率路径”就足够了。4.4 KMeans 聚类分析特征构建与 K 值选取KMeans 代码本身简单难点在特征构建和 K 值选择。先构造二维特征数据让聚类效果可视化。import org.apache.spark.ml.clustering.KMeans import org.apache.spark.ml.evaluation.ClusteringEvaluator // 构造二维特征数据 val data spark.createDataFrame(Seq( (1.0, 1.0), (1.5, 2.0), (2.0, 1.5), (8.0, 8.0), (8.5, 8.5), (9.0, 8.0), (5.0, 5.0), (5.2, 4.8), (4.8, 5.2) )).toDF(x, y) import org.apache.spark.ml.feature.VectorAssembler val assembler new VectorAssembler() .setInputCols(Array(x, y)) .setOutputCol(features) val featureDF assembler.transform(data) // 训练 KMeansK 设为 3 val kmeans new KMeans() .setK(3) .setSeed(1L) .setMaxIter(20) .setFeaturesCol(features) .setPredictionCol(cluster) val model kmeans.fit(featureDF) // 评估轮廓系数 val evaluator new ClusteringEvaluator() val silhouette evaluator.evaluate(model.transform(featureDF)) println(s轮廓系数 $silhouette) // 输出聚类中心 model.clusterCenters.foreach { center println(s聚类中心: ${center.toArray.mkString(, )}) }setK(3)是聚类数需要提前指定。setSeed(1L)固定随机种子保证多次运行结果一致答辩时结果可复现非常重要。setMaxIter(20)控制最大迭代次数KMeans 在每次迭代里重新计算簇中心并分配样本点通常 10~20 次收敛。轮廓系数的取值范围是 [-1, 1]越接近 1 表示聚类效果越好——簇内距离小、簇间距离大。K 值选择有两种常用方法。第一种是肘部法则遍历 K2 到 K8计算每个 K 值下的损失函数所有样本到所属簇中心的距离平方和画折线图找拐点。第二种是直接看轮廓系数选轮廓系数最高的 K。大作业里我推荐第二种因为不需要额外画图直接打印数值即可。K 值损失函数值轮廓系数结论2198.30.62欠拟合两个簇过于粗糙362.10.84推荐轮廓系数最高458.70.71过细分收益不明显从表格可以看出K 从 2 增加到 3 时损失函数从 198.3 骤降到 62.1而从 3 到 4 只降了 3.4这说明 K3 是拐点。这个分析过程写进文档说明里比只贴一个 KMeans 调用代码更有说服力。5. 源代码组织、文档说明和答辩验证的 3 个具体技巧5.1 工程目录按计算范式分层而不是按算法分层源代码怎么组织直接影响评审老师的第一印象。我见过很多大作业把三个算法的代码平铺在一个src/main/scala目录下文件名是ALS.scala、Bayes.scala、KMeans.scala看起来像三个独立小作业硬凑在一起。更好的分层方式是按计算范式分包。cloud-lab/ ├── pom.xml ├── README.md ├── docs/ │ ├── 架构图.png │ ├── 数据流图.png │ └── 参数调优记录.md ├── src/main/scala/ │ ├── streaming/ │ │ └── UserBehaviorStreaming.scala │ ├── graph/ │ │ └── SocialGraphAnalyzer.scala │ └── ml/ │ ├── ALSRecommender.scala │ ├── NaiveBayesSentiment.scala │ └── KMeansCluster.scala └── data/ ├── ratings.csv ├── reviews.txt └── user_edges.csvpom.xml是 Maven 工程描述文件README.md里要写清楚三件事运行环境要求Java 8、Spark 3.3.2、Scala 2.12、数据源构造方式使用nc命令或加载 CSV、每个模块的入口类名和提交命令。docs目录下的参数调优记录是一个很容易被忽略的加分项把第 2、3、4 章里提到的参数调整过程和结果对比整理成表格能直接回应“这些参数你是怎么确定的”这类问题。data目录里的三个文件对应三个任务的数据源ratings.csv给 ALSreviews.txt给朴素贝叶斯user_edges.csv给图计算。5.2 文档说明里必须写清楚的三张图和一张表文档说明不要写成完整的代码注释复述而是要用图表把系统设计讲清楚。第一张是总体架构图画三个模块如何共用一个 SparkSession数据如何从 Socket 流入流式计算模块流式计算的结果如何落盘为 CSV 供图计算和机器学习模块读取。大作业里这种“计算结果下游复用”的设计非常加分比三个模块完全独立要高级得多。第二张是数据流图标注清楚每个阶段的数据格式转换——流式计算产出的是用户行为明细图计算需要的是用户间关系边ALS 需要的是评分三元组这三者要在数据流图里形成闭环。第三张是集群部署图如果用了三台云主机画清楚哪台跑 NameNode/ResourceManager、哪台跑 DataNode/NodeManager。一张参数调优表放在最后列出每个算法的核心参数、实验过的值、最终选定的值和理由。5.3 答辩前必做的验证命令三个模块能否独立重跑答辩现场最容易出的状况是某个模块跑不出来所以要准备一条命令能快速验证每个模块。我把这三条命令固化在 README 里每次答辩前按顺序执行一遍。# 1. 启动流数据源 while true; do echo user_${RANDOM:0:4} actionclick item_id$((RANDOM % 100)); sleep 1; done | nc -lk 9999 # 2. 跑流式计算观察控制台持续输出聚合结果 spark-submit --class streaming.UserBehaviorStreaming --master local[2] cloud-lab-1.0.jar # 3. 跑图计算输出 PageRank 和连通分量结果 spark-submit --class graph.SocialGraphAnalyzer --master local[2] cloud-lab-1.0.jar # 4. 跑机器学习三件套打印 RMSE、预测概率和聚类中心 spark-submit --class ml.ALSRecommender --master local[2] cloud-lab-1.0.jar spark-submit --class ml.NaiveBayesSentiment --master local[2] cloud-lab-1.0.jar spark-submit --class ml.KMeansCluster --master local[2] cloud-lab-1.0.jar这里有一个容易忽略的细节三个机器学习模块共用了spark.ml包在同一个 JVM 里运行没问题但如果你把三个作业串行提交到同一个 SparkContext就会报SparkContext already exists错误。解决方案是在每个对象的main方法里都调用spark.stop()或者用一个Main.scala按顺序调用三个任务的入口函数。最后一个技巧是启动spark-shell时注意不要和本地已有的 Spark 作业抢端口。如果你在前台启动了一个spark-submit作业占用了 4040 端口再起一个spark-shell的时候新的作业会自动使用 4041。这个现象不是报错但答辩时如果发现 Web UI 打不开要能反应过来是端口冲突。用lsof -i :4040查端口占用用SPARK_UI_PORT环境变量指定新的端口就能快速恢复。本文还有配套的精品资源点击获取
返回列表