PredictionIO 批量持久化评估器实战:用pio eval为一组查询批量产出推荐预测结果
【免费下载链接】predictionioPredictionIO, a machine learning server for developers and ML engineers.项目地址: https://gitcode.com/gh_mirrors/pred/predictionio
本文基于 Apache PredictionIO 的 Recommendation 模板(v0.3.2),完整讲解如何改造 DataSource 的
readEval()生成批量查询、编写一个不做指标计算而是把 (Query, PredictedResult) 直接落盘的BatchPersistableEvaluator,并配合Evaluation与EngineParamsGenerator通过一条pio eval命令输出批量预测文件。读完本文你将掌握:如何为任意一批用户/参数组合批量执行推荐预测、如何把结果以 JSON 文本文件形式持久化、以及该方案在源码层的运行链路(BaseEvaluator、Evaluation、CoreWorkflow.runEvaluation)。
适用前提与阅读前置
本文对应的模板版本为Recommendation template v0.3.2。文中使用$pio eval配合自定义 Evaluator 将一组查询的预测结果持久化到输出目录,属于实验性/开发者特性(experimental and developer features),未来版本可能发生变更。
在动手之前,建议先阅读 Evaluation Explained (Recommendation),理解两件事:
- DataSource 中
readEval()的职责——返回Seq[(TrainingData, EmptyEvaluationInfo, RDD[(Query, ActualResult)])],即训练数据、空的评估信息、以及「查询 → 实际结果」的配对 RDD; - Evaluation 组件的用法——
Evaluation定义了「用哪个引擎 + 哪个评估器」跑评估。
常规的推荐评估(如 evaluation.html.md.erb 中介绍的MetricEvaluator+PrecisionAtK)会计算准确率等指标,并把最优参数与指标分数打印出来。而本文要做的恰恰相反:不关心指标分数,只关心「给定一批查询,每个查询会得到什么推荐结果」——这正是批量离线预测(batch predict)的典型场景。
整体思路
整个方案由三部分组成,缺一不可:
| 组成 | 文件(建议命名) | 作用 |
|---|---|---|
| 改造后的 DataSource | DataSource.scala | 在readEval()中构造我们想要批量预测的Query列表,并配上 dummy 的ActualResult |
| 自定义 Evaluator | BatchPersistableEvaluator.scala | 继承BaseEvaluator,接收评估流水线给出的 (Query, PredictedResult, ActualResult) RDD,序列化为 JSON 后落盘 |
| 评估入口对象 | BatchEvaluation.scala | 定义Evaluation(绑定新 Evaluator)与EngineParamsGenerator(指定引擎参数),供pio eval调用 |
它们的配合方式:pio eval <Evaluation> <EngineParamsGenerator>启动评估工作流 → 工作流用EngineParams实例化引擎 → 引擎调用 DataSource 的readEval()得到批量查询 → 训练完成后对每个查询做预测 → 把三元组 RDD 交给 Evaluator → Evaluator 写入输出目录。
第 1 步:改造 DataSource 生成批量查询
1.1 覆写readEval()
在你模板的DataSource.scala中,把readEval()改为返回一批你希望批量预测的查询。下面这段代码是文档给出的示例实现:
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) }要点解读:
- 返回值结构与常规评估一致:仍然是
Seq[(TrainingData, EmptyEvaluationInfo, RDD[(Query, ActualResult)])],因此评估工作流无需任何额外适配即可消费它; - 查询内容完全自定义:示例用
sc.parallelize(Seq(...))硬编码了 3 个查询(用户 "1" 取 10 个推荐、用户 "3" 取 15 个、用户 "5" 取 20 个)。文档注释指出,更常见的做法是从训练数据中取出所有去重后的 user id 来构造查询(例如把readTraining得到的TrainingData中的用户集合映射为Query列表); ActualResult只放占位数据:因为批量预测不关心真实结果,这里用ActualResult(Array())填充空的评分数组即可,它的作用是让三元组类型完整、通过评估流水线的类型检查。
1.2 对比:常规评估的readEval()长什么样
为了理解上面的改动,可以对照 customize-serving 示例的 DataSource.scala。常规实现的readEval()是一个 k-fold 切分过程:用getRatings(sc).zipWithUniqueId给每条评分打上唯一 id,然后按idx % kFold把数据划分成训练集与测试集,再groupBy(_.user)为每个用户构造一条查询,并携带该用户在验证集中真实的ActualResult(ratings.toArray)。
对比两者可以看出核心差异:
- 常规评估:查询数量 = 验证集用户数,
ActualResult是真实评分,用于计算Precision@K等指标; - 批量预测:查询数量与内容完全由你指定,
ActualResult是空占位,评估器完全忽略它。
1.3 备选做法:新建一个 DataSource 子类
文档特别提示:也可以不修改原有 DataSource,而是新建一个继承原 DataSource 的类来覆写readEval()。这样原始模板代码保持不动,只在需要跑批量预测时切换数据源。具体步骤为:
- 新建子类并覆写
readEval(); - 在
Engine.scala中把该子类注册进Engine(例如new Engine(classOf[BatchDataSource], ...)); - 在
engine.json中指定使用该 Engine 配置。
(文档原处标注了 “TODO add more details”,即这一做法在文档中属于提示性内容,具体注册细节可参考引擎默认配置自行扩展。)
第 2 步:编写BatchPersistableEvaluator
2.1 为什么需要一个新的 Evaluator
PredictionIO 默认的MetricEvaluator会计算指标分数并把结果写进数据库,其工作方式见 Evaluation.scala:engineMetric_=会把Metric包装成MetricEvaluator。而我们不需要任何指标计算,只需要把「查询 + 预测结果」原样写盘,因此要新建一个 Evaluator。
新建文件BatchPersistableEvaluator.scala,完整代码如下:
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(s"Result can be found in $outputDir") new BatchPersistableEvaluatorResult() } }2.2 逐段理解这个 Evaluator
类型参数:BaseEvaluator[EmptyEvaluationInfo, Query, PredictedResult, ActualResult, BatchPersistableEvaluatorResult]。对照 BaseEvaluator.scala 的定义,五个类型参数依次是评估信息类EI、查询类Q、预测结果类P、实际结果类A、评估结果类ER。这里EI用EmptyEvaluationInfo,ER是自定义的BatchPersistableEvaluatorResult(继承BaseEvaluatorResult)。
evaluateBase方法:这是BaseEvaluator中唯一需要实现的方法。它的入参中,engineEvalDataSet是Seq[(EngineParams, Seq[(EI, RDD[(Q, P, A)])])]——外层对应一组引擎参数,内层对应一组评估数据集,最内层的RDD[(Q, P, A)]就是「查询、预测结果、实际结果」的三元组 RDD。
严格约束输入规模:代码里有两个require:
engineEvalDataSet.size == 1:只允许一组引擎参数。因为本示例只为单套参数(如rank=10, numIterations=20, lambda=0.01)输出一个结果目录,若传多组参数会直接抛异常;evalDataSet.size == 1:只允许一个 RDD。因为我们在readEval()中只Seq(evalDataSet)返回了一份数据。
文档注释说明:如果希望放宽限制,可以改造输出逻辑,让不同的 engine params 写到不同的目录。
落盘逻辑:
qpaRDD.map { case (q, p, a) => Row(q, p) }:丢弃不需要的a(dummy 实际结果),只保留Row(query, predictedResult);- 借助 json4s 的
Serialization.write(row)把每条记录序列化为 JSON 字符串。注意implicit val formats: Formats = DefaultFormats声明在 map 内部,配合import org.json4s.DefaultFormats / Formats / native.Serialization使用; saveAsTextFile(outputDir):把整个 RDD 以文本文件形式写入outputDir(outputDir = "batch_result",由局部变量指定,在集群上运行时建议改为绝对路径)。
关于输出目录的扩展:saveAsTextFile走的是 Spark 的 Hadoop 文件接口。因此如果存储系统支持 Hadoop 接口,可以把同样的逻辑改写成向数据库并行写入(见代码注释);HDFS 等支持 Hadoop 接口的存储天然可用。
2.3 源码依据:BaseEvaluator与BaseEvaluatorResult
BaseEvaluator.scala 是 PredictionIO 所有评估器的基类,被标注为@DeveloperApi。关键点:
evaluateBase(...)由评估工作流(Evaluation Workflow)调用,引擎开发者一般不要直接调用它;BaseEvaluatorResult提供toOneLiner()/toHTML()/toJSON()三个默认为空串的方法,用于把评估结果呈现到评估 UI,以及noSave标志控制结果是否写入数据库。
在本文的BatchPersistableEvaluatorResult中这些方法都保持默认(空实现),因此CoreWorkflow.runEvaluation更新评估实例时拿到的toOneLiner等均为空字符串,日志里展示的就是对象默认的toString(如org.template.recommendation.BatchPersistableEvaluatorResult@2f886889)。
第 3 步:定义Evaluation与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)))))) }3.1BatchEvaluation:绑定引擎与评估器
Evaluationtrait 的定义见 Evaluation.scala。它通过engineEvaluator这个 setter 接收「引擎 + 评估器」二元组,内部会校验「评估器最多只能设置一次」。这里绑定的是RecommendationEngine()(模板自带的引擎工厂,见模板 Engine.scala,内部把DataSource、Preparator、ALSAlgorithm、Serving组装在一起)和新建的BatchPersistableEvaluator。
对比默认模板:常规的RecommendationEvaluation绑定的是MetricEvaluator(metric = PrecisionAtK(...), otherMetrics = ...)(见 customize-serving 示例的 Evaluation.scala),而这里换成了不做指标计算的BatchPersistableEvaluator。
3.2BatchEngineParamsList:指定引擎参数
EngineParamsGenerator是一个包含engineParamsList的 trait,pio eval的第二个参数就指向它。这里的engineParamsList只包含一个EngineParams:
dataSourceParams:记得把appName从"INVALID_APP_NAME"改成你自己的应用名(即你导入事件数据时使用的 app 名),evalParams = None表示不启用 k-fold 切分(因为我们不需要DataSourceEvalParams,批量查询完全由覆写后的readEval()提供);algorithmParamsList:ALS 算法的参数,rank = 10(隐因子数量)、numIterations = 20(迭代次数)、lambda = 0.01(正则化系数)、seed = Some(3L)(随机种子)。这些参数可以直接沿用你训练时验证过的一组值。
3.3 参数解析的源码依据
EngineParams与WorkflowParams在 core 模块 与 WorkflowParams.scala 中定义。其中WorkflowParams还暴露了batch(本次运行的批次标签)、verbose(日志级别)、saveModel(是否持久化模型)、skipSanityCheck、stopAfterRead、stopAfterPrepare等参数,意味着评估工作流本身也可以通过命令行开关做细粒度控制(例如调试数据源时用--stop-after-read提前中止)。engineParamsList是一个Seq,常规调参场景下可以放多组参数做网格搜索(如 EngineParamsList 中对rank、numIterations的组合遍历),而批量预测场景通常只需要一组——这也正是BatchPersistableEvaluator里require(engineEvalDataSet.size == 1)的前提。
第 4 步:构建并运行批量评估
4.1 构建
在模板根目录执行:
$ pio build构建成功后,控制台应输出:
[INFO] [Console$] Your engine is ready for training.4.2 运行
执行pio eval,第一个参数是Evaluation对象全名,第二个参数是EngineParamsGenerator对象全名:
$ pio eval org.template.recommendation.BatchEvaluation org.template.recommendation.BatchEngineParamsList4.3 预期输出
运行成功后,你应该看到类似下面的日志:
[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.BatchPersistableEvaluatorResult@2f886889 [INFO] [CoreWorkflow$] runEvaluation completed解读这四行日志:
- 前两行来自
BatchPersistableEvaluator自身的logger.info,表明落盘开始与完成; - 后两行来自评估工作流入口 CoreWorkflow.runEvaluation:先打印
runEvaluation started,随后把评估实例写入数据库并更新其状态(EVALCOMPLETED),最后打印runEvaluation completed。由于BatchPersistableEvaluatorResult没有覆写toOneLiner等方法,日志中展示的是对象的默认字符串表示。
4.4 查看结果
在输出目录batch_result/下,你可以找到批量查询及其预测结果。saveAsTextFile产生的文件内容大致为每条记录一行 JSON,例如:
{"query":{"user":"1","num":10},"predictedResult":{"itemScores":[{"item":"i123","score":4.82},...]}} {"query":{"user":"3","num":15},"predictedResult":{"itemScores":[{"item":"i456","score":3.91},...]}} {"query":{"user":"5","num":20},"predictedResult":{"itemScores":[{"item":"i789","score":3.05},...]}}每个 JSON 对象对应一个Row(query, predictedResult):query是你传入的批量查询,predictedResult.itemScores是该用户按得分降序排列的推荐条目列表(条目数量不超过查询中的num,条目结构对应模板Engine.scala中的ItemScore(item: String, score: Double))。拿到这个文件后,你可以用任意脚本(如 Python/awk)解析,把它导入业务数据库或用于离线分析。
运行链路:从命令到结果文件的源码级复盘
最后把整条链路在源码层面对齐,方便你排查问题或做二次开发:
- 命令入口:
pio eval <Eval> <ParamsList>定位到CoreWorkflow.runEvaluation(CoreWorkflow.scala),它先创建评估实例记录,再调用EvaluationWorkflow.runEvaluation(...)把evaluation、engine、engineParamsList、evaluator交给评估工作流; - 数据准备:评估工作流按
EngineParams实例化引擎,引擎内的 DataSource 调用你覆写后的readEval()产出批量Query及 dummyActualResult; - 训练与预测:引擎使用
ALSAlgorithmParams(rank=10, numIterations=20, lambda=0.01, seed=Some(3))训练模型,并对每个Query产出PredictedResult; - 评估落盘:评估工作流把
RDD[(Query, PredictedResult, ActualResult)]交给BatchPersistableEvaluator.evaluateBase,后者过滤掉实际结果、序列化Row(query, predictedResult)并saveAsTextFile("batch_result"); - 收尾:
evaluateBase返回BatchPersistableEvaluatorResult,CoreWorkflow.runEvaluation更新评估实例状态为EVALCOMPLETED并打印runEvaluation completed(CoreWorkflow.scala)。
常见问题与注意事项
appName未修改:BatchEngineParamsList里默认是"INVALID_APP_NAME",不改成你的事件应用名会导致读取不到数据;- 相对路径的输出目录:
outputDir = "batch_result"是相对路径。在本地单机运行时没有问题;在集群上运行(如 Spark on YARN)时,任务可能在不同节点执行,务必改用绝对路径,或使用支持 Hadoop 接口的存储(如 HDFS)并通过绝对路径写入; require约束:如果修改了readEval()让它返回多个数据集,或让engineParamsList包含多组参数,BatchPersistableEvaluator会因require失败而中止——要么保持单数据集、单参数组,要么按注释改造输出逻辑(为不同参数组写不同目录);- 特性稳定性:
BaseEvaluator属于@DeveloperApi,本文方案依赖的接口均为实验性 API,升级 PredictionIO 版本后需要回归验证; - 与常规评估的取舍:如果你关心的是「哪组参数更好」,请继续使用
MetricEvaluator与PrecisionAtK(参考 Evaluation.scala 示例);如果你关心的是「给定一组参数,为这批用户批量产出推荐结果」,本文的BatchPersistableEvaluator正是为此设计。
【免费下载链接】predictionioPredictionIO, a machine learning server for developers and ML engineers.项目地址: https://gitcode.com/gh_mirrors/pred/predictionio
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考