构建你的第一个 Watermill 应用:用 Go 与 Kafka 实现消息消费、转换与再发布
2026/9/15 23:45:47 网站建设 项目流程

构建你的第一个 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.modGo modules 依赖声明
go.sumGo modules 校验和文件

依赖方面,go.mod 声明了两个直接依赖:核心库github.com/ThreeDotsLabs/watermill v1.5.1与 Kafka 适配器github.com/ThreeDotsLabs/watermill-kafka/v3 v3.1.2(底层基于IBM/sarama)。也就是说,Watermill 本身只提供消息抽象(message.Publishermessage.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方法):

  1. 捕获 handler 中的 panic 并转为Nack()
  2. 调用 handler 拿到产物消息与错误;
  3. 出错 →msg.Nack(),消息进入重投递;
  4. 成功 → 将产物消息发布到publishTopic,发布失败同样Nack()
  5. 全部成功 →msg.Ack(),Kafka 消费组 offset 前移。

Message的 Ack/Nack 本身是幂等且非阻塞的(message/message.go),通过内部 channel 的关闭状态记录确认结果。这套机制保证了“消息至少被处理一次”的语义:处理失败的消息绝不会被静默吞掉。

如果想改变出错时的默认行为(如重试若干次、进入死信队列),可以像示例挂载Recoverer一样添加更多中间件,例如middleware.Retrymiddleware.PoisonQueuemiddleware.Timeout等,它们都位于 message/router/middleware 目录;也可以按HandlerMiddleware的签名实现自定义中间件。

docker-compose 环境说明

示例的 docker-compose.yml 值得留意几个细节:

  • server服务通过volumes把示例目录挂载进容器,command中先go install .../tools/mill@latestgo 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询