ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Apache PredictionIO 技术指南:基于 Spark 与 Lambda 架构的机器学习服务器全解析

Apache PredictionIO 技术指南:基于 Spark 与 Lambda 架构的机器学习服务器全解析 Apache PredictionIO 技术指南基于 Spark 与 Lambda 架构的机器学习服务器全解析【免费下载链接】predictionioPredictionIO, a machine learning server for developers and ML engineers.项目地址: https://gitcode.com/gh_mirrors/pred/predictionio导读Apache PredictionIO 是一个面向开发者、数据科学家与终端用户的开源机器学习框架它把「事件收集、算法训练与部署、离线评估、REST API 在线查询预测结果」整合为一条完整的流水线并基于 Hadoop、HBase或其他数据库、Elasticsearch、Spark 等可扩展开源服务实现 Lambda 架构。本文以仓库根目录 README.md 为主体结合 core、data、tools 等源码与官方文档深入讲解其架构原理、安装方式、DASE 引擎开发模型、CLI 命令、存储配置与快速上手路径帮助你从零开始理解并实际运行 PredictionIO。一、PredictionIO 是什么定位与核心能力根据 README.md 的官方描述Apache PredictionIO 是一个开源机器学习框架面向开发者developers、数据科学家data scientists和终端用户end users核心能力包括事件收集event collection通过 Event Server 收集应用的用户行为事件数据算法部署deployment of algorithms将训练好的模型部署为可查询的预测服务评估evaluation支持对算法/模型进行离线评估与调参REST API 查询预测结果querying predictive results via REST APIs对外提供统一的 HTTP 查询接口。其底层构建在可水平扩展的开源服务之上包括Hadoop、HBase以及其他数据库、Elasticsearch、Spark并实现了业界所称的Lambda ArchitectureLambda 架构——即同时兼顾批量batch与实时speed两条数据处理路径从而在可扩展性、容错性和低延迟查询之间取得平衡。仓库内 docs/manual/source 目录存放了完整的官方文档使用 Middleman 构建core/src/main/scala/org/apache/predictionio 与 data/src/main/scala/org/apache/predictionio 则是框架的核心实现本文后续将结合这些源码逐步展开。二、系统架构Event Server、引擎实例与存储层下图来自 docs/manual/source/images/system-overview.png直观展示了 PredictionIO 的系统级架构从图中可以看出系统的三个关键组成部分Your App业务应用多个应用App 1、App 2持续向 PredictionIO 上报事件数据Event Server 与 Event Datastore事件服务器统一接收事件事件数据在存储层按app_id例如app_id:1、app_id:2隔离存放Engine Instance引擎实例接收查询请求Queries运行训练好的模型返回预测结果Prediction Results。2.1 Event Server事件收集的入口事件收集是 PredictionIO 数据管道的起点。Event Server 的实现位于 data/src/main/scala/org/apache/predictionio/data/api/EventServer.scala其默认配置如下case class EventServerConfig( ip: String localhost, port: Int 7070, plugins: String plugins, stats: Boolean false)即默认绑定localhost:7070。该实现基于 Akka HTTP 构建使用 json4s 处理事件序列化见Json4sProtocol并提供了批量事件请求支持其单次批量请求上限为 50 条MaxNumberOfEventsPerBatchRequest 50EventServer.scala。事件数据通过LEvents/PEvents两类接口写入不同的存储后端详见「存储抽象」一节。2.2 引擎实例训练与部署的统一载体每个 Engine Instance 对应一次「引擎定义 参数组合」的运行单元训练时它负责读取数据、训练模型部署时它对外响应查询。这一设计在 core/src/main/scala/org/apache/predictionio/controller/Engine.scala 中有清晰的体现——Engine类将整个数据处理链路串联起来其train方法与eval方法分别驱动训练与评估两条工作流。三、核心抽象DASE 引擎开发模型PredictionIO 的引擎Engine由四个标准组件构成即DASE 模型组件职责对应源码基类coreData Source读取并转换训练/评估数据BaseDataSourcecore/core/BaseDataSource.scalaAlgorithm基于准备好的数据训练模型、做批量预测BaseAlgorithmcore/core/BaseAlgorithm.scalaServing组合多个算法结果生成最终预测BaseServingcore/core/BaseServing.scalaEngine定义引擎与各组件之间的装配关系Enginecontroller/Engine.scala此外还有可选的Preparator数据预处理对应 core/core/BasePreparator.scala负责把训练数据转换为算法可直接消费的 prepared data。3.1 引擎定义示例以仓库内置的推荐引擎示例为例examples/scala-parallel-recommendation/train-with-view-event/src/main/scala/Engine.scalaobject RecommendationEngine extends EngineFactory { def apply() { new Engine( classOf[DataSource], classOf[Preparator], Map(als - classOf[ALSAlgorithm]), classOf[Serving]) } }这里Map(als - classOf[ALSAlgorithm])表示引擎注册了一个名为als的算法本例为基于 Spark MLlib 的 ALS 协同过滤。从 Engine.scala 的源码可以看到Engine类接受四个映射dataSourceClassMap、preparatorClassMap、algorithmClassMap、servingClassMap组件之间通过泛型类型[TD, EI, PD, Q, P, A]约束数据契约TD训练数据类型EI评估信息类型PD预处理后数据类型Q查询类型REST 查询入参P预测结果类型A真实值类型评估用。3.2 训练流水线的源码级解读Engine.trainEngine.scala的调用链揭示了训练过程的真实顺序根据EngineParams中的名字实例化 DataSource、Preparator 与一个或多个 AlgorithmDoer(...)调用dataSource.readTrainingBase(sc)读取训练数据调用preparator.prepareBase(sc, td)生成预处理数据对每个算法执行algo.trainBase(sc, pd)得到模型列表通过makeSerializableModels持久化模型engineInstanceId作为模型 ID 的一部分格式为engineInstanceId-算法索引-算法名。同时训练过程支持数据健全性检查SanityCheck见 controller/SanityCheck.scala可通过工作流参数关闭。Engine.evalEngine.scala则实现了评估流水线对每个评估切分eval split训练模型、批量预测batchPredictBase、由 Serving 组合多算法结果serveBase最终产出(评估信息, RDD[(Q, P, A)])供 Metric 计算指标。3.3 模型持久化机制PredictionIO 对训练产出的模型采用分层持久化策略Engine.scala实现了PersistentModel接口controller/PersistentModel.scala的模型走自定义持久化逻辑本地local算法模型会被整体序列化保存并行parallel算法若模型由巨大 RDD 构成则返回Unit下次部署时按需从零重训prepareDeploy中检测到Unit模型会自动重训。四、引擎参数与 engine.json如何配置一次训练引擎的参数组合通过engine.json描述。仓库示例 examples/scala-parallel-recommendation/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 } } ] }各字段含义与默认值说明字段说明id变体variant标识default为常用值engineFactory引擎工厂类全限定名指向EngineFactory实现datasource.params.appName读取哪个 App 的事件数据对应 Event Server 中注册的 Appalgorithms[].name算法注册名必须与引擎定义中的Mapkey 一致如alsalgorithms[].params算法参数。以 ALS 为例rank10隐因子数量、numIterations20迭代次数、lambda0.01正则化系数、seed3随机种子从源码看jValueToEngineParamsEngine.scala负责解析engine.json分别从datasource、preparator、algorithms、serving字段中提取参数并反序列化为对应的Params对象。若algorithms数组为空则回退为(, EmptyParams())。pio build时pio.sbt会依据引擎目录下的engine.json自动生成见 tools/commands/Engine.scala。五、安装方式README 提供了两条官方安装路径详细文档位于 docs/manual/source/install5.1 从二进制包 / 源码安装完整步骤见 安装文档核心流程包括从 Apache 镜像下载二进制发行包apache-predictionio-版本-bin.tar.gz或源码包解压验证发行版使用项目 KEYS 与签名校验文件进行 GPG 校验$ gpg --import KEYS $ gpg --verify apache-predictionio-版本-bin.tar.gz.asc apache-predictionio-版本-bin.tar.gz安装依赖Spark 为硬依赖详见 pio-env.sh.template复制conf/pio-env.sh.template为pio-env.sh并按站点实际情况编辑。5.2 Docker 安装仓库内提供了完整的 Docker 编排资源docker/docker-compose.yml一键拉起 PredictionIO 及其依赖服务docker/pio/DockerfilePredictionIO 主镜像docker/elasticsearch、docker/mysql、docker/pgsql、docker/localfs不同存储后端的 docker-compose 配置含 event / meta / model 三套仓库的拆分部署docker/jupyter内置 Jupyter 的分析环境可参考 docker/JUPYTER.md。Docker 安装的详细步骤见 Docker 安装文档。5.3 环境变量配置pio-env.sh安装后必须按需修改 conf/pio-env.sh.template复制为pio-env.sh。该文件分为两大部分核心配置节选关键项# SPARK_HOME: Apache Spark is a hard dependency and must be configured. SPARK_HOME$PIO_HOME/vendors/spark-2.1.1-bin-hadoop2.6 POSTGRES_JDBC_DRIVER$PIO_HOME/lib/postgresql-42.0.0.jar MYSQL_JDBC_DRIVER$PIO_HOME/lib/mysql-connector-java-5.1.41.jar # ES_CONF_DIR: 高级 Elasticsearch 配置 # HADOOP_CONF_DIR: 使用 Hadoop 2 时配置 # HBASE_CONF_DIR: 使用远程 HBase 集群时配置 # Filesystem paths where PredictionIO uses as block storage. PIO_FS_BASEDIR$HOME/.pio_store PIO_FS_ENGINESDIR$PIO_FS_BASEDIR/engines PIO_FS_TMPDIR$PIO_FS_BASEDIR/tmp存储配置PredictionIO 使用「仓库Repository 数据源Source」两级模型组织存储默认全部指向 PostgreSQLPIO_STORAGE_REPOSITORIES_METADATA_NAMEpio_meta PIO_STORAGE_REPOSITORIES_METADATA_SOURCEPGSQL PIO_STORAGE_REPOSITORIES_EVENTDATA_NAMEpio_event PIO_STORAGE_REPOSITORIES_EVENTDATA_SOURCEPGSQL PIO_STORAGE_REPOSITORIES_MODELDATA_NAMEpio_model PIO_STORAGE_REPOSITORIES_MODELDATA_SOURCEPGSQL PIO_STORAGE_SOURCES_PGSQL_TYPEjdbc PIO_STORAGE_SOURCES_PGSQL_URLjdbc:postgresql://localhost/pio PIO_STORAGE_SOURCES_PGSQL_USERNAMEpio PIO_STORAGE_SOURCES_PGSQL_PASSWORDpio三大仓库的职责为METADATA引擎实例、App 等元数据、EVENTDATA事件数据、MODELDATA训练模型。切换存储只需更改对应的PIO_STORAGE_SOURCES_*系列变量模板中已给出 MySQL、Elasticsearch含可选 HTTP Basic Auth、Local File System、HBase、AWS S3 的完整示例。六、CLI 命令速查PredictionIO 的所有交互均通过命令行完成格式为pio command [options] args...。命令分为三大类CLI 文档6.1 通用命令命令说明pio help显示用法摘要pio help command查看子命令详情pio version显示已安装的 PredictionIO 版本pio status显示安装路径及 PredictionIO 与其依赖服务的运行状态6.2 Event Server 相关命令命令说明pio eventserver启动 Event Serverpio app管理 Event Server 使用的 Apppio app>case class WorkflowParams( batch: String , verbose: Int 2, saveModel: Boolean true, sparkEnv: Map[String, String] MapString, String, skipSanityCheck: Boolean false, stopAfterRead: Boolean false, stopAfterPrepare: Boolean false)其中saveModeltrue表示训练后的模型会被持久化skipSanityCheck可跳过数据健全性检查以加快调试stopAfterRead/stopAfterPrepare便于分阶段排查数据管道问题。七、快速开始从模板引擎跑通第一个推荐系统README 提供了三个官方模板的快速开始入口对应文档均位于 docs/manual/source/templatesRecommendation Engine Template推荐引擎快速开始文档Similar Product Engine Template相似商品引擎快速开始文档Classification Engine Template分类引擎快速开始文档。同时examples 目录提供了多套可直接运行的 Scala 并行引擎示例例如 scala-parallel-recommendation含 blacklist-items、customize-data-prep、customize-serving、reading-custom-events、train-with-view-event 等变体、scala-parallel-similarproduct 与 scala-parallel-classification每套均包含完整的engine.json、template.json与数据导入脚本如 import_eventserver.py。典型的完整流程为# 1. 创建 App 并获取 Access Key pio app new MyApp1 # 2. 启动 Event Server 并导入事件数据 pio eventserver python data/import_eventserver.py --access_key 你的access_key # 3. 在引擎目录中构建、训练、部署 cd examples/scala-parallel-recommendation/train-with-view-event pio build pio train pio deploy # 4. 通过 REST API 查询预测结果 curl -H Content-Type: application/json -d {user: u1, num: 4} \ http://localhost:8000/queries.json集成测试脚本位于 tests/pio_tests/scenarios如quickstart_test.py、eventserver_test.py可用作端到端验证的参考。八、评估与调参不只是训练除了训练与部署PredictionIO 还内置了完整的评估体系见 docs/manual/source/evaluationEvaluation把数据切分为训练集与测试集对每个切分训练模型并批量预测Metric基于(Q, P, A)元组计算指标如 Precision、Recall、MAP实现位于 controller/Metric.scala 与 controller/MetricEvaluator.scalaEngine Params Generatorcontroller/EngineParamsGenerator.scala用于参数网格搜索可自动化地遍历超参数组合配合 Evaluation Dashboard 直观对比结果。评估流水线的核心逻辑由Engine.evalEngine.scala承载它对每个评估切分(TD, EI, RDD[(Q, A)])独立完成「预处理 → 训练 → 批量预测 → Serving 组合」最终产出带评估信息的RDD[(Q, P, A)]交给 Metric 计算从而实现算法对比与参数调优的科学决策。九、存储抽象与扩展后端PredictionIO 的事件与元数据存储通过 data/src/main/scala/org/apache/predictionio/data/storage 下的统一接口如LEvents、PEvents、Apps、AccessKeys、EngineInstances等抽象具体实现分散于storage/elasticsearchElasticsearch 后端storage/hbaseHBase 后端storage/jdbcJDBCPostgreSQL/MySQL后端storage/localfs本地文件系统后端storage/s3AWS S3 后端storage/hdfsHDFS 后端。通过修改 conf/pio-env.sh.template 中的PIO_STORAGE_SOURCES_*与PIO_STORAGE_REPOSITORIES_*变量即可切换后端例如将 METADATA 与 EVENTDATA 保留在 PostgreSQL、把 MODELDATA 切换到本地文件系统或 HBase。官方文档 切换存储后端指南 对此有系统讲解docker 目录下的 elasticsearch/mysql/pgsql/localfs 子目录也分别给出了各后端的完整 docker-compose 示例。十、参与社区与贡献报告 Bug / 请求特性使用 Apache JIRA项目代号 PIO提交社区与开发动态订阅用户邮件列表user-subscribepredictionio.apache.org与开发邮件列表dev-subscribepredictionio.apache.org贡献代码参见 CONTRIBUTING.md 与社区文档 贡献代码指南贡献文档文档以 Middleman 构建于 docs/manual可参考 贡献文档指南发布自己的模板/项目可在社区项目页登记。本项目基于Apache License 2.0开源见 LICENSE.txtApache 品牌与版权声明详见 NOTICE.txt。结语从架构到源码Apache PredictionIO 将「事件收集 → 数据预处理 → 算法训练 → 模型评估 → REST 部署」这条机器学习流水线抽象为一套清晰、可组合的 DASE 组件模型并借助 Spark 的分布式计算能力与可插拔的多存储后端落地 Lambda 架构。无论你是想快速验证一个推荐/分类/相似商品引擎还是需要深入定制算法与存储本文所梳理的仓库文档、示例与源码路径都可以作为后续探索的起点。【免费下载链接】predictionioPredictionIO, a machine learning server for developers and ML engineers.项目地址: https://gitcode.com/gh_mirrors/pred/predictionio创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表