
人工智能机器学习后端模型推理服务大数据【免费下载链接】predictionioPredictionIO, a machine learning server for developers and ML engineers.项目地址https://gitcode.com/gh_mirrors/pre/predictionio点击查看免费下载导读Apache PredictionIO 的 Event Server 是机器学习数据管线的数据入口负责接收并持久化应用产生的各类事件。本文基于官方文档 docs/manual/source/datacollection/plugin.html.md 并结合仓库源码系统讲解 Event Server 插件机制两类插件Input Blocker与Input Sniffer的区别与适用场景、完整的插件工程搭建步骤、插件加载原理ServiceLoaderPIO_HOME/plugins目录、以及事件服务器提供的插件 REST API。读完本文你将能够独立开发、打包并部署一个可用于数据校验、日志记录、统计与转发的 Event Server 插件。Event Server 插件能做什么事件服务器插件用于处理进入 Event Server 的输入数据常见的应用场景包括拦截非法数据在事件写入事件存储Event Store之前校验数据发现非法事件时抛出异常以阻止其入库日志记录把流入的事件记录到日志系统便于审计与追踪统计实时汇总事件流量、类型分布等指标转发将事件转发到其他处理系统如消息队列、流处理平台、下游数据仓库。根据处理方式的不同Event Server 插件分为两类二者行为差异显著类型常量值执行时机是否阻塞写入是否可改写事件典型用途Input Blockerinputblocker事件进入事件存储之前会阻塞抛异常即拒绝入库不能改写事件输入数据校验、黑白名单过滤Input Snifferinputsniffer事件写入成功后并行广播不阻塞不能改写事件日志、统计、转发到下游系统官方文档明确指出两类插件的处理语义plugin.html.mdInput Blocker存在时进入事件服务器的事件会先依次经过所有已加载且激活的 Blocker再到达真正的事件存储。处理顺序未定义事件可能以任意顺序通过这些插件。典型用途是校验输入数据并抛出异常阻止坏数据入库。这类插件不能转换事件。Input Sniffer存在时事件会被并行广播给所有 Sniffer。它们不会阻止事件到达事件存储适合用于日志、统计和转发到其他处理系统。需要注意两类插件都不能转换改写事件内容它们只能观察或拦截同时 Blocker 的执行顺序未定义因此不要在多个 Blocker 之间依赖先后次序。开发一个 Event Server 插件1. 创建 sbt 工程首先创建一个独立的 sbt 项目build.sbt内容如下与官方文档示例一致name : pio-plugin-example version : 1.0 scalaVersion : 2.11.12 libraryDependencies org.apache.predictionio %% apache-predictionio-core % 0.14.0依赖坐标org.apache.predictionio %% apache-predictionio-core对应仓库中的 core 模块core/src/main/scala插件 API 的实体类则位于 data 模块的 data/src/main/scala/org/apache/predictionio/data/api 包内。2. 实现EventServerPlugin接口事件服务器插件必须继承实现EventServerPlugintrait。官方文档给出的完整示例package com.example import org.apache.predictionio.data.api._ class MyEventServerPlugin extends EventServerPlugin { val pluginName my-eventserver-plugin val pluginDescription an example of event server plug-in // inputBlocker or inputSniffer val pluginType EventServerPlugin.inputBlocker // Plug-in can handle input data in this method. // If plug-in found invalid data, its possible to block them // by throwing an exception in this method. override def process( eventInfo: EventInfo, context: EventServerPluginContext): Unit { println(eventInfo) } // Plug-in can handle requests to /plugins/pluginType/pluginName/* // on the event server in this method. override def handleREST( appId: Int, channelId: Option[Int], arguments: Seq[String]): String { {pluginName: my-eventserver-plugin} } }对照源码 EventServerPlugin.scala该 trait 定义了四个必须实现的成员pluginName: String插件唯一名称也是后续 REST API 路径中的标识符pluginDescription: String插件描述会出现在/plugins.json的响应中pluginType: String插件类型必须取EventServerPlugin.inputBlocker值inputblocker或EventServerPlugin.inputSniffer值inputsniffer常量定义见源码object EventServerPluginprocess(eventInfo: EventInfo, context: EventServerPluginContext): Unit事件处理入口。Input Blocker 若在此方法中抛出异常即可阻止该事件写入事件存储Input Sniffer 则在此方法中做观察式处理日志、统计、转发不能抛出异常阻断流程handleREST(appId: Int, channelId: Option[Int], arguments: Seq[String]): String处理指向/plugins/pluginType/pluginName/*的 HTTP 请求appId/channelId由事件服务器在鉴权后注入arguments为 URL 中插件名之后的路径分段。其中process方法收到的EventInfo定义在 EventInfo.scala结构为case class EventInfo( appId: Int, channelId: Option[Int], event: Event)即事件所属应用 ID、可选 channel ID以及完整的Event对象事件名、实体、属性、时间戳等。EventServerPluginContextEventServerPluginContext.scala则提供了inputBlockers与inputSniffers两个只读视图插件可以在process内感知当前已加载的其他插件。3. 通过META-INF/services注册插件插件由 JDK 标准的ServiceLoader机制加载因此必须在 jar 内创建服务描述文件META-INF/services/org.apache.predictionio.data.api.EventServerPlugin文件内容写入插件类的完整限定名每行一个实现类com.example.MyEventServerPlugin服务加载的实现见 EventServerPluginContext.scalaEventServerPluginContext.apply通过ServiceLoader.load(classOf[EventServerPlugin])扫描 classpath 上的所有实现并按pluginType分组存放到inputblocker与inputsniffer两个映射中。如果同一 jar 内存在多个插件只需在服务描述文件中逐行列出全部类名即可。4. 打包并部署到PIO_HOME/plugins执行sbt package将插件打成 jar 包。按文档示例产物位于target/scala-2.11/pio-plugin-example_2.11-1.0.jar将该 jar 复制到PIO_HOME/plugins目录然后启动或重启事件服务器插件即会被加载并启用cp target/scala-2.11/pio-plugin-example_2.11-1.0.jar $PIO_HOME/plugins/ pio eventserver为什么放在PIO_HOME/plugins从仓库工具代码可见启动引擎服务与批量预测时都会把该目录下的所有 jar 追加到 classpathOption(new File(pioHome, plugins).listFiles()).getOrElse(Array.empty[File]).map(_.toURI)见 RunServer.scala 与 RunBatchPredict.scala。pio eventserver命令最终调用Management.eventserverManagement.scala再经由EventServer.createEventServer构造路由并创建PluginsActor插件上下文在该过程中被一次性初始化。提示插件加载发生在事件服务器启动时因此新增/修改插件后必须重启事件服务器才能生效EventServerConfig中默认的插件目录名为plugins见 EventServer.scala与文档所述的PIO_HOME/plugins一致。事件服务器的插件相关 API事件服务器对外暴露了三类与插件相关的端点路由定义见 EventServer.scalaGET /plugins.json列出所有已启用的插件含 Blocker 与 Sniffer需要携带accessKey鉴权GET /plugins/inputblocker/pluginName/*由对应的 Input Blocker 插件的handleREST处理GET /plugins/inputsniffer/pluginName/*由对应的 Input Sniffer 插件的handleREST处理。查询已启用插件发送如下请求curl -XGET http://localhost:7070/plugins.json?accessKey$ACCESS_KEY事件服务器将返回如下 JSON与官方文档示例一致{ plugins: { inputblockers: { my-eventserver-plugin: { name: my-eventserver-plugin, description: an example of event server plug-in, class: com.example.MyEventServerPlugin } }, inputsniffers: {} } }该响应的字段name、description、class由源码在路由中直接组装pluginContext.inputBlockers.map { case (n, p) n - Map(name - p.pluginName, description - p.pluginDescription, class - p.getClass.getName) }EventServer.scala。若没有加载任何 Sniffer则inputsniffers为空对象如上所示。调用插件的自定义 REST 端点插件可以在handleREST中实现任意业务逻辑如查询插件内部的统计状态、动态切换过滤规则等外部通过如下模式调用GET /plugins/inputblocker/my-eventserver-plugin/arguments... GET /plugins/inputsniffer/my-eventserver-plugin/arguments...其中arguments...是插件名之后的路径分段会以Seq[String]形式传入handleREST。该端点同样需要accessKey鉴权可通过 URL 查询参数?accessKey...或 HTTPAuthorization: Basic头携带鉴权逻辑见 EventServer.scala鉴权成功后appId与channelId由服务器注入并一并传给插件。从源码理解插件在事件流水线中的位置Input Blocker写入前的同步校验以单事件写入端点POST /events.json为例路由中的关键调用顺序EventServer.scalaentity(as[Event]) { event if (events.isEmpty || authData.events.contains(event.event)) { pluginContext.inputBlockers.values.foreach( _.process(EventInfo(appId, channelId, event), pluginContext)) onSuccess(eventClient.futureInsert(event, appId, channelId)){ id pluginsActorRef ! EventInfo(appId, channelId, event) // 通知 Sniffer ... } } }可见先做事件名白名单校验AccessKey 允许的事件类型见authData.events依次同步调用所有 Input Blocker 的process——这正是“抛异常即可拦截”的实现点若某个 Blocker 抛出异常futureInsert不会执行事件被拒绝写入写入成功后才把EventInfo发送给PluginsActor触发 Sniffer 处理。批量写入端点POST /batch/events.json的流程同理先对每个可写事件逐个调用所有 Blocker 的processEventServer.scala批量插入成功后再逐一向PluginsActor广播EventInfoEventServer.scala。另外批量接口单次最多接收 50 个事件常量MaxNumberOfEventsPerBatchRequestEventServer.scala超出会返回400 Bad Request。Input Sniffer写入后的并行广播Sniffer 的process由专门的 Akka Actor——PluginsActor——负责调用。该 Actor 的receive逻辑PluginsActor.scalacase e: EventInfo pluginContext.inputSniffers.values.foreach(_.process(e, pluginContext))由于 Actor 的receive在同一时刻只处理一条消息而foreach依次调用每个 Sniffer 的process整体上 Sniffer 与主 HTTP 请求处理是异步解耦的——主流程在事件写入后立即返回不会等待 Sniffer 完成这正是“不阻塞事件到达事件存储”的工程实现。若需要在多个 Sniffer 之间并行执行可在各插件的process内部自行使用异步/线程池。此外/plugins/inputsniffer/pluginName/*的 REST 调用也经由PluginsActor.HandleREST消息路由到对应插件PluginsActor.scalaSniffer 的handleREST执行异常时会被捕获并返回{message:...}JSON而 Blocker 的 REST 调用则在路由中直接执行EventServer.scala。实战建议与注意事项用 Blocker 做校验用 Sniffer 做旁路数据合法性校验必填字段、事件名白名单、实体 ID 格式、数值范围等应放在 Input Blocker 中通过抛异常拒绝坏数据需要不影响主链路的日志、指标统计、消息转发则放在 Input Sniffer 中避免拖慢写入吞吐。Blocker 之间不要依赖顺序官方文档明确“处理顺序未定义”多个 Blocker 应彼此独立逻辑上不能假设先校验 A 再校验 B。不要试图改写事件两类插件均只能观察或拒绝不能修改EventInfo.event的内容如需清洗/增强数据应在数据进入事件服务器之前完成如 SDK 侧预处理。版本与目录保持一致插件构建需使用与你部署的 PredictionIO 版本匹配的apache-predictionio-core依赖与 Scala 版本官方示例为 0.14.0 Scala 2.11.12并把 jar 放入PIO_HOME/plugins后重启事件服务器。验证加载结果重启事件服务器后先请求GET /plugins.json确认插件名称、描述与类名是否正确出现在对应分组中再联调process与handleREST行为。扩展阅读事件数据的写入与查询完整指南Event Server API 文档数据采集整体流程与通道Channel机制数据采集总览插件 API 核心源码EventServerPlugin.scala、EventServerPluginContext.scala、PluginsActor.scala事件服务器路由与启动入口EventServer.scala事件服务器 CLI 命令实现Management.scala赞分享人工智能机器学习后端模型推理服务大数据【免费下载链接】predictionioPredictionIO, a machine learning server for developers and ML engineers.项目地址https://gitcode.com/gh_mirrors/pre/predictionio点击查看免费下载相关推荐PredictionIO Event Server 插件开发实战Input Blocker 与 Input Sniffer 数据接入治理全解析PredictionIO Event Server 插件开发实战Input Blocker 与 Input Sniffer 数据接入治理全解析 本文基于 Ap机器学习后端大数据PredictionIO 事件服务器插件开发指南Input Blocker 与 Input Sniffer 的完整实现与原理剖析PredictionIO 事件服务器插件开发指南Input Blocker 与 Input Sniffer 的完整实现与原理剖析 本文围绕 Apache Pr机器学习后端推荐系统PredictionIO 引擎服务器插件Engine Server Plugin开发指南Output Blocker 与 Output Sniffer 实战PredictionIO 引擎服务器插件Engine Server Plugin开发指南Output Blocker 与 Output Sniffer 实机器学习后端推荐系统上一篇Steam成就管理器完整指南如何快速解锁游戏成就的终极解决方案下一篇Miller 表达式语言的复杂度设计从 put-DSL 的简短性到速度与简单性的平衡创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考