数据库后端流处理【免费下载链接】EventStoreKurrentDB is a database thats engineered for modern software applications and event-driven architectures. Its event-native design simplifies data modeling and preserves data integrity while the integrated streaming engine solves distributed messaging challenges and ensures data consistency.项目地址https://gitcode.com/gh_mirrors/ev/EventStore点击查看免费下载KurrentDB 的 Webhook Source 连接器允许外部系统把 Webhook 调用直接投递到 KurrentDB每个入站请求的负载都会被写入一条 KurrentDB 事件流中间无需任何额外中间件。本文围绕 Webhook Source 官方文档 展开完整覆盖连接器的创建、配置参数、路由脚本、签名校验与响应语义并结合仓库源码给出底层实现佐证帮助你快速把 GitHub、Shopify、Slack、Stripe 等平台的 Webhook 接入 KurrentDB 事件驱动架构。工作原理概述Webhook Source 连接器的核心思路非常直接在 KurrentDB 节点上暴露一个POST /webhook/{connectorId}的 HTTP 端点外部系统向该端点发送 JSON 负载连接器收到请求后通过一段 JavaScript 路由脚本routingScript决定把负载写入哪条流、以什么事件类型写入随后将请求体作为事件数据持久化到 KurrentDB。从源码结构看这一功能位于 src/Connectors/KurrentDB.Connectors/Planes/Webhook/ 目录下WebhookPlaneWireUp.cs 负责把POST /webhook/{connectorId}路由注册到 ASP.NET Core 端点并将请求交给WebhookHandler.Ingest处理WebhookHandler.cs 实现请求的读取、路由结果的判定以及完整的 HTTP 状态码映射逻辑。连接器的实例类型名instanceTypeName为webhook-source它被登记在 ConnectorCatalogue.cs 的连接器目录中并需要CONNECTORS_WEBHOOK_SOURCE授权同样在该文件中可以看到WebhookSource实现了ISource接口因此被归类为 Source源类型连接器。前置条件使用 Webhook Source 连接器之前请确认以下条件已满足启用连接器的 KurrentDB 实例Connectors 插件已预装在所有 KurrentDB 二进制文件中且默认启用。若被禁用可在配置文件中通过Connectors: Enabled: false关闭参考 settings.md需要启用时设置为true即可。配置数据保护令牌Data Protection Token连接器需要加密signature:secret这类敏感字段因此 KurrentDB 实例中必须配置数据保护令牌。配置方式有两种详见 features.md#data-protection# 方式一使用令牌文件 Connectors: Enabled: true DataProtection: TokenFile: /path/to/token/file # 方式二直接指定令牌 Connectors: Enabled: true DataProtection: Token: your-secret-token注意一旦设置令牌即为永久值绝不能更改。所有连接器共用同一令牌更改会导致此前加密的数据全部不可读。若同时提供Token和TokenFile系统优先使用令牌文件。可公开访问的 KurrentDB 端点如果上游提供方如 GitHub、Stripe从外部网络发送 Webhook你需要一个公网可达的 KurrentDB 端点若 Webhook 由内网系统发出则无此要求。从 ConnectorDataProtectors.cs 可以看到WebhookSourceConnectorDataProtector把Signature:Secret声明为敏感键它会被连接器的数据保护框架加密存储而不是以明文落盘。快速上手创建并接收第一个 Webhook创建 Webhook Source 连接器只需要两步准备路由脚本、调用管理 API 创建连接器。第一步编写并编码路由脚本Webhook Source 要求提供一个routingScript它是一段base64 编码的 UTF-8 JavaScript 代码其中必须定义route(request)函数。对于每一个入站请求该函数返回目标流名称并可选地返回一个 schema 名称即事件类型。例如下面的函数根据请求体中的resource与action字段动态生成流名和事件类型function route(request) { return { stream: webhook-${request.body.resource}-${request.body.action}, schema: ${request.body.resource}.${request.body.action}, }; }将其 base64 编码后得到ZnVuY3Rpb24gcm91dGUocmVxdWVzdCkgeyByZXR1cm4geyBzdHJlYW06IGB3ZWJob29rLSR7cmVxdWVzdC5ib2R5LnJlc291cmNlfS0ke3JlcXVlc3QuYm9keS5hY3Rpb259YCwgc2NoZW1hOiBgJHtyZXF1ZXN0LmJvZHkucmVzb3VyY2V9LiR7cmVxdWVzdC5ib2R5LmFjdGlvbn1gIH07IH0第二步创建连接器使用管理 API 的POST /connectors/{id}创建连接器{id}替换为你想要的连接器 IDPOST /connectors/{id} Host: localhost:2113 Content-Type: application/json { settings: { instanceTypeName: webhook-source, routingScript: ZnVuY3Rpb24gcm91dGUocmVxdWVzdCkgeyByZXR1cm4geyBzdHJlYW06IGB3ZWJob29rLSR7cmVxdWVzdC5ib2R5LnJlc291cmNlfS0ke3JlcXVlc3QuYm9keS5hY3Rpb259YCwgc2NoZW1hOiBgJHtyZXF1ZXN0LmJvZHkucmVzb3VyY2V9LiR7cmVxdWVzdC5ib2R5LmFjdGlvbn1gIH07IH0 } }第三步发送负载到 Webhook 端点创建并启动连接器后向 Webhook 端点发送 JSON 负载POST /webhook/{id} Host: localhost:2113 Content-Type: application/json { resource: order, action: created, data: { orderId: abc-123 } }上述请求会把负载写入流webhook-order-created事件类型为order.created。默认情况下连接器一旦接受负载即返回202 Accepted仅代表进入处理管线不代表已持久化。管理 API 的完整端点列表见 Manage Connectors包括POST /connectors/{id}/start启动、POST /connectors/{id}/stop停止、GET /connectors列表、PUT /connectors/{id}/settings重配置、DELETE /connectors/{id}删除等操作。配置参数详解下表列出 Webhook Source 的全部专用设置项它同时继承所有 Source 连接器共有的通用配置见 Source Options其中包含instanceTypeName、logging:enabled等名称说明routingScript必填。base64 编码的 UTF-8 JavaScript定义route(request)函数为每个入站请求返回目标流名称与可选的 schema 名称。详见 路由脚本。waitForWrite启用后HTTP 请求会等待事件被持久化写入后才返回响应。默认false。详见 写确认。confirmationTimeout当waitForWrite启用时请求等待持久化写确认的最长时长超时返回504 Gateway Timeout。默认00:00:3030 秒。signature:scheme预配置的 HMAC 签名校验方案。当存在signature配置块时必填必须显式选择一种提供方。每个方案完整定义了该提供方期望的请求头名称、编码方式、签名负载形状与防重放行为。可选值GitHub、Shopify、Slack、Stripe。默认未设置必须指定。signature:secret共享 HMAC 密钥按 UTF-8 解释。当存在signature配置块时必须是非空字符串当整个signature配置块省略时签名校验被禁用。默认未设置校验禁用。signature:timestampTolerance当所选方案要求时间戳校验Slack、Stripe时允许的最大时钟偏差对不含时间戳的方案忽略该配置。默认00:05:005 分钟。allowedHeaders入站 HTTP 请求头名称的允许清单大小写不敏感这些请求头会被持久化到写入的事件上。每个被允许的请求头成为事件记录上独立的元数据条目键为其小写名称例如content-type。不在清单中的请求头绝不存储。Webhook 提供方经常发送携带凭据的请求头Authorization、Cookie、API 密钥等因此只应添加你已验证不含机密信息的请求头。默认[Content-Type, User-Agent, X-Request-Id]。提示signature:secret属于敏感字段由数据保护框架加密存储见 ConnectorDataProtectors.cs 中的WebhookSourceConnectorDataProtector这也是前置条件中必须配置数据保护令牌的原因。摄取端点Ingest EndpointPOST /webhook/{connectorId}Webhook Source 连接器只接受POST请求。每个请求必须使用 JSON 内容类型如application/json包括application/json; charsetutf-8且必须包含非空的 UTF-8 JSON 负载。该 JSON 负载将成为写入 KurrentDB 的事件数据。从 WebhookHandler.cs 的源码可以看到响应码的映射逻辑与文档完全一致响应码返回时机202 Accepted负载已被接受进入连接器处理管线。这是默认行为。201 Created启用了waitForWrite且事件已持久化写入。400 Bad RequestJSON 负载为空或不是合法的 UTF-8 JSON或路由脚本返回了null/空流名请求被跳过。401 Unauthorized签名校验已启用但签名请求头缺失、格式错误或时间戳超出容差窗口。404 Not Found当前节点上不存在该 ID 的活动连接器。415 Unsupported Media Type请求未使用 JSON 内容类型源码中对应!context.Request.HasJsonContentType()分支。500 Internal Server Error路由脚本抛出异常、执行超时、返回非对象值或返回了非法的stream/schema。503 Service Unavailable连接器无法接受负载或在waitForWrite请求等待确认期间连接器停止了源码中WriteConfirmation.Cancelled分支。504 Gateway Timeout启用了waitForWrite且写入在confirmationTimeout内未完成源码中WriteConfirmation.TimedOut分支。源码中WebhookHandler.Ingest通过WebhookIngestionResult的 switch 表达式完成映射SignatureFailed → 401、Rejected → 400、RoutingError → 500、Unavailable → 503而Accepted状态再根据写确认结果细分为201/504/503。路由脚本Routing Script路由脚本是 base64 编码的 UTF-8 JavaScript 代码必须定义一个名为route的函数。该函数在每个入站请求时被调用一次用于计算目标流名称和可选的schema 名称。请求对象结构route函数接收一个request参数包含以下字段字段类型说明bodyany解析后的 JSON 负载对象、数组或原始值。headersobject入站 HTTP 请求头键为小写名称、值为字符串。pathstring请求路径例如/webhook/my-connector。queryobject解析后的查询字符串键为名称、值为字符串。请求头键在传给脚本之前总是被转成小写。返回值route函数必须返回一个包含stream属性的对象即写入事件的目标流名称。流名必须是非空字符串且不能以$开头。可选的schema属性用于为事件打上自定义事件类型标签省略时事件以默认事件类型WebhookReceived写入。路由结果与响应函数的返回值决定了请求的最终去向路由请求返回带非空stream可选带schema的对象请求被写入对应流跳过请求返回null、undefined或带空stream的对象连接器返回400 Bad Request且不写入任何事件路由错误函数抛出异常、执行超时、返回非对象值或返回非法的stream/schema连接器返回500 Internal Server Error。示例按账号路由 Stripe 事件将 Stripe 事件路由到按账号划分的流并用 Stripe 事件类型为每条事件打标签function route(request) { return { stream: stripe-${request.body.account}, schema: stripe.${request.body.type}, }; }写确认Write Confirmation默认情况下连接器在负载进入内存管线后立即返回202 Accepted。202只表示接受处理不表示持久化完成。启用waitForWrite后只有当事件被持久化写入后才会收到201 Created。写入确认的等待时长由confirmationTimeout默认 30 秒控制若在时限内完成写入则返回201 Created超时则返回504 Gateway Timeout若等待期间连接器停止则返回503 Service Unavailable。对应的状态映射同样可在 WebhookHandler.cs 的MapConfirmation中查到。签名校验Signature Validation当配置了signature配置块时连接器会按照所选scheme的规则对每个入站请求校验 HMAC-SHA256 签名。签名缺失、格式错误或无效的请求会被拒绝并返回401 Unauthorized对于包含时间戳的方案超出配置容差窗口的请求同样返回401 Unauthorized。每个方案完整定义了请求头名称、编码、签名负载形状与防重放行为因此你只需选择方案并提供密钥无需关心各提供方的细节差异。各方案对比方案请求头编码签名负载时间戳校验GitHubX-Hub-Signature-256带sha256前缀Hexbody否ShopifyX-Shopify-Hmac-Sha256Base64body否SlackX-Slack-Signature带v0前缀、X-Slack-Request-TimestampHexv0:{timestamp}:{body}是StripeStripe-Signature解析t…,v1…键值对Hex{timestamp}.{body}是GitHubPOST /connectors/{id} Host: localhost:2113 Content-Type: application/json { settings: { instanceTypeName: webhook-source, routingScript: base64-encoded route() function, signature: { scheme: GitHub, secret: your-github-webhook-secret } } }ShopifyPOST /connectors/{id} Host: localhost:2113 Content-Type: application/json { settings: { instanceTypeName: webhook-source, routingScript: base64-encoded route() function, signature: { scheme: Shopify, secret: your-shopify-secret } } }SlackSlack 对v0:{timestamp}:{body}签名并要求时间戳请求头落在容差窗口内POST /connectors/{id} Host: localhost:2113 Content-Type: application/json { settings: { instanceTypeName: webhook-source, routingScript: base64-encoded route() function, signature: { scheme: Slack, secret: your-slack-signing-secret, timestampTolerance: 00:05:00 } } }StripeStripe 在单个Stripe-Signature请求头中发送逗号分隔的键值对t…,v1…签名负载为{timestamp}.{body}POST /connectors/{id} Host: localhost:2113 Content-Type: application/json { settings: { instanceTypeName: webhook-source, routingScript: base64-encoded route() function, signature: { scheme: Stripe, secret: your-stripe-webhook-secret, timestampTolerance: 00:05:00 } } }与管理 API 的衔接创建连接器之后可通过 Connectors 管理 API 完成完整生命周期管理启动POST /connectors/{id}/start可带日志位置POST /connectors/{id}/start/{log_position}从指定位置开始停止POST /connectors/{id}/stop列表与筛选GET /connectors支持state、instanceTypeName、connectorId、includeSettings、page、pageSize查询参数查看/修改设置GET /connectors/{id}/settings与PUT /connectors/{id}/settings重配置需重启后生效重置POST /connectors/{id}/reset或POST /connectors/{id}/reset/{log_position}删除/重命名DELETE /connectors/{id}与PUT /connectors/{id}/rename。总结Webhook Source 连接器让 KurrentDB 可以直接充当 Webhook 接收端外部系统一次POST /webhook/{connectorId}调用即可把 JSON 负载变成一条结构化事件写入流中配合routingScript可灵活实现按请求内容动态路由配合内置的 GitHub/Shopify/Slack/Stripe 签名校验方案可安全对接主流 SaaS 平台配合waitForWrite可精确控制持久化语义。结合仓库源码WebhookPlaneWireUp.cs、WebhookHandler.cs、ConnectorCatalogue.cs可以看到从端点注册、请求处理到敏感字段加密保护整套链路在 KurrentDB 内部均已实现闭环你只需要完成前置配置并按本文步骤创建连接器即可投入使用。赞分享数据库后端流处理【免费下载链接】EventStoreKurrentDB is a database thats engineered for modern software applications and event-driven architectures. Its event-native design simplifies data modeling and preserves data integrity while the integrated streaming engine solves distributed messaging challenges and ensures data consistency.项目地址https://gitcode.com/gh_mirrors/ev/EventStore点击查看免费下载相关推荐KurrentDB Kafka Source 连接器实战从 Kafka Topic 消费消息并写入 KurrentDB 流KurrentDB Kafka Source 连接器实战从 Kafka Topic 消费消息并写入 KurrentDB 流 本文基于 KurrentDB 官方数据库后端流处理PostHog Webhook终极指南事件触发与外部集成完全教程PostHog Webhook终极指南事件触发与外部集成完全教程 PostHog Webhook 是开源产品分析平台PostHog的核心功能之一它允许你将实数据分析后端前端数据可视化大数据Webhook轻量级入站Webhook服务器完全指南Webhook轻量级入站Webhook服务器完全指南 Webhook是一个用Go语言编写的轻量级、可配置的入站webhook服务器它允许开发者在服务器上轻松后端API网关上一篇basic-ftp安全指南FTPS配置与TLS参数优化详解下一篇OmniAuth与太空探索航天器认证系统的特殊挑战创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考