PredictionIO 推荐模板批量预测结果持久化:Batch Persistable Evaluator 从零到实战

发布时间:2026/10/6 2:01:22
PredictionIO 推荐模板批量预测结果持久化:Batch Persistable Evaluator 从零到实战 人工智能机器学习后端模型推理服务大数据【免费下载链接】predictionioPredictionIO, a machine learning server for developers and ML engineers.项目地址https://gitcode.com/gh_mirrors/pre/predictionio点击查看免费下载PredictionIO 的pio eval通常用于计算引擎在验证集上的评估指标如 PrecisionK。本教程介绍一种特殊用法通过自定义一个批量可持久化 EvaluatorBatch Persistable Evaluator让pio eval直接对一批构造好的 Query 执行批量预测并把Query PredictedResult逐行序列化为 JSON 落盘到输出目录实现离线批量推荐结果导出。读完本文你将掌握如何改造DataSource.readEval()生成自定义批量查询、如何编写不计算指标只写结果的BaseEvaluator子类、如何组装Evaluation与EngineParamsGenerator并通过pio build/pio eval一键产出批量预测结果。适用前提本文基于 Recommendation 模板 v0.3.2且涉及的功能属于实验性 / 开发者特性接口在未来版本中可能发生变化。建议先阅读 推荐模板评估机制详解理解readEval()与 Evaluation 组件的基础用法后再继续。一、思路概述为什么需要批量持久化评估常规的pio eval工作流是评估导向的它以MetricEvaluator为核心在验证集上计算 PrecisionK、PositiveCount 等指标用于挑选最优引擎参数。而批量预测场景的目标完全不同——我们并不关心指标得分而是希望一次性对大量 Query 调用已训练模型把每条 Query 与其预测结果批量写出供下游报表、离线分析或二次加工使用。实现这一目标的关键差异在于 Evaluator 的职责MetricEvaluator对(Query, PredictedResult, ActualResult)三元组逐条计算指标汇总成评价分数自写的BatchPersistableEvaluator不做任何指标计算仅将三元组中的Query与PredictedResult映射成一行 JSON通过 Spark 的saveAsTextFile写入输出目录。从 BaseEvaluator 抽象类定义 可以看到所有 Evaluator 的统一入口是evaluateBase(sc, evaluation, engineEvalDataSet, params)方法其中engineEvalDataSet的类型为Seq[(EngineParams, Seq[(EI, RDD[(Q, P, A)])])]即引擎参数 →评估信息预测三元组 RDD的嵌套结构。我们只需要在该方法内消费第一个三元组 RDD 并写盘即可这正是批量持久化的全部秘密。二、步骤 1改造 DataSource用readEval()生成批量查询预测一批什么数据完全由DataSource.readEval()决定。readEval()的返回类型是Seq[(TrainingData, EmptyEvaluationInfo, RDD[(Query, ActualResult)])]即若干个训练数据 评估信息 查询与实际结果对的集合。批量预测场景下我们把它改造为只返回一个评估数据集其中RDD[(Query, ActualResult)]中装的是我们自定义的批量 Query而ActualResult全部用空的 rating 数组占位因为我们不关心真实结果只想要模型输出。override def readEval(sc: SparkContext) : Seq[(TrainingData, EmptyEvaluationInfo, RDD[(Query, ActualResult)])] { // This function only return one evaluation data set // Create your own queries here. Below are provided as examples. // for example, you may get all distinct user id from the trainingData to create the Query val batchQueries: RDD[Query] sc.parallelize( Seq( Query(user 1, num 10), Query(user 3, num 15), Query(user 5, num 20) ) ) val queryAndActual: RDD[(Query, ActualResult)] batchQueries.map (q // the ActualResult contain dummy empty rating array // because we not interested in Actual result for batch predict purpose. (q, ActualResult(Array())) ) val evalDataSet ( readTraining(sc), new EmptyEvaluationInfo(), queryAndActual ) Seq(evalDataSet) }这段代码的要点Query 的构造完全自由。示例中硬编码了三个 Query用户 1/3/5分别要 10/15/20 条推荐文档注释指出更常见的做法是从训练数据中取出所有去重后的 user id 来批量生成 Query。ActualResult(Array())只是占位符因为批量预测不关心实际答案但它满足了三元组类型的完整性要求。训练数据照常复用readTraining(sc)保证 Evaluator 拿到的是训练好的同一套数据。返回结构用Seq(...)包裹单个数据集对应下面 Evaluator 中engineEvalDataSet.size 1的约束。替代方案更解耦也可以不改动原 DataSource而是新建一个 DataSource 类继承原始 DataSource、仅覆写readEval()然后在Engine.scala的引擎工厂中注册新类并在engine.json的datasource段指定使用哪一个。这样原始 DataSource 的线上训练逻辑完全不受影响。以仓库中的 推荐模板 DataSource 实现 为参照其readEval()默认实现是 k 折切分由DataSourceEvalParams(kFold, queryNum)控制覆写时建议保留readTraining/getRatings的复用逻辑只替换测试用户 → Query的生成策略。三、步骤 2编写BatchPersistableEvaluator.scala新建文件BatchPersistableEvaluator.scala放在模板工程的src/main/scala/org/template/recommendation/包下。与MetricEvaluator不同这个 Evaluator不执行任何指标计算只是把 Query 与对应的 PredictedResult 写入输出目录package org.template.recommendation import org.apache.predictionio.controller.EmptyEvaluationInfo import org.apache.predictionio.controller.Engine import org.apache.predictionio.controller.EngineParams import org.apache.predictionio.controller.EngineParamsGenerator import org.apache.predictionio.controller.Evaluation import org.apache.predictionio.controller.Params import org.apache.predictionio.core.BaseEvaluator import org.apache.predictionio.core.BaseEvaluatorResult import org.apache.predictionio.workflow.WorkflowParams import org.apache.spark.SparkContext import org.apache.spark.rdd.RDD import org.json4s.DefaultFormats import org.json4s.Formats import org.json4s.native.Serialization import grizzled.slf4j.Logger class BatchPersistableEvaluatorResult extends BaseEvaluatorResult {} class BatchPersistableEvaluator extends BaseEvaluator[ EmptyEvaluationInfo, Query, PredictedResult, ActualResult, BatchPersistableEvaluatorResult] { transient lazy val logger Logger[this.type] // A helper object for the json4s serialization case class Row(query: Query, predictedResult: PredictedResult) extends Serializable def evaluateBase( sc: SparkContext, evaluation: Evaluation, engineEvalDataSet: Seq[( EngineParams, Seq[(EmptyEvaluationInfo, RDD[(Query, PredictedResult, ActualResult)])])], params: WorkflowParams): BatchPersistableEvaluatorResult { /** Extract the first data, as we are only interested in the first * evaluation. It is possible to relax this restriction, and have the * output logic below to write to different directory for different engine * params. */ require( engineEvalDataSet.size 1, There should be only one engine params) val evalDataSet engineEvalDataSet.head._2 require(evalDataSet.size 1, There should be only one RDD[(Q, P, A)]) val qpaRDD evalDataSet.head._2 // qpaRDD contains 3 queries we specified in readEval, the corresponding // predictedResults, and the dummy actual result. /** The output directory. Better to use absolute path if you run on cluster. * If your database has a Hadoop interface, you can also convert the * following to write to your database in parallel as well. */ val outputDir batch_result logger.info(Writing result to disk) qpaRDD .map { case (q, p, a) Row(q, p) } .map { row // Convert into a json implicit val formats: Formats DefaultFormats Serialization.write(row) } .saveAsTextFile(outputDir) logger.info(sResult can be found in $outputDir) new BatchPersistableEvaluatorResult() } }对照 BaseEvaluator 源码 逐点拆解泛型参数BaseEvaluator[EI, Q, P, A, ER]依次为评估信息类型、Query 类型、PredictedResult 类型、ActualResult 类型和结果类型这里分别填入EmptyEvaluationInfo、Query、PredictedResult、ActualResult与自定义的BatchPersistableEvaluatorResult。evaluateBase签名必须与父类抽象方法完全一致params: WorkflowParams承载工作流参数见 WorkflowParams 定义包含batch、verbose、saveModel、sparkEnv等。require双重校验要求恰好一组引擎参数、恰好一个RDD[(Q, P, A)]与我们readEval()返回Seq单元素的结构严格对应若未来要支持多组引擎参数可按文档注释的思路为不同引擎参数写不同目录来放开该约束。Row(query, predictedResult)一个用于 json4s 序列化的辅助 case class并且extends Serializable以适配 Spark 分布式执行。序列化与落盘implicit val formats: Formats DefaultFormats提供 json4s 默认格式Serialization.write(row)把每行转成 JSON 字符串最后saveAsTextFile(outputDir)以 Hadoop 文本文件形式目录下通常生成part-00000、_SUCCESS等文件写入输出目录。输出目录默认相对路径batch_result文档注释明确建议在集群上运行时改用绝对路径若底层存储支持 Hadoop 接口也可以把写盘逻辑改造成并行写入数据库。四、步骤 3定义BatchEvaluation与EngineParamsGenerator新建文件BatchEvaluation.scala把上面写好的 Evaluator 与推荐引擎绑定并提供引擎参数列表package org.template.recommendation import org.apache.predictionio.controller.EngineParamsGenerator import org.apache.predictionio.controller.EngineParams import org.apache.predictionio.controller.Evaluation object BatchEvaluation extends Evaluation { // Define Engine and Evaluator used in Evaluation /** * Specify the new BatchPersistableEvaluator. */ engineEvaluator (RecommendationEngine(), new BatchPersistableEvaluator()) } object BatchEngineParamsList extends EngineParamsGenerator { // We only interest in a single engine params. engineParamsList Seq( EngineParams( dataSourceParams DataSourceParams(appName INVALID_APP_NAME, evalParams None), algorithmParamsList Seq((als, ALSAlgorithmParams( rank 10, numIterations 20, lambda 0.01, seed Some(3L)))))) }关键点engineEvaluator二元组(RecommendationEngine(), new BatchPersistableEvaluator())明确告诉pio eval使用哪个引擎工厂、哪个 Evaluator。这与常规评估中(RecommendationEngine(), MetricEvaluator(...))的组装方式完全一致可对比 Evaluation.scala 中的 MetricEvaluator 用法。engineParamsList只含一组参数批量预测不需要做参数搜索因此只给出一组EngineParams与 Evaluator 中engineEvalDataSet.size 1的约束呼应。DataSourceParams(appName, evalParams None)appName必须改为你导入数据时使用的应用名文档中示例值为INVALID_APP_NAME占位直接照抄会读不到事件数据由于我们已把批量 Query 写死在readEval()里evalParams传None即可无需 k 折等评估参数。ALS 算法参数ALSAlgorithmParams(rank, numIterations, lambda, seed)与 ALSAlgorithm.scala 的参数定义 一致也可与 engine.json 的算法配置 保持同一组超参保证离线批量结果与线上训练口径一致。五、步骤 4构建并运行先构建引擎$ pio build构建成功后应看到[INFO] [Console$] Your engine is ready for training.随后以全限定类名分别指定Evaluation对象与EngineParamsGenerator对象运行$ pio eval org.template.recommendation.BatchEvaluation org.template.recommendation.BatchEngineParamsList预期控制台输出[INFO] [BatchPersistableEvaluator] Writing result to disk [INFO] [BatchPersistableEvaluator] Result can be found in batch_result [INFO] [CoreWorkflow$] Updating evaluation instance with result: org.template.recommendation.BatchPersistableEvaluatorResult2f886889 [INFO] [CoreWorkflow$] runEvaluation completed运行完成后在输出目录batch_result/下即可找到批量查询与预测结果——每行是一个 JSON 对象形如{query:{...},predictedResult:{itemScores:[...]}}。日志中的BatchPersistableEvaluatorResult2f886889是结果对象的 toString 输出说明整个评估工作流runEvaluation已正常收尾。关于pio eval的两个必填参数推荐模板评估机制详解 给出了权威解释Evaluation 对象告诉 PredictionIO 评估使用哪个引擎与哪个 EvaluatorEngineParamsGenerator包含一组或多组待评估/待执行的引擎参数。六、原理纵深批量预测从 Query 到 JSON 的完整链路把本文各部分串起来整个批量预测持久化流程是pio eval启动CoreWorkflow.runEvaluation调用BatchEngineParamsList.engineParamsList取出唯一一组EngineParams引擎按EngineFactory见 Engine.scala 中的 RecommendationEngine实例化 DataSource / Preparator / ALSAlgorithm / Serving 四个组件DataSource.readEval()生成RDD[(Query, ActualResult)]批量 Query 空 ActualResult 占位引擎对每个 Query 执行批量预测——ALS 算法同时实现了单条predict与面向评估的批量batchPredict见 ALSAlgorithm.scala 的 batchPredict后者通过cartesian展开用户 × 物品全组合并调用 MLlib ALS 的model.predict再按用户分组取 Top-N三元组 RDDRDD[(Query, PredictedResult, ActualResult)]传入BatchPersistableEvaluator.evaluateBaseEvaluator 剥离 ActualResult把(Query, PredictedResult)包装成Row用 json4s 序列化为 JSON 文本行saveAsTextFile(batch_result)落盘返回BatchPersistableEvaluatorResultrunEvaluation完成并输出日志。值得注意的工程细节Batch 预测与在线预测的差异在线predict逐个 Query 处理评估/批量场景使用batchPredict一次处理RDD[(Long, Query)]效率更高这也是pio eval能支撑大批量查询的原因。写盘位置与规模saveAsTextFile会按 Spark 分区并行写出因此批量既指查询条数多也指写入本身是分布式的如果结果要落数据库可在 Evaluator 内改用数据库的并行写入 API。结果校验batch_result/目录下除数据文件外还有_SUCCESS标记文件可通过cat batch_result/part-*或pio eval后读取该目录来核对每行 JSON 是否与readEval()中构造的 Query 一一对应。七、常见问题与注意事项appName配错DataSourceParams(appName ...)必须指向真实应用名否则PEventStore.find读不到事件数据批量预测会返回空结果。请参照 DataSource.scala 的 getRatings 事件读取逻辑 核对事件名与实体类型。集群运行时输出目录相对路径batch_result在集群模式下各 Executor 的工作目录可能不一致务必改为绝对路径如hdfs://.../batch_result或本地绝对路径。require约束engineEvalDataSet.size 1与evalDataSet.size 1是写死的前提若readEval()返回多个评估数据集、或engineParamsList含多组参数会直接抛异常终止需要按注释改造输出逻辑按引擎参数区分目录。实验性接口BaseEvaluator、evaluateBase、WorkflowParams均标注为 Developer API未来版本可能调整生产使用前请锁定 PredictionIO 版本并回归验证。本教程对应的完整模板代码可在仓库的 scala-parallel-recommendation 示例目录 中找到其DataSource、Evaluation、ALSAlgorithm等文件可作为本文全部代码片段的落地参照。赞分享人工智能机器学习后端模型推理服务大数据【免费下载链接】predictionioPredictionIO, a machine learning server for developers and ML engineers.项目地址https://gitcode.com/gh_mirrors/pre/predictionio点击查看免费下载相关推荐PredictionIO 批量持久化评估器Batch Persistable Evaluator使用 pio eval 批量生成推荐预测结果实战指南PredictionIO 批量持久化评估器Batch Persistable Evaluator使用 pio eval 批量生成推荐预测结果实战指南 本指机器学习后端大数据PredictionIO 批量持久化评估器实战用 pio eval 为一组查询批量产出推荐预测结果PredictionIO 批量持久化评估器实战用 pio eval 为一组查询批量产出推荐预测结果 本文基于 Apache PredictionIO 的 Re机器学习后端推荐系统PredictionIO 批量预测Batch Predict实战指南用 Spark 并行处理大规模查询PredictionIO 批量预测Batch Predict实战指南用 Spark 并行处理大规模查询 批量预测Batch Predict是 Pred人工智能机器学习后端模型推理服务大数据上一篇大麦自动抢票脚本实操指南一次配好 Appium 环境10 分钟跑通完整链路下一篇Umi-OCR 完整介绍免费离线、支持批量处理的 OCR 文字识别工具创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考