用 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 外呼)。示例由 producer(Kafka 生产者)、router(事件路由与转换)、webhooks-server(Webhook 接收端)与 Redpanda(Kafka 兼容消息后端)四个服务构成,读完本文你将掌握 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.md):
producer:向 Kafka 持续发布消息,消息按随机顺序取Foo、Bar、Baz三种类型之一,事件类型写入元数据键event_type;webhooks_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: topic=kafka_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:9092:
pub, 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 Webhook
router/main.go 是整套示例的核心。它创建了一个HTTP Publisher:
publisher, 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, url+"foo", publisher, filterMessages("Foo")) router.AddHandler("foo_or_bar", topic, subscriber, url+"foo_or_bar", publisher, filterMessages("Foo", "Bar")) router.AddHandler("all", topic, subscriber, url+"all", 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 don't 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),使得按Ctrl+C发送中断信号时 Router 能优雅关闭。
webhooks-server:Webhook 接收端
webhooks-server/main.go 是标准的 Gonet/http服务,监听:8001端口,把所有路径的请求体读取后打印到 stdout,并返回200 OK:
func 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.go | kafka |
producer | /app/producer/ | go run main.go | kafka、webhooks-server、router |
kafka | - | 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-level=warn抑制日志噪音(logging.driver: none直接关闭了该服务的日志收集); - 重启策略:均设置为
unless-stopped,保证进程意外退出后自动拉起。
示例中三个 Go 模块的依赖版本可从各自go.mod确认:producer/go.mod 使用watermill v1.5.1+watermill-kafka/v3 v3.1.2;router/go.mod 额外引入watermill-http v1.1.4(HTTP 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(..., url+"foo", 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),仅供参考