ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

用 Watermill 构建 Kafka 到 HTTP 的 Webhook 推送:sending-webhooks 示例全解析

用 Watermill 构建 Kafka 到 HTTP 的 Webhook 推送:sending-webhooks 示例全解析 用 Watermill 构建 Kafka 到 HTTP 的 Webhook 推送sending-webhooks 示例全解析【免费下载链接】watermillBuilding event-driven applications the easy way in Go.项目地址: https://gitcode.com/GitHub_Trending/wa/watermill本文基于仓库中的 sending-webhooks 示例讲解如何使用 Watermill 的消息路由能力将 Kafka 中的事件转换为 HTTP POST 请求即 Webhook 外呼。示例由 producerKafka 生产者、router事件路由与转换、webhooks-serverWebhook 接收端与 RedpandaKafka 兼容消息后端四个服务构成读完本文你将掌握 HTTP Publisher 的接入方式、按消息元数据分流到多个 Webhook 的实战手法以及整套服务的本地运行与日志观测方法。示例要解决的问题Kafka 事件 → 多个 Webhook 外呼在事件驱动架构中一个典型场景是业务事件写入消息中间件如 Kafka下游的第三方系统如 CRM、通知服务并不直接消费 Kafka而是以 HTTP Webhook 的形式被动接收通知。本示例正是这一模式的落地演示——从 Kafka 消费事件再以 HTTP POST 请求的形式把事件推送给外部接收方整个过程由 Watermill 的消息路由Router无缝衔接。示例围绕三种事件类型展开事件的类型通过消息元数据metadata的event_type键进行编码事件类型含义Foo类型 A 的事件Bar类型 B 的事件Baz类型 C 的事件整个示例由三个服务加一个消息后端组成详见 README.mdproducer向 Kafka 持续发布消息消息按随机顺序取Foo、Bar、Baz三种类型之一事件类型写入元数据键event_typewebhooks_server一个极简 HTTP 服务器监听请求并把路径path、方法method与请求体payload打印到标准输出router消费 Kafka 消息使用HTTP Publisher向webhooks_server发送请求。为了演示一条消息可以派生出多个 Webhook路由会根据event_type调用不同的路径/foo仅接收Foo类型事件/foo_or_bar接收Foo或Bar类型事件/all接收所有类型事件。此外kafka服务一个 Redpanda broker为 Kafka 生产者与订阅者提供消息后端。提示该示例同时被 docs/content/pubsubs/http.md 引用作为把 Kafka 消息转换为 HTTP Webhook 请求的官方配套示例对应的反向场景HTTP 收 Webhook 写入 Kafka可参考仓库中的 receiving-webhooks 示例。整体运行流程docker-compose up启动后数据流如下producer ──(Publish: topickafka_to_http_example)──▶ Kafka(Redpanda) │ ▼ router ──(Subscribe: 同一 topic)─────────────────── Kafka Subscriber │ ├─ 处理器 foo ──(POST http://webhooks-server:8001/foo)──▶ webhooks-server ├─ 处理器 foo_or_bar ──(POST http://webhooks-server:8001/foo_or_bar)──▶ webhooks-server └─ 处理器 all ──(POST http://webhooks-server:8001/all)──▶ webhooks-server由于/foo_or_bar与/all的过滤条件与/foo存在交集Foo事件会命中全部三条规则一条Foo事件最终会触发三次 Webhook 请求这正是示例想要演示的一个消息扇出fan-out到多个 Webhook效果。源码拆解三个服务各司其职producer向 Kafka 发布带元数据的事件producer/main.go 中生产者通过kafka.NewPublisher创建发布器Brokers指向 docker-compose 网络中的kafka:9092pub, err : kafka.NewPublisher( kafka.PublisherConfig{ Brokers: brokers, // []string{kafka:9092} Marshaler: kafka.DefaultMarshaler{}, }, logger, )随后进入无限循环每秒随机发布一条Foo/Bar/Baz事件。关键点在于事件类型不是写在消息体里而是写进元数据eventTypes : []eventType{Foo, Bar, Baz} for { eventType : eventTypes[rand.Intn(3)] msg : message.NewMessage(watermill.NewUUID(), []byte(message)) msg.Metadata.Set(event_type, string(eventType)) fmt.Printf(%s Publishing %s\n\n, time.Now().String(), eventType) if err : pub.Publish(kafka_to_http_example, msg); err ! nil { panic(err) } time.Sleep(time.Second) }消息 ID 由watermill.NewUUID()生成消息体统一为字节串message主题固定为kafka_to_http_example。把事件类型放入 metadata 而非消息体可以让消费端无需反序列化消息体即可完成分流见下文 router 的注释。router消费 Kafka 并按元数据分发 HTTP Webhookrouter/main.go 是整套示例的核心。它创建了一个HTTP Publisherpublisher, err : watermill_http.NewPublisher(watermill_http.PublisherConfig{ MarshalMessageFunc: watermill_http.DefaultMarshalMessageFunc, }, logger)MarshalMessageFunc决定了消息如何被翻译成 HTTP 请求URL、方法、请求头、请求体DefaultMarshalMessageFunc的行为是向配置好的具体 URL 发送 POST 请求。关于该配置项的底层细节见下文HTTP Publisher 原理一节。Kafka 订阅者与 Router 的创建同样直白subscriber, err : kafka.NewSubscriber( kafka.SubscriberConfig{ Brokers: []string{kafka:9092}, Unmarshaler: kafka.DefaultMarshaler{}, }, logger, ) router, err : message.NewRouter(message.RouterConfig{}, logger)接着定义 Webhook 目标地址并注册三个处理器topic : kafka_to_http_example url : http://webhooks-server:8001/ router.AddHandler(foo, topic, subscriber, urlfoo, publisher, filterMessages(Foo)) router.AddHandler(foo_or_bar, topic, subscriber, urlfoo_or_bar, publisher, filterMessages(Foo, Bar)) router.AddHandler(all, topic, subscriber, urlall, publisher, filterMessages(Foo, Bar, Baz)) router.AddPlugin(plugin.SignalsHandler)router.AddHandler的参数依次为处理器名称、订阅主题、订阅者、发布目标这里是url路径、发布者、处理函数。每个处理器的处理函数由filterMessages生成// filterMessages passes the message along if its event type is one of acceptedTypes. func filterMessages(acceptedTypes ...string) message.HandlerFunc { return func(msg *message.Message) ([]*message.Message, error) { // the kafka producer sets this metadata so that we dont have to unmarshal the body // just sort the messages based on event type metadata msgEventType : msg.Metadata.Get(event_type) for _, typ : range acceptedTypes { if typ msgEventType { return message.Messages{msg}, nil } } return nil, nil } }filterMessages是一个返回处理函数的高阶函数当消息的event_type命中任一acceptedTypes时把消息原样返回Watermill 的 HandlerFunc 返回值会作为新的消息继续交给发布者发送未命中则返回nil, nil表示丢弃。正因为事件类型放在 metadata 中处理器无需解析消息体即可完成过滤。最终通过router.Run(context.Background())启动路由并注册了plugin.SignalsHandler见 message/router/plugin/signals.go使得按CtrlC发送中断信号时 Router 能优雅关闭。webhooks-serverWebhook 接收端webhooks-server/main.go 是标准的 Gonet/http服务监听:8001端口把所有路径的请求体读取后打印到 stdout并返回200 OKfunc handler(w http.ResponseWriter, r *http.Request) { body, err : ioutil.ReadAll(r.Body) if err ! nil { w.WriteHeader(http.StatusBadRequest) return } fmt.Printf( [%s] %s %s: %s\n\n, time.Now().String(), r.Method, r.URL.String(), string(body), ) w.WriteHeader(http.StatusOK) } func main() { http.HandleFunc(/, handler) http.ListenAndServe(:8001, http.DefaultServeMux) }在真实项目中这个角色通常由第三方服务的 HTTP 接口、或者网关/负载均衡器扮演。docker-compose 配置详解docker-compose.yml 定义了四个服务全部使用golang:1.25镜像kafka除外并把示例目录挂载进容器、在各子目录下直接go run main.go服务工作目录启动命令依赖webhooks-server/app/webhooks-server/go run main.go无router/app/router/go run main.gokafkaproducer/app/producer/go run main.gokafka、webhooks-server、routerkafka-redpanda start ...无几个值得注意的细节卷挂载.示例根目录挂载到容器/app同时将宿主机$GOPATH/pkg/mod挂载到/go/pkg/mod以复用 Go 模块缓存避免每次启动都重新下载依赖启动顺序depends_on保证router在kafka就绪后启动producer则在消息链路kafka、webhooks-server、router全部启动后再开始发布避免消息发到尚未就绪的订阅端Redpanda 配置kafka服务使用redpandadata/redpanda:v26.1.7以--mode dev-container开发模式运行通过--kafka-addr与--advertise-kafka-addr同时暴露容器内地址kafka:9092与宿主机地址localhost:19092--smp 1限制 CPU 核心数以节省资源--default-log-levelwarn抑制日志噪音logging.driver: none直接关闭了该服务的日志收集重启策略均设置为unless-stopped保证进程意外退出后自动拉起。示例中三个 Go 模块的依赖版本可从各自go.mod确认producer/go.mod 使用watermill v1.5.1watermill-kafka/v3 v3.1.2router/go.mod 额外引入watermill-http v1.1.4HTTP Publisher 的出处。运行与日志观测前置要求运行本示例需要安装 Docker 与 docker-compose官方安装指南见 README.md 引用的 Docker 文档仓库根目录 README.md 中也介绍了相关的依赖准备方式。启动全部服务在_examples/real-world-examples/sending-webhooks/目录下执行docker-compose up启动后producer会以每秒一条的频率发布事件router消费后按event_type分发 HTTP 请求webhooks-server则持续打印收到的请求。由于示例中的webhooks-server服务会不断向 stdout 输出日志你通常不需要额外干预即可看到完整的 Webhook 外呼链路。按服务过滤日志当多个服务同时输出日志时可以在另一个终端窗口中单独查看某个服务的输出# 只查看 router 的日志 docker-compose logs router # 带 -f 标志模拟 tail -f 行为持续跟踪输出 docker-compose logs -f router-f标志等价于tail -f会持续跟随输出。把{service}替换为producer、router、webhooks-server或kafka即可切换观察对象。通过对比producer发布的event_type、router转发的路径与webhooks-server实际收到的 POST 请求三份日志可以直观验证前面提到的扇出规则。HTTP Publisher 原理消息到 HTTP 请求的翻译本示例的灵魂是watermill_http.NewPublisher。官方文档 docs/content/pubsubs/http.md 对其机制有明确说明HTTP publisher 按照配置把消息翻译成 HTTP 请求并发送。消息的 topic 与 body 如何映射为 HTTP 请求的 URL、方法、请求头与请求体完全由MarshalMessageFunc决定watermill_http.DefaultMarshalMessageFunc会向构造时指定的具体 URL发送POST请求示例中正是通过router.AddHandler(..., urlfoo, publisher, ...)把目标 URL 逐处理器地传给发布者你也可以自定义MarshalMessageFunc按业务需要改写 URL、方法、请求头或载荷例如把消息元数据拼进 header、把 topic 映射到路径等HTTP 客户端默认使用 Go 的http.Client也可以传入自定义的http.Client以控制超时、重试与连接池行为。从该文档的Characteristics表可知HTTP Publisher 支持ExactlyOnceDelivery配合幂等键等机制与GuaranteedOrder消息顺序发送但不支持 ConsumerGroups也不具备持久化能力——它的职责是发出请求持久化保证依赖消息来源本例中的 Kafka。另外注意HTTP Publisher 所在的外部模块watermill-http并不在本仓库内本文所有结论均基于示例代码与仓库文档的公开使用方式。延伸思考这套模式如何落地到真实项目示例虽小却勾勒出一个可复用的事件外呼骨架实际落地时可从以下几处演进过滤与扇出filterMessages的高阶函数写法可直接复用——把是否转发的判定抽象成接受可变参数的白名单既保持了每个处理器独立可读又避免了重复代码目标地址管理示例把 URL 硬编码为http://webhooks-server:8001/真实场景应替换为外部 Webhook 地址并可考虑通过MarshalMessageFunc从消息元数据中动态解析目标 URL实现一条消息投递到不同第三方可靠性Kafka 端天然具备持久化与至少一次交付语义HTTP 外呼若需增强可靠性可叠加仓库内置的 retry 中间件、recoverer 中间件 或结合 requeuer 组件 处理失败消息反向场景如果需要收 Webhook → 写 Kafka可参考仓库的 receiving-webhooks 示例HTTP Subscriber 的用法见 docs/content/pubsubs/http.md 中StartHTTPServer()与-r.Running()的启动时序。小结sending-webhooks 示例完整演示了 Watermill 在Kafka 事件 → HTTP Webhook 外呼场景下的标准姿势用kafka包订阅事件用watermill-http包发布请求用Router.AddHandler把两者粘合并用 metadata 过滤函数实现多路径扇出。它既是 HTTP Publisher 的入门范本也是事件驱动系统中消息中间件与外部 HTTP 系统解耦这一常见需求的参考实现。配合docker-compose up一条命令即可在本地复现完整链路建议按上文日志观测一节实际运行一遍观察Foo事件如何同时触发/foo、/foo_or_bar、/all三条 Webhook。【免费下载链接】watermillBuilding event-driven applications the easy way in Go.项目地址: https://gitcode.com/GitHub_Trending/wa/watermill创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表