ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

PredictionIO 内置算法库解析:基于 Spark MLlib 的 Algorithm 体系与多算法引擎实战

PredictionIO 内置算法库解析:基于 Spark MLlib 的 Algorithm 体系与多算法引擎实战 机器学习后端推荐系统【免费下载链接】predictionioPredictionIO, a machine learning server for developers and ML engineers.项目地址https://gitcode.com/gh_mirrors/pred/predictionio点击查看免费下载Apache PredictionIO 的 Engine 采用 DASE 架构DataSource、Preparator、Algorithm、Serving其中Algorithm算法是决定模型质量的核心组件。本文围绕官方文档 Built-in Algorithm Libraries 展开系统讲解 PredictionIO 的算法抽象类体系、对 Spark MLlib 的原生支持并结合仓库中的分类、相似商品等示例给出切换算法与组合多算法的完整实战路径。读完本文你将掌握如何理解并选用内置算法、如何把引擎默认算法替换为其他 MLlib 算法以及如何在同一引擎中集成多个算法并在 Serving 层融合预测结果。一、Algorithm 组件在引擎中的定位PredictionIO 将一条完整的机器学习流水线抽象为四个可插拔组件DataSource读取并转换原始数据、Preparator生成可直接训练的数据、Algorithm训练模型并产出预测、Serving对查询返回最终预测结果。官方文档明确指出An engine can virtually call any algorithm in the Algorithm class.也就是说引擎可以调用 Algorithm 类体系中的任意算法实现算法与引擎之间通过统一的接口解耦——这正是 DASE 架构允许换算法不换引擎的根本原因。Serving 组件接收的predictedResults参数类型为Seq[PredictedResult]这从接口层面决定了一个引擎可以挂载多个算法并将多个预测结果交给 Serving 统一融合详见 Combining Multiple Algorithms 文档。二、算法基类PAlgorithm 与 P2LAlgorithm在源码 core/src/main/scala/org/apache/predictionio/controller/PAlgorithm.scala 与 core/src/main/scala/org/apache/predictionio/controller/P2LAlgorithm.scala 中PredictionIO 定义了两种核心算法基类基类训练特性模型特性典型场景PAlgorithm[PD, M, Q, P]可在集群上并行训练模型可包含RDD可随集群分布式保存大规模模型、需要继续以分布式方式预测的算法P2LAlgorithm[PD, M, Q, P]可并行训练产出单机可承载的本地模型MLlib 的多数模型如NaiveBayesModel、RandomForestModel以P2LAlgorithm为例其类型参数含义为PDPrepared Data准备好的训练数据、M训练产出的模型、Q查询输入、P预测输出。开发者只需实现两个方法train(sc: SparkContext, pd: PD): M—— 从准备好的数据训练出模型predict(model: M, query: Q): P—— 使用模型对单条查询做出预测。P2LAlgorithm还提供了默认的batchPredict实现对RDD[(Long, Q)]逐条调用predict并定义了模型持久化策略默认模型会被自动序列化保存若模型混入PersistentModeltrait则改由PersistentModel.save手动持久化并返回PersistentModelManifest供部署时通过PersistentModelLoader加载。此外仓库中还有面向 Java 引擎开发者的 PJavaAlgorithm.scala 作为并行算法的 Java 接口。三、对 Spark MLlib 的原生支持PredictionIO 目前对Spark MLlib机器学习库提供原生支持这也是引擎模板库Template Gallery 配置中众多模板的算法基础。仓库示例中的算法实现直接印证了这一支持分类NaiveBayes朴素贝叶斯与RandomForest随机森林见 NaiveBayesAlgorithm.scala 与 RandomForestAlgorithm.scala推荐ALS交替最小二乘协同过滤见 scala-parallel-similarproduct/multi-events-multi-algos 示例。以朴素贝叶斯算法为例NaiveBayesAlgorithm.scala 展示了标准的算法封装方式case class AlgorithmParams( lambda: Double ) extends Params // extends P2LAlgorithm because the MLlibs NaiveBayesModel doesnt contain RDD. class NaiveBayesAlgorithm(val ap: AlgorithmParams) extends P2LAlgorithm[PreparedData, NaiveBayesModel, Query, PredictedResult] { override def train(sc: SparkContext, data: PreparedData): NaiveBayesModel { // MLLib NaiveBayes cannot handle empty training data. require(data.labeledPoints.take(1).nonEmpty, sRDD[labeledPoints] in PreparedData cannot be empty. Please check if DataSource generates TrainingData and Preparator generates PreparedData correctly.) NaiveBayes.train(data.labeledPoints, ap.lambda) } override def predict(model: NaiveBayesModel, query: Query): PredictedResult { val label model.predict(Vectors.dense( Array(query.attr0, query.attr1, query.attr2) )) PredictedResult(label) } }这段代码揭示了三个关键事实算法参数以case class ... extends Params定义与engine.json中的params字段一一对应PredictionIO 会按 JSON 反序列化注入选择P2LAlgorithm而非PAlgorithm的判断依据是模型是否含 RDD——MLlib 的NaiveBayesModel、RandomForestModel都是本地模型因此继承P2LAlgorithmtrain中通过require做训练数据非空校验这是官方示例一贯的防御性写法避免 MLlib 在空数据上崩溃。四、切换算法从 NaiveBayes 到 RandomForest官方文档 Switching to Another Algorithm 指出每个引擎模板都带有默认算法切换算法只需修改 Algorithm 类无需改动 DataSource、Preparator 与 Serving。分类模板的 add-algorithm 示例说明文档完整演示了这一过程第 1 步实现新算法类。新建RandomForestAlgorithm将模型类型替换为RandomForestModel并在train中调用 MLlib 的RandomForest.trainClassifiercase class RandomForestAlgorithmParams( numClasses: Int, numTrees: Int, featureSubsetStrategy: String, impurity: String, maxDepth: Int, maxBins: Int ) extends Params class RandomForestAlgorithm(val ap: RandomForestAlgorithmParams) extends P2LAlgorithm[PreparedData, RandomForestModel, Query, PredictedResult] { override def train(sc: SparkContext, data: PreparedData): RandomForestModel { // Empty categoricalFeaturesInfo indicates all features are continuous. val categoricalFeaturesInfo Map[Int, Int]() RandomForest.trainClassifier( data.labeledPoints, ap.numClasses, categoricalFeaturesInfo, ap.numTrees, ap.featureSubsetStrategy, ap.impurity, ap.maxDepth, ap.maxBins) } override def predict(model: RandomForestModel, query: Query): PredictedResult { val label model.predict(Vectors.dense( Array(query.attr0, query.attr1, query.attr2) )) PredictedResult(label) } }完整源码见 RandomForestAlgorithm.scala代码中以// CHANGED注释标出了相对默认算法的全部改动点便于对照学习。第 2 步注册算法到引擎。在EngineFactory中把新算法加入算法映射此处同时保留了naive与randomforest两个实现见 Engine.scalaobject ClassificationEngine extends EngineFactory { def apply() { new Engine( classOf[DataSource], classOf[Preparator], Map(naive - classOf[NaiveBayesAlgorithm], randomforest - classOf[RandomForestAlgorithm]), classOf[Serving]) } }第 3 步更新engine.json。将默认算法替换为randomforest并填写其参数完整文件见 engine.json{ id: default, engineFactory: org.apache.predictionio.examples.classification.ClassificationEngine, datasource: { params: { appName: MyApp1 } }, algorithms: [ { name: randomforest, params: { numClasses: 4, numTrees: 5, featureSubsetStrategy: auto, impurity: gini, maxDepth: 4, maxBins: 100 } } ] }完成上述三步后重新执行pio build、pio train与pio deploy引擎即以随机森林替代默认的朴素贝叶斯对外服务。切换过程中需要注意params中的字段名必须与算法参数类RandomForestAlgorithmParams的字段一一对应否则反序列化会失败。五、组合多算法在同一引擎中融合多个模型官方文档 Combining Multiple Algorithms 说明了多算法机制引擎可以用多个算法训练出多个模型预测结果由 Serving 类统一融合。相似商品模板的 multi-events-multi-algos 示例文档 给出了端到端实现完整示例源码位于 examples/scala-parallel-similarproduct/multi-events-multi-algos。该示例在原有基于 view 事件训练 ALS 模型的基础上新增一个处理 like/dislike 事件的LikeAlgorithm演示了多事件 多算法的完整链路5.1 DataSource读取多类事件在DataSource.scala中新增LikeEvent数据结构含布尔字段like表示喜欢/不喜欢并通过PEventStore.find同时拉取like与dislike两种事件val likeEventsRDD: RDD[LikeEvent] PEventStore.find( appName dsp.appName, entityType Some(user), eventNames Some(List(like, dislike)), targetEntityType Some(Some(item)))(sc) .map { event event.event match { case like | dislike LikeEvent( user event.entityId, item event.targetEntityId.get, t event.eventTime.getMillis, like (event.event like)) case _ throw new Exception(sUnexpected event ${event} is read.) } }.cache()TrainingData相应增加likeEvents: RDD[LikeEvent]字段。5.2 Preparator透传新数据Preparator.scala只需把TrainingData中的likeEvents原样传给PreparedData无需额外转换逻辑。5.3 新算法用隐式正负反馈训练新的LikeAlgorithm直接继承原ALSAlgorithm并重写train()这是复用既有算法逻辑最经济的方式。其核心处理包括同一用户对同一商品多次 like/dislike 时按事件时间戳取最新一次reduceByKey保留t更大的值利用ALS.trainImplicit()支持负偏好的特性将 dislike 映射为评分-1、like 映射为1从而在无显式评分的场景下表达正负反馈注意负偏好不适用于显式评分的ALS.train()显式评分场景如 rate 事件应使用ALS.train()。class LikeAlgorithm(ap: ALSAlgorithmParams) extends ALSAlgorithm(ap) { override def train(sc: SparkContext, data: PreparedData): ALSModel { // ... require 非空校验 ... val mllibRatings data.likeEvents .map { r ((uindex, iindex), (r.like, r.t)) } .filter { case ((u, i), v) (u ! -1) (i ! -1) } .reduceByKey { case (v1, v2) // keep the latest value if (t1 t2) v1 else v2 } .map { case ((u, i), (like, t)) // With ALS.trainImplicit(), we can use negative value to indicate dislike MLlibRating(u, i, if (like) 1 else -1) } .cache() val m ALS.trainImplicit( ratings mllibRatings, rank ap.rank, iterations ap.numIterations, lambda ap.lambda, blocks -1, alpha 1.0, seed seed) // ... } }5.4 Serving标准化并融合多算法结果部署后查询 Query 会同时发给引擎内所有算法各算法返回的PredictedResult以Seq[PredictedResult]形式进入Serving.serve()。示例采用z-score 标准化后求和的融合策略先对每个算法的得分序列计算均值与标准差标准差为 0 时取 0表示所有项排名相等将每个算法的得分标准化为 z-score消除不同算法得分量纲的差异对相同 item 的标准化得分求和按总分降序取query.num个作为最终结果class Serving extends LServing[Query, PredictedResult] { override def serve(query: Query, predictedResults: Seq[PredictedResult]): PredictedResult { val standard: Seq[Array[ItemScore]] if (query.num 1) { predictedResults.map(_.itemScores) } else { // Standardize the score before combine (z-score) ... } // sum the standardized score if same item val combined standard.flatten .groupBy(_.item) .mapValues(itemScores itemScores.map(_.score).reduce(_ _)) .toArray .sortBy(_._2)(Ordering.Double.reverse) .take(query.num) .map { case (k, v) ItemScore(k, v) } PredictedResult(combined) } }官方文档同时提示多算法结果的融合方式并非唯一完全可以根据业务需求自定义如加权、取最大、级联过滤等。5.5 Engine.scala 与 engine.json注册多算法EngineFactory中注册三个算法并分别命名engine.json的algorithms数组按序列出每个算法的名称与参数Map( als - classOf[ALSAlgorithm], cooccurrence - classOf[CooccurrenceAlgorithm], likealgo - classOf[LikeAlgorithm]){ algorithms: [ { name: als, params: { rank: 10, numIterations: 20, lambda: 0.01, seed: 3 } }, { name: likealgo, params: { rank: 8, numIterations: 15, lambda: 0.01, seed: 3 } } ] }likealgo与als参数结构相同是因为LikeAlgorithm复用了ALSAlgorithmParams若新算法有独立参数类在params中按其字段配置即可。示例还附带导入脚本data/import_eventserver.py可通过python data/import_eventserver.py --access_key 你的 Access Key快速导入含 like/dislike 事件的样本数据若报TypeError: __init__() got an unexpected keyword argument access_key需升级 Python SDK。六、自定义算法从模板出发官方文档 Adding your own Algorithms 目前标注为 (Coming soon)尚无独立教程但从仓库示例可以推断出自定义算法的标准范式——本质上就是实现一个继承算法基类的类 在EngineFactory注册 在engine.json配置参数三步与上文切换算法流程一致。可参考的现成实现包括分类模板的 NaiveBayes / RandomForest 双算法实现add-algorithm阅读自定义属性做特征工程的变体reading-custom-properties多事件多算法的相似商品引擎multi-events-multi-algos算法抽象类的权威定义PAlgorithm.scala、P2LAlgorithm.scala。七、小结PredictionIO 的算法体系以PAlgorithm/P2LAlgorithm抽象类为统一契约天然支持 Spark MLlib 的各类算法并通过算法注册映射 engine.json参数配置实现算法与引擎的灵活解耦。基于这一体系开发者既可以一条命令切换到不同算法如 NaiveBayes → RandomForest也可以在一个引擎内组合多个算法并在 Serving 层融合结果如 ALS 隐式正负反馈 ALS为多事件、多信号的真实推荐与分类场景提供了可落地的工程方案。赞分享机器学习后端推荐系统【免费下载链接】predictionioPredictionIO, a machine learning server for developers and ML engineers.项目地址https://gitcode.com/gh_mirrors/pred/predictionio点击查看免费下载相关推荐PredictionIO 引擎定制化实战基于 DASE 架构修改数据源、算法与服务逻辑PredictionIO 引擎定制化实战基于 DASE 架构修改数据源、算法与服务逻辑 PredictionIO 提供的引擎模板均随附完整源码并遵循统一的机器学习后端推荐系统Apache Spark MLlib 分类与回归RDD 版 API全解析算法选型、原理与实战Apache Spark MLlib 分类与回归RDD 版 API全解析算法选型、原理与实战 本文基于 Apache Spark 仓库中的 docs/ml大数据数据分析批处理流处理机器学习图计算Apache Spark MLlib 聚类算法实战指南K-means、二分 K-means、GMM、LDA 与 PIC 全解析Apache Spark MLlib 聚类算法实战指南K means、二分 K means、GMM、LDA 与 PIC 全解析 本篇技术指南以 Apache大数据数据分析批处理流处理机器学习图计算创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表