ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

PredictionIO DASE 架构深度解析:从数据源到评估的引擎组件开发指南

2026/10/6 2:02:12 拓冰建站 浏览量
PredictionIO DASE 架构深度解析:从数据源到评估的引擎组件开发指南 人工智能机器学习后端模型推理服务大数据【免费下载链接】predictionioPredictionIO, a machine learning server for developers and ML engineers.项目地址https://gitcode.com/gh_mirrors/pre/predictionio点击查看免费下载导读DASE 是 Apache PredictionIO 引擎的核心架构范式它将一个预测引擎的代码解耦为Data Source数据源、Algorithm算法、Serving服务与Evaluation评估四个独立组件使开发者能够像搭积木一样自由组合、替换与复用引擎的各个部分。本文以官方《Learning DASE》文档为主线结合仓库内 controller 包的真实源码实现系统讲解 DASE 各组件的职责、引擎的训练与实时查询工作流以及基于PDataSource、PPreparator、P2LAlgorithm/PAlgorithm、LServing等基类的完整实现方法。读完本文你将能够独立编写一个可训练、可部署、可评估的自定义 DASE 引擎。什么是 DASE引擎代码的四大组成在 Apache PredictionIO 中一个引擎Engine的代码由 D-A-S-E 四个组件构成[D] Data Source数据源与 Data Preparator数据预备器Data Source负责从输入源读取数据并转换为期望的格式Data Preparator对数据进行预处理并将其转发给算法用于模型训练。前者解决读什么、怎么读后者解决怎么加工。[A] Algorithm算法Algorithm 组件包含机器学习算法本身及其参数设置它决定了预测模型是如何被构建出来的。算法是引擎的大脑直接产出可预测的模型。[S] Serving服务Serving 组件接收预测queries查询并返回预测结果。如果引擎配置了多个算法Serving 会将多个结果合并为一个同时业务特定的逻辑也可以在 Serving 中加入以进一步定制最终返回的结果。[E] Evaluation Metrics评估指标评估指标用数值分数量化预测的准确程度可用于比较不同算法或同一算法的不同参数设置是引擎调优与评测的标尺。Apache PredictionIO 帮助你模块化这些组件例如你可以为一个 Engine 构建多个 Serving 组件并在创建引擎时选择部署哪一个。引擎的角色训练模型与实时响应一个引擎的主要功能有两项使用训练数据训练模型并作为 Web 服务部署上线实时响应预测查询引擎通过指定以下内容将全部 DASE 组件整合为一个可部署的状态一个Data Source一个Data Preparator一个或多个Algorithm(s)一个ServingINFO如果指定了多个算法每个算法的模型预测结果都会被传递给 Serving 进行集成ensembling。每个引擎独立处理数据、独立构建预测模型因此每个引擎各自服务一套预测结果。例如你可以为移动应用部署两个引擎一个用于向用户推荐新闻另一个用于向用户推荐新朋友。训练视角pio train时 DASE 的工作流当执行pio train时DASE 组件按以下流程协作Data Source 从事件存储中读取并筛选事件产出TrainingDataData Preparator 对TrainingData做特征选择与数据加工产出算法所需的PreparedDataAlgorithm 基于PreparedData调用train()训练出模型随后模型被持久化存储供部署阶段加载使用。查询视角REST 查询到达已部署引擎时的 DASE 工作流当已部署的引擎收到 REST 查询时工作流则变为查询先进入 Serving必要时可先做 query 补充处理再被分发到各算法执行predict()实时预测最终由 Serving 将多个预测结果合并并返回给调用方。关于 DASE 的具体实现细节请继续阅读 Implement DASE下文将结合仓库源码逐组件展开。Data Source读取事件、构造训练数据Data Source 从 Event StoreEvent Server 的数据存储读取并筛选有用数据返回TrainingData。其并行版本基类为PDataSource定义在 PDataSource.scalaabstract class PDataSource[TD, EI, Q, A] extends BaseDataSource[TD, EI, Q, A] { def readTraining(sc: SparkContext): TD def readEval(sc: SparkContext): Seq[(TD, EI, RDD[(Q, A)])] Seq[(TD, EI, RDD[(Q, A)])]() }其中类型参数TD为训练数据类、EI为评估信息类、Q为查询类、A为实际结果类。除readTraining()外readEval()用于提供引擎的评估功能默认返回空序列即引擎可以不实现评估功能也能编译运行。实现readTraining()你需要实现PDataSource的readTraining()在其中使用 PEventStore 引擎 API 读取事件并基于事件创建TrainingData。下面示例读取用户的 view 和 buy 商品事件过滤出特定类型的事件供后续处理并据此返回TrainingDataclass DataSource(val dsp: DataSourceParams) extends PDataSource[TrainingData, EmptyEvaluationInfo, Query, EmptyActualResult] { transient lazy val logger Logger[this.type] override def readTraining(sc: SparkContext): TrainingData { val eventsRDD: RDD[Event] PEventStore.find( appName dsp.appName, entityType Some(user), eventNames Some(List(view, buy)), // targetEntityType is optional field of an event. targetEntityType Some(Some(item)))(sc) .cache() val viewEventsRDD: RDD[ViewEvent] eventsRDD .filter { event event.event view } .map { ... } ... new TrainingData(...) } }使用 PEventStore 引擎 APIPEventStore.find()是读取事件的入口其完整签名见 PEventStore.scala支持丰富的过滤条件参数含义appName读取该应用的事件必填channelName读取该通道的事件None表示默认通道startTime/untilTime按事件时间范围过滤eventTime startTime且 untilTimeentityType/entityId按实体类型/ID 过滤eventNames返回这些事件名中的任意一种targetEntityTypeNone不限制Some(None)表示事件无目标实体Some(Some(x))表示目标实体类型须为 xtargetEntityId同上语义按目标实体 ID 过滤例如假设你有如下事件{ event: myEvent, entityType: user, entityId: u0, targetEntityType: item, targetEntityId: i0, properties : { a : 3, b : some_string, c : [a, b, c], d : [1.2, 3.4, 5.6], e : 6 } }下面的代码可读取这些事件提取事件的properties字段并转换为MyEvent对象val myEvents: RDD[MyEvent] PEventStore.find( appName dsp.appName, entityType Some(user), eventNames Some(List(myEvent)), // targetEntityType is optional field of an event. targetEntityType Some(Some(item)))(sc) .map { event try { MyEvent( entityId event.entityId, targetEntityId event.targetEntityId.get, a event.properties.getInt, b event.properties.getString, c event.properties.get[List[String]](c), d event.properties.get[List[Double]](d), e event.properties.getOptInt // use getOpt for optional data ) } catch { case e: Exception logger.error(sCannot convert ${event}. Exception: ${e}.) throw e } }值得注意的细节properties.get[T]用于必填属性而properties.getOpt[T]用于可选属性当属性缺失或类型不匹配时get会抛出异常因此示例中用try/catch记录错误日志后重新抛出避免静默丢失脏数据。使用aggregateProperties()聚合实体属性如果你使用了特殊事件$set/$unset/$delete来设置实体的属性则可以用PEventStore.aggregateProperties()按实体聚合出当前属性快照。该方法在 PEventStore.scala 中定义签名如下def aggregateProperties( appName: String, entityType: String, channelName: Option[String] None, startTime: Option[DateTime] None, untilTime: Option[DateTime] None, required: Option[Seq[String]] None) (sc: SparkContext): RDD[(String, PropertyMap)]其中required参数可指定仅保留这些必选属性已定义的实体。返回值是(entityId, PropertyMap)的 RDD例如// create a RDD of (entityID, Item) val itemsRDD: RDD[(String, Item)] PEventStore.aggregateProperties( appName dsp.appName, entityType item )(sc).map { case (entityId, properties) try { val item Item( a properties.getInt, b properties.getString, c properties.get[List[String]](c), d properties.get[List[Double]](d), e properties.getOptInt // use getOpt for optional data ) (entityId, item) } catch { case e: Exception logger.error(sFailed to get properties ${properties} of ${entityId}. Exception: ${e}.) throw e } }关于$set/$unset/$delete特殊事件的语义与事件建模可参见 eventmodel 文档。示例更多 DataSource 写法参见 Similar Product 模板的 DataSource。Preparator预处理训练数据、产出 PreparedDataPreparator 负责对TrainingData进行预处理以完成必要的特征选择与数据处理任务生成算法所需的PreparedData。其并行版本基类PPreparator定义在 PPreparator.scalaabstract class PPreparator[TD, PD] extends BasePreparator[TD, PD] { def prepare(sc: SparkContext, trainingData: TD): PD }Preparator 的典型用途特征提取feature extraction当引擎有多个算法时承担公共的预处理逻辑对于简单场景Preparator 也可以直接将同一个TrainingData作为PreparedData传给 Algorithm实现prepare()你需要实现PPreparator的prepare()方法来完成上述任务。仓库中的两个典型案例Lead Scoring 模板的 Preparator对TrainingData进行预处理生成算法所需的特征向量Similar Product 模板的 Preparator仅简单地将TrainingData作为PreparedData传给算法即 Identity 预处理。Algorithm训练模型与实时预测Algorithm 类有两个核心方法train()与predict()。train()负责训练预测模型在执行pio train时被调用Apache PredictionIO 会存储该模型predict()负责使用模型进行预测当你向引擎发送 JSON 查询时被调用。注意predict()是实时调用的其性能直接影响线上响应延迟。Apache PredictionIO 支持两种类型的算法P2LAlgorithmparallel-to-local训练的模型不包含 RDD可以装入单机内存PAlgorithmparallel训练的模型包含 RDD模型本身也可分布式存放于集群。两种算法的公共抽象BaseAlgorithm[PD, M, Q, P]中PD为 PreparedData 类、M为模型类、Q为查询类、P为预测结果类。P2LAlgorithm本地模型自动持久化对于P2LAlgorithm模型在训练完成后会被 Apache PredictionIO自动序列化并持久化。从 P2LAlgorithm.scala 的makePersistentModel()实现可以看到默认情况下直接返回模型对象m交由框架自动持久化仅当模型混入了PersistentModeltrait 时才会调用m.save(...)走自定义持久化路径。因此对于P2LAlgorithm来说实现IPersistentModel与IPersistentModelLoader是可选的——默认自动持久化已经足够。示例参见 Similar Product 模板的 Algorithm。PAlgorithmRDD 模型手动持久化当你的模型包含 RDD 时应使用PAlgorithm。由PAlgorithm产生的模型默认不会被持久化见 PAlgorithm.scalacase _ ()分支直接返回空值意味着部署时需要重新训练。要实现模型持久化你需要做两件事模型类应扩展IPersistentModeltrait并实现save()方法用于保存模型。该 trait 需要一个类型参数即算法参数类在 PersistentModel.scala 中trait PersistentModel[AP : Params]save(id: String, params: AP, sc: SparkContext): Boolean返回true表示保存成功框架据此生成持久化清单返回false则部署时重新训练模型工厂对象应扩展IPersistentModelLoadertrait并实现apply()用于加载模型。该 trait 需要两个类型参数算法参数类与算法产出的模型类apply(id: String, params: AP, sc: Option[SparkContext]): Msc仅在模型由PAlgorithm产生时注入。注IPersistentModel/IPersistentModelLoader自 0.9.2 起已废弃推荐直接使用PersistentModel/PersistentModelLoader见 PersistentModel.scala二者 API 完全一致仅命名不同。示例Recommendation 模板的 Algorithm实现了PAlgorithm以及IPersistentModel与IPersistentModelLoaderVanilla 模板完整走查了P2LAlgorithm与PAlgorithm两种写法。在predict()中使用 LEventStore 做实时查询增强你可以在predict()中使用 LEventStore.findByEntity() 检索某个具体实体的近期事件——例如获取查询中指定用户的近期行为并利用这些近期事件进行实时预测如过滤用户已看过的商品、实现个性化重排。例如下面的代码读取query.user最近的 10 条 view 事件val recentEvents try { LEventStore.findByEntity( appName ap.appName, // entityType and entityId is specified for fast lookup entityType user, entityId query.user, eventNames Some(List(view)), targetEntityType Some(Some(item)), limit Some(10), latest true, // set time limit to avoid super long DB access timeout Duration(200, millis) ) } catch { case e: scala.concurrent.TimeoutException logger.error(sTimeout when read recent events. s Empty list is used. ${e}) Iterator[Event]() case e: Exception logger.error(sError when read recent events: ${e}) throw e }两点工程细节值得借鉴一是entityType与entityId指明按实体快速查找二是设置了timeout Duration(200, millis)的时间上限避免查询拖慢实时预测链路超时后回退为空列表而非直接失败。示例E-Commerce Recommendation 模板的 Algorithm 在predict()中用LEventStore.findByEntity()取出用户看过的所有商品并将其从推荐结果中过滤掉。Serving整合预测结果、定制返回逻辑Serving 组件的基类为LServing[Q, P]定义在 LServing.scala。除核心的serve()外它还提供可选的supplement(q: Q): Q方法标注为 Experimental用于在查询发给算法之前对 query 做补充处理。实现serve()你需要实现LServing的serve()方法def serve(query: Q, predictions: Seq[P]): Pserve()处理预测结果当你有多个预测模型时它也负责将多个预测结果合并为一个。注意query是发送给引擎的原始查询而非supplement()产生的补充后查询predictions是各算法的预测结果列表顺序与 engine.json 中算法声明顺序一致。典型示例Similar Product 模板的 Serving直接返回唯一的预测结果Similar Product 模板的多算法示例将多个算法的结果合并后返回。组装引擎engine.json 与引擎工厂DASE 各组件最终通过引擎工厂类engineFactory与engine.json配置组装成一个可部署的引擎。以仓库中的推荐引擎为例参见 train-with-view-event 示例的 engine.json{ id: default, description: Default settings, engineFactory: org.apache.predictionio.examples.recommendation.RecommendationEngine, datasource: { params : { appName: MyApp1 } }, algorithms: [ { name: als, params: { rank: 10, numIterations: 20, lambda: 0.01, seed: 3 } } ] }可以看到datasource.params.appName指定读取哪个应用的事件对应PDataSource中dsp.appName的用法algorithms数组可以声明一个或多个算法每个算法携带自己的超参数如 ALS 的rank、numIterations、lambda、seed它们会注入到对应算法类的参数对象中并在train()时生效——这也印证了原文档中一个引擎可指定多个算法、多算法结果交由 Serving 集成的设计。学习路线从模板到自定义引擎掌握了 DASE 原理后最快的学习路径是研读仓库中官方模板的 DASE 实现并在此基础上改造Recommendation 模板的 DASESimilar Product 模板的 DASEClassification 模板的 DASELead Scoring 模板的 DASE每个模板目录下都有完整的engine.json、template.json以及src/main/scala下的 DataSource / Preparator / Algorithm / Serving 实现例如 scala-parallel-recommendation 目录配合本文对PDataSource、PPreparator、P2LAlgorithm/PAlgorithm、LServing等基类源码的解析即可按 D-A-S-E 的骨架逐步替换组件构建属于自己的预测引擎。赞分享人工智能机器学习后端模型推理服务大数据【免费下载链接】predictionioPredictionIO, a machine learning server for developers and ML engineers.项目地址https://gitcode.com/gh_mirrors/pre/predictionio点击查看免费下载相关推荐PredictionIO 核心术语指南从 DASE 引擎架构到事件采集与离线评估PredictionIO 核心术语指南从 DASE 引擎架构到事件采集与离线评估 导读 PredictionIO 是一个面向开发者与机器学习工程师的机器学习服机器学习后端大数据PredictionIO 术语表DASE 引擎架构、Event Server 与数据评估体系核心概念详解PredictionIO 术语表DASE 引擎架构、Event Server 与数据评估体系核心概念详解 本篇技术指南围绕 Apache Prediction机器学习后端推荐系统PredictionIO 核心术语全解析DASE 引擎架构、Event Server 与离线/在线评估实战指南PredictionIO 核心术语全解析DASE 引擎架构、Event Server 与离线/在线评估实战指南 Apache PredictionIO 是一套人工智能机器学习后端模型推理服务大数据上一篇Better Auth SSO 插件 1.7 深度解读SAML/OIDC 企业单点登录的安全加固与账户身份模型重构下一篇P2Pool合并挖矿教程提升门罗币挖矿收益的高级技巧创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考