Spark大数据分析与实战笔记(第八章 Spark MLlib 机器学习算法库-01)

📅 2026/8/21 20:55:29
Spark大数据分析与实战笔记(第八章 Spark MLlib 机器学习算法库-01)
文章目录每日一句正能量章节概要第8章 Spark MLlib 机器学习算法库8.1 初识机器学习8.1.1 什么是机器学习8.1.2 Spark机器学习工作流程8.1.3 机器学习的应用8.2 Spark 机器学习库MLlib的概述8.2.1 MLlib的简介8.2.2 Spark机器学习工作流程1. 数据准备与特征工程2. 模型训练与调优3. 模型评估与选择4. 模型部署与预测MLlib 工作流程示例PySpark每日一句正能量“在蝉鸣里午睡在雪夜里早归跟着节气过日子。”主动的聆听与融入让自然的节奏成为自己生活的节拍器。跟着自然的节奏过日子不是被动顺应而是主动与天地共呼吸。春生夏长秋收冬藏人在其中便也成了节气的一部分。章节概要MLlib是Spark提供的处理机器学习方面的功能库该库包含了许多机器学习算法开发者可以不需要深入了解机器学习算法就能开发出相关程序。本章将介绍Spark MLlib基本知识以及使用方法最后通过构建推荐引擎了解机器学习系统的构建思路及流程。第8章 Spark MLlib 机器学习算法库8.1 初识机器学习8.1.1 什么是机器学习随着互联网的高速发展被收集并应用于分析的数据量呈现出爆发式的增长面对如此量级的数据以及常见的实时利用该数据的需求单单依靠人工处理难免力不从心这就催生了所谓的大数据和机器学习系统。机器学习是一门多领域交叉学科涉及概率论、统计学、逼近论、凸分析、算法杂度理论等多门学科专门研究计算机如何模拟或实现人类的学习行为以获取新的知识或技能重新组织已有的知识结构使之不断改善自身的性能。通俗的讲传统计算机工作时需要接收指令并按照指令逐步执行最终得到计算结果机器学习是通过某种算法将历史数据进行训练得出某种模型当有新的数据提供时可以使用训练产生的模型对未来进行预测。机器学习是一种能沙赋予机器进行自主学习不依靠人工进行自主判断的技术它和人类对历史经验归的过程有着相似之处接下来通过图对机器学习和人类思考进行对比。在图中左图是机器学习的过程右图则是人类思考的过程。人类在学习成长的过程中积累了很多历史经验将经验进行归纳总结得到规律因此当我们遇到一些问题时总能从事物的发展规律找到方向进行推测而机器学习中的训练和预测过程可以近似看作人类的归纳和推测的过程从图中可以发现机器学习思想并不复杂仅仅是对人类学习成长的过程一个模拟 由于机器学习不是通过编程的形式得出结果因此它的处理过程不是因果的逻辑而是通过归纳思想得出的相关结论。根据数据类型和需求不同建模方式也会不同。在机器学习领域中按照学习方式分类可以让研究人员在建模和算法选择的时候考虑根据输入数据来选择合适的算法从而得到更好的效果通常机器学习可以分为下面几类有监督学习通过已有的训练样本(即已知数据以及其对应的输出)训练得到个最优模型 再利用这个模型将所有的输入映射为相应的输出对输出进行简单的判断从而实现分类的目的。例如分类、回归和推荐算法都属于有监督学习。无监督学习根据类别未知没有被标记的训练样本而需要直接对数据进行建模我们无法知道要预测的答案。例如聚类、降维和文本处理的某些特征提取都属于无监督学习。8.1.2 Spark机器学习工作流程Spark中的机器学习流程大致分为三个阶段即数据准备阶段、训练模型评估阶段以及部署预测阶段。数据准备阶段如图所示在数据准备阶段需要将数据收集系统采集的原始数据进行数据预处理清洗后的数据便于提取特征字段与标签字段从而生产机器学习所需的数据格式然后将数据随机分为3个部分即训练数据模块、验证数据模块和测试数据模块。训练模型评估阶段如图所示通过Spark MLlib库中的函数将训练数据转换为一种适合机器学习模型的表现形式对于许多模型来说可以将其理解为包含数值数据的向量或者矩阵。然后使用验证数据集对模型进行测试来判断准确率这个过程需要重复许多次才能得出最佳模型最后使用测试数据集再次检验最佳模型以避免过渡拟合的问题如果训练评估阶段阶段准确率很高而使用测试数据阶段准确率低就说明可能有过拟合的问题。部署预测阶段通过多次训练测试得到最佳模型后就可以部署到生产系统中在该阶段的生产系统数据经过特征提取产生数据特征使用最佳模型进行预测最终得到预测结果。这个过程也是重复检验最佳模型的阶段可以使生产系统环境下的预测更加准确。8.1.3 机器学习的应用电子商务机器学习在电商领域的应用主要涉及搜索、广告、推荐三个方面在机器学习的参与下搜索引擎能够更好的理解语义对用户搜索的关键词进行匹配同时它可以对点击率与转化率进行深度分析从而利于用户选择更加符合自己需求的商品。医疗普通医疗体系并不能永远保持精准且快速的诊断在目前研究阶段中技术人员利用机器学习对上百万个病例数据库的医学影像进行图像识别分析数据并训练模型帮助医生做出更精准高效的诊断。金融机器学习正在对金融行业产生重大的影像例如在金融领域最常见的应用是过程自动化该技术可以替代体力劳动从而提高生产力例如摩根大通推出了利用自然语言处理技术的智能合同的解决方案该解决方案可以从文件合同中提取重要数据大大节省了人工体力劳动成本机器学习还可以应用于风控领域银行通过大数据技术监控账户的交易参数分析持卡人的用户行为从而判断该持卡人信用级别。8.2 Spark 机器学习库MLlib的概述8.2.1 MLlib的简介MLlib采用Scala语言编写借助了函数式编程设计思想开发人员在开发的过程中只需要关注数据而不需要关注算法本身所有要做的就是传递参数和调试参数。其中MLlib库中包含了一些通用的机器学习算法和工具类包括分类、回归、聚类、降维等具体如图所示。MLlib主要包含两部分分别是底层基础和算法库。其中底层基础包括Spark的运行库、矩阵库和向量库向量8.2.2 Spark机器学习工作流程Spark MLlib 遵循一个清晰、分阶段的机器学习工作流程该流程与 8.1.2 节中概述的通用流程一脉相承但深度集成了 Spark 的分布式计算能力和 MLlib 提供的专用 API。其核心流程可以概括为数据准备与特征工程 - 模型训练与调优 - 模型评估与选择 - 模型部署与预测。1. 数据准备与特征工程这是流程的基石。MLlib 提供了丰富的工具来处理存储在 Spark DataFrame 或 RDD 中的原始数据。数据加载使用SparkSession从 HDFS、Hive、本地文件系统或各类数据库如 MySQL中读取数据。数据清洗处理缺失值、异常值、重复数据等。MLlib 的Imputer、StringIndexer等转换器Transformer可自动化部分清洗工作。特征提取与转换这是 MLlib 的核心优势所在。开发者使用VectorAssembler将多个特征列组合成一个特征向量使用StandardScaler、MinMaxScaler进行特征缩放使用PCA进行降维等。所有转换操作都通过Pipeline管道进行链式封装确保训练和预测阶段的一致性。2. 模型训练与调优在此阶段使用处理好的训练数据来构建机器学习模型。选择算法根据问题类型分类、回归、聚类等从 MLlib 的算法库中选择合适的估计器Estimator例如LinearRegression、RandomForestClassifier、KMeans。设置超参数为选定的算法设置初始参数如maxIter迭代次数、regParam正则化参数。交叉验证与网格搜索为了找到最优的超参数组合MLlib 提供了CrossValidator和TrainValidationSplit等工具可以自动进行多轮训练和验证避免过拟合并输出性能最佳的模型。3. 模型评估与选择使用独立的测试数据集来客观评估训练出的模型性能。评估指标MLlib 提供了丰富的评估器Evaluator如BinaryClassificationEvaluator用于二分类评估 AUC、PR曲线下面积、MulticlassClassificationEvaluator用于多分类评估准确率、F1-score、RegressionEvaluator用于回归评估 RMSE、R²。模型选择根据评估指标如 AUC 越高越好RMSE 越低越好从多个候选模型例如不同算法或不同超参数组合训练出的模型中选择性能最优的一个作为最终模型。4. 模型部署与预测将训练好的最佳模型投入生产环境对新数据进行预测。模型保存与加载使用model.save(path)将训练好的模型包括整个PipelineModel持久化到分布式存储如 HDFS或本地。在预测服务中使用PipelineModel.load(path)加载模型无需重新训练。批量/流式预测加载模型后调用其transform()方法即可对新的 DataFrame 进行批量预测。该 API 与 Spark Structured Streaming 无缝集成同样支持流式数据的实时预测。MLlib 工作流程示例PySpark以下是一个简化的逻辑回归分类示例展示了上述流程的关键步骤frompyspark.sqlimportSparkSessionfrompyspark.mlimportPipelinefrompyspark.ml.featureimportVectorAssembler,StringIndexerfrompyspark.ml.classificationimportLogisticRegressionfrompyspark.ml.evaluationimportBinaryClassificationEvaluatorfrompyspark.ml.tuningimportParamGridBuilder,CrossValidator# 1. 初始化SparkSession并加载数据sparkSparkSession.builder.appName(MLlibWorkflow).getOrCreate()dfspark.read.csv(path/to/your/data.csv,headerTrue,inferSchemaTrue)# 2. 特征工程将类别特征索引化并将数值特征组合成向量indexerStringIndexer(inputColcategory,outputColcategoryIndex)assemblerVectorAssembler(inputCols[feature1,feature2,categoryIndex],outputColfeatures)# 3. 定义算法估计器lrLogisticRegression(featuresColfeatures,labelCollabel)# 4. 构建管道pipelinePipeline(stages[indexer,assembler,lr])# 5. 划分训练集和测试集train_df,test_dfdf.randomSplit([0.7,0.3],seed42)# 6. 模型训练modelpipeline.fit(train_df)# 7. 在测试集上进行预测并评估predictionsmodel.transform(test_df)evaluatorBinaryClassificationEvaluator(labelCollabel,rawPredictionColrawPrediction,metricNameareaUnderROC)aucevaluator.evaluate(predictions)print(fTest AUC {auc})# 可选8. 超参数调优paramGridParamGridBuilder()\.addGrid(lr.regParam,[0.01,0.1,1.0])\.addGrid(lr.elasticNetParam,[0.0,0.5,1.0])\.build()crossvalCrossValidator(estimatorpipeline,estimatorParamMapsparamGrid,evaluatorevaluator,numFolds3)cv_modelcrossval.fit(train_df)best_modelcv_model.bestModelprint(fBest model AUC on test set:{evaluator.evaluate(best_model.transform(test_df))})# 9. 模型保存与加载用于部署best_model.write().overwrite().save(hdfs://path/to/saved_model)# loaded_model PipelineModel.load(hdfs://path/to/saved_model)流程总结MLlib 通过DataFrameAPI、Pipeline、Transformer、Estimator和Evaluator这一套高级抽象将复杂的分布式机器学习流程标准化、模块化极大地提升了开发效率和模型的可维护性。开发者只需按照此流程组织代码即可高效地构建可扩展的机器学习应用。转载自https://blog.csdn.net/u014727709/article/details/163802600欢迎 点赞✍评论⭐收藏欢迎指正