
人工智能机器学习后端模型推理服务大数据【免费下载链接】predictionioPredictionIO, a machine learning server for developers and ML engineers.项目地址https://gitcode.com/gh_mirrors/pre/predictionio点击查看免费下载导读本文讲解如何为 Apache PredictionIO 的 Event Server 编写并集成Webhooks Connector把 SegmentIO、MailChimp 等第三方站点的 Webhooks 数据转换为标准的 Event JSON 并写入事件存储Event Store。读完本文你将掌握两种内置 Connector 抽象JsonConnector与FormConnector的接口契约、Event Server 的 Webhooks URL 约定、一个完整的 Connector 实现范例以及如何把新 Connector 注册进WebhooksConnectors并通过测试验证。1. Webhooks Connector 的定位与工作原理Event Server 除了直接接收 App 上报的事件还可以通过第三方站点的 Webhooks 服务收集数据例如 SegmentIO、MailChimp。要让这些异构数据进入 Event Store需要为每个第三方数据源实现一个Webhooks Connector——它的职责非常简单把第三方数据转换成 Event JSON仅此而已。从源码结构看这一设计把数据格式适配与事件写入彻底解耦Connector 只负责格式转换输出统一的 Event JSON写入、统计等职责由 Webhooks.scala 统一完成。目前项目内置两类 Connector 抽象位于 data/src/main/scala/org/apache/predictionio/data/webhooks/JsonConnector负责接收 JSON 数据FormConnector负责接收表单提交Form Submission数据。1.1 JsonConnector 接口package org.apache.predictionio.data.webhooks import org.json4s.JObject /** Connector for Webhooks connection */ private[predictionio] trait JsonConnector { // TODO: support conversion to multiple events? /** Convert from original JObject to Event JObject * param data original JObject recevived through webhooks * return Event JObject */ def toEventJson(data: JObject): JObject }即 JsonConnector.scala 中定义的核心方法输入第三方 Webhooks 的原始 JSONjson4s 的JObject输出 Event JSON同样是JObject。注释中的TODO: support conversion to multiple events?说明当前设计上一个 Webhooks 请求只转换成一个事件属于可扩展方向而非既有能力。Event Server 收集 Webhooks JSON 数据的 URL 路径http://EVENT SERVER URL/webhooks/CONNECTOR_NAME.json?accessKeyYOUR_ACCESS_KEYchannelCHANNEL_NAME1.2 FormConnector 接口package org.apache.predictionio.data.webhooks import org.json4s.JObject /** Connector for Webhooks connection with Form submission data format */ private[predictionio] trait FormConnector { // TODO: support conversion to multiple events? /** Convert from original Form submission data to Event JObject * param data Map of key-value pairs in String type received through webhooks * return Event JObject */ def toEventJson(data: Map[String, String]): JObject }即 FormConnector.scala 中定义的方法输入为Map[String, String]表单字段均为字符串输出 Event JSON。由于表单值全部是字符串Connector 内部需要自行完成toDouble、toBoolean等类型转换。Event Server 收集 Webhooks 表单提交数据的 URL 路径注意没有.json后缀http://EVENT SERVER URL/webhooks/CONNECTOR_NAME?accessKeyYOUR_ACCESS_KEYchannelCHANNEL_NAME1.3 关于 channel 参数的建议你可以不传channel参数让 Webhooks 数据进入默认 channel。但官方强烈建议为每种 Webhooks 数据创建专用 Channel例如为 SegmentIO 建一个 segmentio channel为 MailChimp 建一个 mailchimp channel这样数据管理与查询更简单Webhooks 数据不会与 App 的其他正常业务数据混杂。Channel 的创建与使用方式参见 docs/manual/source/datacollection/channel.html.md.erb。2. 完整示例为 ExampleJson 站点实现 JSON Connector假设有一个名为 ExampleJson 的第三方网站通过其 Webhooks 服务发送如下 JSON 数据我们希望把它采集进 Event Store。UserActionItem 原始 Webhooks JSON{ type: userActionItem, userId: as34smg4, event: do_something_on, itemId: kfjd312bc, context: { ip: 1.23.4.56, prop1: 2.345, prop2: value1 }, anotherPropertyA: 4.567, anotherPropertyB: false, timestamp: 2015-01-15T04:20:23.567Z }2.1 实现步骤一编写ExampleJsonConnector因为该站点发送的是 JSON 格式我们实现一个继承JsonConnector的object ExampleJsonConnectorpackage org.apache.predictionio.data.webhooks.examplejson import org.apache.predictionio.data.webhooks.JsonConnector import org.apache.predictionio.data.webhooks.ConnectorException import org.json4s.Formats import org.json4s.DefaultFormats import org.json4s.JObject private[predictionio] object ExampleJsonConnector extends JsonConnector { implicit val json4sFormats: Formats DefaultFormats override def toEventJson(data: JObject): JObject { val common try { data.extract[Common] } catch { case e: Exception throw new ConnectorException( sCannot extract Common field from ${data}. ${e.getMessage()}, e) } val json try { common.type match { case userAction toEventJson(common common, userAction data.extract[UserAction]) case userActionItem toEventJson(common common, userActionItem data.extract[UserActionItem]) case x: String throw new ConnectorException( sCannot convert unknown type ${x} to Event JSON.) } } catch { case e: ConnectorException throw e case e: Exception throw new ConnectorException( sCannot convert ${data} to eventJson. ${e.getMessage()}, e) } json } // Convert the UserAction JSON to Event JSON def toEventJson(common: Common, userAction: UserAction): JObject { import org.json4s.JsonDSL._ // map to EventAPI JSON val json (event - userAction.event) ~ (entityType - user) ~ (entityId - userAction.userId) ~ (eventTime - userAction.timestamp) ~ (properties - ( (context - userAction.context) ~ (anotherProperty1 - userAction.anotherProperty1) ~ (anotherProperty2 - userAction.anotherProperty2) )) json } // Convert the UserActionItem JSON to Event JSON def toEventJson(common: Common, userActionItem: UserActionItem): JObject { import org.json4s.JsonDSL._ // map to EventAPI JSON val json (event - userActionItem.event) ~ (entityType - user) ~ (entityId - userActionItem.userId) ~ (targetEntityType - item) ~ (targetEntityId - userActionItem.itemId) ~ (eventTime - userActionItem.timestamp) ~ (properties - ( (context - userActionItem.context) ~ (anotherPropertyA - userActionItem.anotherPropertyA) ~ (anotherPropertyB - userActionItem.anotherPropertyB) )) json } // Common required fields case class Common( type: String ) // User Actions fields case class UserAction ( userId: String, event: String, context: Option[JObject], anotherProperty1: Int, anotherProperty2: Option[String], timestamp: String ) // UserActionItem fields case class UserActionItem ( userId: String, event: String, itemId: String, context: JObject, anotherPropertyA: Option[Double], anotherPropertyB: Option[Boolean], timestamp: String ) }完整实现见 data/src/main/scala/org/apache/predictionio/data/webhooks/examplejson/ExampleJsonConnector.scala。这段代码体现了一个典型 Connector 的三段式结构提取公共字段先用data.extract[Common]取出type字段作为路由依据Common只有type: String一个必填字段解析失败时抛出ConnectorException按 type 分发common.type匹配userAction与userActionItem两种类型各自调用对应的转换方法未知类型抛出ConnectorException(Cannot convert unknown type ...)映射为 Event JSON利用 json4s 的JsonDSL构造目标 JObject。2.2 理解转换目标Event JSON 的字段语义从上面toEventJson的输出可以看到Event Server 期望的 Event JSON 包含以下核心字段这也是 ConnectorUtil.scala 中EventJson4sSupport.APISerializer能够反序列化的标准格式字段含义示例event事件名称do_something_onentityType主体类型userentityId主体 IDas34smg4targetEntityType目标类型可选itemtargetEntityId目标 ID可选kfjd312bceventTime事件时间ISO 86012015-01-15T04:20:23.567Zproperties自定义属性 JObject见上文示例注意可选字段的处理anotherPropertyA与anotherPropertyB声明为Option[Double]、Option[Boolean]在 json4s 中None会被序列化省略从而支持不含可选字段的 Webhooks 载荷而userActionItem的context是必填的JObjectuserAction的context则是Option[JObject]两者对可选性的要求不同编写 Connector 时需按第三方实际数据约定。2.3 错误处理约定ConnectorException上述代码反复抛出ConnectorException它定义在 ConnectorException.scalaprivate[predictionio] class ConnectorException(message: String, cause: Throwable) extends Exception(message, cause) { /** Webhooks Connnector Exception with cause being set to null */ def this(message: String) this(message, null) }它提供两个构造器带 cause 的用于包装底层解析异常保留原始堆栈和不带 cause 的。最佳实践是所有无法转换的情况都应抛出ConnectorException并尽量附带原始数据片段与底层异常信息便于排查第三方数据格式问题。2.4 目录组织约定官方明确要求每个站点一个独立目录。例如 SegmentIO 的 Connector 代码放在data/src/main/scala/org.apache.predictionio/data/webhooks/segmentio/对应测试放在data/src/test/scala/org/apache/predictionio/data/webhooks/segmentio/当前仓库中已经按此约定组织可参见 segmentio/ 与 mailchimp/ 目录。2.5 为 Connector 编写测试官方同时提供了测试规范data/src/test/scala/org/apache/predictionio/data/webhooks/examplejson/ExampleJsonConnectorSpec.scala。测试使用 specs2通过共享的ConnectorTestUtil完成输入 → 输出断言class ExampleJsonConnectorSpec extends Specification with ConnectorTestUtil { ExampleJsonConnector should { convert userAction to Event JSON in { // webhooks input val userAction { type: userAction, userId: as34smg4, event: do_something, ... } // expected converted Event JSON val expected { event: do_something, entityType: user, entityId: as34smg4, ... } check(ExampleJsonConnector, userAction, expected) } ... } }测试覆盖了四类典型场景userAction完整字段、userAction不含可选字段、userActionItem完整字段、userActionItem不含可选字段。建议在真实项目中至少覆盖完整载荷与缺省可选字段载荷两类确保Option字段的处理符合预期。其他真实 Connector 的测试可参考 SegmentIOConnectorSpec.scala 与 MailChimpConnectorSpec.scala。2.6 表单数据FormConnector示例对于表单提交数据可参考完整实现 data/src/main/scala/org/apache/predictionio/data/webhooks/exampleform/ExampleFormConnector.scala。其输入形如typeuserActionItem userIdas34smg4 eventdo_something_on itemIdkfjd312bc context[ip]1.23.4.56 context[prop1]2.345 context[prop2]value1 anotherPropertyA4.567 // optional anotherPropertyBfalse // optional timestamp2015-01-15T04:20:23.567Z注意两个与JsonConnector不同的实现要点嵌套字段用context[...]前缀表达表单是扁平的键值对嵌套属性通过data.get(context[ip])等键名拼接实现userActionToEventJson中还会用data.exists(_._1.startsWith(context[))判断是否存在整个context块字符串到类型的手动转换data(anotherProperty1).toInt、data(context[prop1]).toDouble、data.get(anotherPropertyB).map(_.toBoolean)——由于Map[String, String]中全是字符串必须显式转换且可选字段用Option.map保持None语义。对应测试见 data/src/test/scala/org/apache/predictionio/data/webhooks/exampleform/ExampleFormConnectorSpec.scala。3. 将 Connector 集成进 Event ServerConnector 实现完成后需要注册到 WebhooksConnectors.scala 中Event Server 才能通过 URL 路由到它package org.apache.predictionio.data.api import org.apache.predictionio.data.webhooks.JsonConnector import org.apache.predictionio.data.webhooks.FormConnector import org.apache.predictionio.data.webhooks.examplejson.ExampleJsonConnector // ADDED import org.apache.predictionio.data.webhooks.segmentio.SegmentIOConnector import org.apache.predictionio.data.webhooks.mailchimp.MailChimpConnector private[predictionio] object WebhooksConnectors { val json: Map[String, JsonConnector] Map( segmentio - SegmentIOConnector, examplejson - ExampleJsonConnector // ADDED ) val form: Map[String, FormConnector] Map( mailchimp - MailChimpConnector ) }注册表语义jsonMap 的 key 就是 JSON 类 Webhooks URL 中的CONNECTOR_NAMEformMap 的 key 对应表单类 URL 的CONNECTOR_NAME。Connector 名称将直接成为 Webhooks URL 的一部分。注册完成后重新编译 Apache PredictionIO向以下 URL 发送 ExampleJson 数据即可把数据写入对应 Access Key 所属的 Apphttp://EVENT SERVER URL/webhooks/examplejson.json?accessKeyYOUR_ACCESS_KEYchannelCHANNEL_NAME对FormConnectorURL 不带.json例如http://EVENT SERVER URL/webhooks/mailchimp?accessKeyYOUR_ACCESS_KEYchannelCHANNEL_NAME4. 源码视角Webhooks 请求在 Event Server 中的处理链路了解请求如何被路由有助于理解 Connector 在整个链路上的位置。路由与写入逻辑集中在 Webhooks.scala核心是四个方法postJson/postForm接收数据与getJson/getForm健康/握手探测。以postJson为例其执行流程可概括为在Future中通过WebhooksConnectors.json.get(web)按 URL 中的 Connector 名称查表若查不到返回StatusCodes.NotFound与消息webhooks connection for ${web} is not supported.若查到调用ConnectorUtil.toEvent(connector, data)完成转换再由eventClient.futureInsert(event, appId, channelId)写入事件存储写入成功返回StatusCodes.Created与eventId若开启 stats 统计还会向statsActorRef发送Bookkeeping(appId, status, event)消息。其中关键的一环是 ConnectorUtil.scala 的转换逻辑private[predictionio] object ConnectorUtil { implicit val eventJson4sFormats: Formats DefaultFormats new EventJson4sSupport.APISerializer // intentionally use EventJson4sSupport.APISerializer to convert // from JSON to Event object. Dont allow connector directly create // Event object so that the Event object formation is consistent // by enforcing JSON format def toEvent(connector: JsonConnector, data: JObject): Event { readEvent)) } def toEvent(connector: FormConnector, data: Map[String, String]): Event { readEvent)) } }postForm与之对称区别仅在于输入来自 Akka HTTP 的FormData通过data.fields.toMap转成Map[String, String]后交给FormConnector。从这段实现可以推断几个设计约束Connector 不得直接构造Event对象统一走EventJson4sSupport.APISerializer做JSON → Event反序列化从而强制 Connector 输出合法的 Event JSON 格式保证事件对象构建的一致性转换失败不会污染写入toEventJson抛出的ConnectorException会在Future中体现为请求失败而不会写入半成品事件。5. 参考现有真实 ConnectorSegmentIO 与 MailChimp在动手写自己的 Connector 之前建议先阅读仓库中两个生产级 Connector它们是接口契约 实现范式的最佳参照SegmentIOJSON 类data/src/main/scala/org/apache/predictionio/data/webhooks/segmentio/SegmentIOConnector.scala——展示如何把 SegmentIO 的多种 track/identify 载荷映射为 Event JSON其注册名segmentio对应 URL/webhooks/segmentio.jsonMailChimp表单类data/src/main/scala/org/apache/predictionio/data/webhooks/mailchimp/MailChimpConnector.scala——展示表单数据如何解析与类型转换其注册名mailchimp对应 URL/webhooks/mailchimp。两者的测试文件SegmentIOConnectorSpec.scala、MailChimpConnectorSpec.scala也演示了面向真实第三方载荷的测试写法。6. 接入清单与常见注意事项完成一个 Webhooks Connector 的全过程总结如下确定数据类型第三方 Webhooks 发送的是 JSON 还是表单据此选择继承JsonConnector或FormConnector编写转换逻辑在toEventJson中实现第三方数据 → Event JSON的字段映射包含event、entityType、entityId、eventTime、properties等核心字段所有异常路径抛出ConnectorException按站点建目录源码放data/src/main/scala/org/apache/predictionio/data/webhooks/站点名/测试放data/src/test/scala/org/apache/predictionio/data/webhooks/站点名/编写测试至少覆盖完整载荷与缺省可选字段载荷两种场景复用ConnectorTestUtil的check断言注册 Connector在 WebhooksConnectors.scala 的json或formMap 中增加名称 - Connector重新编译并验证向http://EVENT SERVER URL/webhooks/名称.jsonJSON 类或http://EVENT SERVER URL/webhooks/名称表单类发送数据带accessKey与可选的channel参数确认返回201 Created与eventId建议使用专用 Channel为每个数据源创建独立 channel便于管理与查询避免 Webhooks 数据与正常业务数据混在一起参见 Channel 文档。常见陷阱FormConnector的嵌套属性必须用context[key]这类带前缀的键名表达并自行完成字符串到数值/布尔的转换Option字段用于表达可选不要把必填字段声明为Option以免静默丢失数据Connector 名称一旦注册即成为公开 URL 的一部分命名应简洁、唯一如segmentio、mailchimp若第三方载荷包含多种业务类型如示例中的userAction/userActionItem务必用type等公共字段做显式分发未知类型一律抛ConnectorException避免静默丢弃。赞分享人工智能机器学习后端模型推理服务大数据【免费下载链接】predictionioPredictionIO, a machine learning server for developers and ML engineers.项目地址https://gitcode.com/gh_mirrors/pre/predictionio点击查看免费下载相关推荐PredictionIO Webhooks Connector 开发指南为 Event Server 编写第三方数据接入连接器PredictionIO Webhooks Connector 开发指南为 Event Server 编写第三方数据接入连接器 PredictionIO 的机器学习后端推荐系统PredictionIO Webhooks Connector 开发实战为 Event Server 编写第三方数据接入连接器PredictionIO Webhooks Connector 开发实战为 Event Server 编写第三方数据接入连接器 PredictionIO 的机器学习后端大数据为 Apache PredictionIO 贡献 SDK 的完整开发指南Event Client 与 Engine Client 的 REST 实现规范为 Apache PredictionIO 贡献 SDK 的完整开发指南Event Client 与 Engine Client 的 REST 实现规范 导读机器学习后端推荐系统上一篇RenderDoc Event ID 机制详解事件 ID、Action 与 API 参数的结构化数据映射下一篇使用 Deployer 零停机部署 TYPO3 项目完整实战指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考