构建你的第一个 Watermill 应用:用 Go 与 Kafka 实现消息消费、转换与再发布
【免费下载链接】watermillBuilding event-driven applications the easy way in Go.项目地址: https://gitcode.com/GitHub_Trending/wa/watermill
本指南基于仓库中的_examples/basic/1-your-first-app示例,带你完整走一遍 Watermill 的“第一个应用”:程序在一个循环中从 Kafka 的eventstopic 消费事件,经 Handler 处理后发布到events-processedtopic。读完本文,你将掌握 Watermill Router 的核心用法——如何创建 Publisher/Subscriber、注册 Handler、挂载中间件与插件,并能在本地用 Docker Compose 一键启动整套环境,通过内置mill命令行工具直接观测 Kafka 中的消息流转。
示例项目概览
这个示例演示的是 Watermill 最经典的一条数据链路:
Kafka topic: events ──▶ Router Handler ──▶ Kafka topic: events-processed (原始事件) (转换处理) (处理结果)应用启动后,后台会以每秒一条的速度向eventstopic 生产模拟事件;Router 中的 Handler 订阅该 topic,将事件反序列化、打印、加上处理时间戳后重新序列化,并发布到events-processedtopic。
整个示例只有四个文件,结构非常精简:
| 文件 | 作用 |
|---|---|
| main.go | 示例源码,整个示例的核心,本文重点剖析的对象 |
| docker-compose.yml | 本地环境编排,包含 Go 应用容器与 Kafka(Redpanda)容器 |
| go.mod | Go modules 依赖声明 |
| go.sum | Go modules 校验和文件 |
依赖方面,go.mod 声明了两个直接依赖:核心库github.com/ThreeDotsLabs/watermill v1.5.1与 Kafka 适配器github.com/ThreeDotsLabs/watermill-kafka/v3 v3.1.2(底层基于IBM/sarama)。也就是说,Watermill 本身只提供消息抽象(message.Publisher、message.Subscriber、Router 等),与具体消息中间件的对接全部通过独立的 pub/sub 适配包完成。
环境要求
运行该示例需要本机安装 Docker 与 docker-compose,Kafka 与整个编译运行环境都由容器提供,无需在宿主机安装 Go 工具链。
Docker Compose 中定义了两个服务(见 docker-compose.yml):
- server:基于
golang:1.26镜像,将当前目录挂载到/app,启动时先执行go install github.com/ThreeDotsLabs/watermill/tools/mill@latest安装命令行工具mill,再执行go run main.go; - kafka:使用
redpandadata/redpanda:v26.1.7镜像以开发模式启动,对外暴露9092(容器网络内)与19092(宿主机映射)两个 Kafka 协议端口。Redpanda 与 Kafka 协议完全兼容,因此代码与命令都无需任何改动。
运行示例
1. 启动完整环境
在示例目录下执行:
docker-compose up启动后会看到 server 容器持续打印received event {...}日志,表明 Handler 正在逐条消费eventstopic 中的消息:
> docker-compose up [some initial logs] server_1 | 2019/08/29 19:41:23 received event {ID:0} server_1 | 2019/08/29 19:41:23 received event {ID:1} server_1 | 2019/08/29 19:41:23 received event {ID:2} server_1 | 2019/08/29 19:41:23 received event {ID:3} server_1 | 2019/08/29 19:41:24 received event {ID:4} server_1 | 2019/08/29 19:41:25 received event {ID:5} server_1 | 2019/08/29 19:41:26 received event {ID:6} server_1 | 2019/08/29 19:41:27 received event {ID:7} server_1 | 2019/08/29 19:41:28 received event {ID:8} server_1 | 2019/08/29 19:41:29 received event {ID:9}2. 观察 Kafka 中的原始事件
打开另一个终端,使用mill命令直接查看 Kafka 的eventstopic。mill是仓库自带的多中间件命令行工具(源码位于 tools/mill),这里通过docker-compose exec在 server 容器内调用它的 kafka 子命令:
docker-compose exec server mill kafka consume -b kafka:9092 --topic events输出如下:
{"id":12} {"id":13} {"id":14} {"id":15} {"id":16} {"id":17}注意这里展示的是命令启动之后新到达的消息(ID 12 开始)。这与mill的默认读取偏移有关——查看 tools/mill/cmd/kafka.go 可知,只有显式传入--from-beginning标志时才会将偏移设为OffsetOldest(等价于auto.offset.reset: earliest),默认从最新偏移开始消费。而 ID 0~11 的事件已被 consumer grouphandler_1消费确认。
3. 查看处理后的消息
再观察events-processedtopic(--topic可简写为-t):
docker-compose exec server mill kafka consume -b kafka:9092 -t events-processed输出的是 Handler 转换后的消息,每条都包含了原事件 ID 与处理时间戳:
{"processed_id":21,"time":"2019-08-29T19:42:31.4464598Z"} {"processed_id":22,"time":"2019-08-29T19:42:32.4501767Z"} {"processed_id":23,"time":"2019-08-29T19:42:33.4530692Z"} {"processed_id":24,"time":"2019-08-29T19:42:34.4561694Z"} {"processed_id":25,"time":"2019-08-29T19:42:35.4608918Z"}对比两个 topic 的输出即可直观看到:消息被完整地“消费 → 转换 → 再发布”了一遍。
源码剖析:main.go 逐段解读
main.go 是示例的核心,约 150 行,我们按逻辑顺序拆解。
全局配置与消息结构体
var ( brokers = []string{"kafka:9092"} consumeTopic = "events" publishTopic = "events-processed" logger = watermill.NewStdLogger( true, // debug false, // trace ) marshaler = kafka.DefaultMarshaler{} )brokers指向容器网络内的 Kafka 地址kafka:9092;logger通过watermill.NewStdLogger创建标准日志适配器,两个布尔参数分别控制是否输出 debug 与 trace 级别日志。其实现见 log.go:底层包装标准库log.Logger,未开启的级别对应的 Logger 为 nil,日志会被直接丢弃;marshaler使用kafka.DefaultMarshaler,同时承担消息体序列化(Marshal)与反序列化(Unmarshal)职责。
两个 JSON 结构体分别对应输入与输出消息的格式:
type event struct { ID int `json:"id"` } type processedEvent struct { ProcessedID int `json:"processed_id"` Time time.Time `json:"time"` }创建 Kafka Publisher 与 Subscriber
func createPublisher() message.Publisher { kafkaPublisher, err := kafka.NewPublisher( kafka.PublisherConfig{ Brokers: brokers, Marshaler: marshaler, }, logger, ) if err != nil { panic(err) } return kafkaPublisher } func createSubscriber(consumerGroup string) message.Subscriber { kafkaSubscriber, err := kafka.NewSubscriber( kafka.SubscriberConfig{ Brokers: brokers, Unmarshaler: marshaler, ConsumerGroup: consumerGroup, // every handler will use a separate consumer group }, logger, ) if err != nil { panic(err) } return kafkaSubscriber }两个工厂函数的返回值类型都是接口message.Publisher/message.Subscriber——这正是 Watermill 的关键抽象:业务代码只依赖接口,底层是 Kafka、RabbitMQ 还是 Go channel 都不影响上层逻辑。createSubscriber接收一个consumerGroup参数,示例以"handler_1"作为消费组名,注释明确说明“每个 handler 应使用独立的 consumer group”,这样不同的 handler 可以各自维护消费进度而互不干扰。
组装 Router 并注册 Handler
router, err := message.NewRouter(message.RouterConfig{}, logger) if err != nil { panic(err) } router.AddPlugin(plugin.SignalsHandler) router.AddMiddleware(middleware.Recoverer)Router 是 Watermill 的调度核心(实现见 message/router.go)。NewRouter接收一个RouterConfig,示例传了空配置——由setDefaults兜底:当CloseTimeout为 0 时默认设为 30 秒,即优雅关闭时最多等待 handler 处理这么久。
接着挂载了两类扩展点:
- 插件(Plugin):
plugin.SignalsHandler在 Router 启动时注册对SIGINT/SIGTERM的监听(见 message/router/plugin/signals.go),收到信号后调用router.Close()优雅关闭,保证 Ctrl+C 退出时正在处理的消息能正常收尾; - 中间件(Middleware):
middleware.Recoverer包装 HandlerFunc,把 handler 中抛出的 panic 捕获并转换为RecoveredPanicError(附带完整堆栈)返回(见 message/router/middleware/recoverer.go),避免单个消息的处理崩溃拖垮整个进程。
随后是示例最重要的部分——注册业务 Handler:
router.AddHandler( "handler_1", // handler name, must be unique consumeTopic, // topic from which messages should be consumed subscriber, publishTopic, // topic to which messages should be published publisher, func(msg *message.Message) ([]*message.Message, error) { consumedPayload := event{} err := json.Unmarshal(msg.Payload, &consumedPayload) if err != nil { return nil, err } fmt.Printf("received event %+v\n", consumedPayload) newPayload, err := json.Marshal(processedEvent{ ProcessedID: consumedPayload.ID, Time: time.Now(), }) if err != nil { return nil, err } newMessage := message.NewMessage(watermill.NewUUID(), newPayload) return []*message.Message{newMessage}, nil }, )AddHandler有六个参数:handler 名称(必须全局唯一)、订阅 topic、订阅者、发布 topic、发布者、处理函数。关于 HandlerFunc 的语义,源码注释(message/router.go 中HandlerFunc定义处)说得很清楚:
- handler 返回 nil 错误时,Router 自动对消息调用
Ack(); - handler 返回错误时,Router 自动调用
Nack(),消息会被重新投递处理; - 若 handler 内部已手动
Ack(),即使再返回错误也不会重复 Nack。
处理函数内做了三件事:json.Unmarshal解析输入;fmt.Printf打印;构造processedEvent后通过message.NewMessage(watermill.NewUUID(), newPayload)创建新消息并返回。这里watermill.NewUUID()为消息生成调试用的 UUID(见 uuid.go)。需要注意的是,处理函数本身不直接调用 publisher 发布——它只需把要发布的消息作为返回值返回,Router 会自动把它们发布到publishTopic。这一点从 message/router.go 中handleMessage的实现可以印证:先执行handler(msg)拿到producedMessages,全部发布成功后才对消费的消息Ack();若发布失败同样走Nack(),保证消息不会被丢。
模拟事件生产者
go simulateEvents(publisher) func simulateEvents(publisher message.Publisher) { i := 0 for { e := event{ID: i} payload, err := json.Marshal(e) if err != nil { panic(err) } err = publisher.Publish(consumeTopic, message.NewMessage( watermill.NewUUID(), payload, )) if err != nil { panic(err) } i++ time.Sleep(time.Second) } }simulateEvents以独立 goroutine 运行,每秒向eventstopic 发布一条{"id":N}事件。这演示了message.Publisher接口的最小用法:Publish(topic, msgs...)一行即可完成消息发布。
启动 Router
if err := router.Run(context.Background()); err != nil { panic(err) }Run是一个阻塞调用(见 message/router.go):先执行所有插件,再启动全部已注册 handler(订阅 topic、拉起消费循环),直到所有 handler 停止或Close()被调用才会返回。
深入原理:Router 的自动确认机制
AddHandler之所以能省去手动 ack/nack,是因为 Router 在handleMessage中封装了完整的消息生命周期(message/router.go 中handleMessage方法):
- 捕获 handler 中的 panic 并转为
Nack(); - 调用 handler 拿到产物消息与错误;
- 出错 →
msg.Nack(),消息进入重投递; - 成功 → 将产物消息发布到
publishTopic,发布失败同样Nack(); - 全部成功 →
msg.Ack(),Kafka 消费组 offset 前移。
Message的 Ack/Nack 本身是幂等且非阻塞的(message/message.go),通过内部 channel 的关闭状态记录确认结果。这套机制保证了“消息至少被处理一次”的语义:处理失败的消息绝不会被静默吞掉。
如果想改变出错时的默认行为(如重试若干次、进入死信队列),可以像示例挂载Recoverer一样添加更多中间件,例如middleware.Retry、middleware.PoisonQueue、middleware.Timeout等,它们都位于 message/router/middleware 目录;也可以按HandlerMiddleware的签名实现自定义中间件。
docker-compose 环境说明
示例的 docker-compose.yml 值得留意几个细节:
server服务通过volumes把示例目录挂载进容器,command中先go install .../tools/mill@latest再go run main.go,因此docker-compose up后即可在容器内直接使用mill命令;kafka服务选用 Redpanda 而非原生 Kafka,是为了降低资源占用、加快启动;--advertise-kafka-addr分别配置了容器内地址(internal://kafka:9092)与宿主机地址(external://localhost:19092),供不同场景连接;- 如需在宿主机直接运行程序,只需把
main.go中的brokers改为localhost:19092即可。
下一步
本示例展示了 Watermill 的最小闭环:Publisher → Topic → Subscriber → Router Handler → Publisher。你可以在此基础上继续深入:
- 阅读 docs/learn/getting-started.md 了解 Router、消息与中间件的完整背景知识;
- 对照 message/router.go 与 message/message.go 研读 Router 生命周期与消息确认机制的源码细节;
- 探索 docs/pubsubs 目录,将同一套 Publisher/Subscriber 抽象替换为 Go channel、NATS、RabbitMQ 等其他中间件,体会 Watermill “更换中间件不改业务代码”的设计;
- 查看 docs/advanced 目录下的高级主题,如延迟消息、消息转发(Forwarder)、指标监控与错误重试等。
【免费下载链接】watermillBuilding event-driven applications the easy way in Go.项目地址: https://gitcode.com/GitHub_Trending/wa/watermill
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考