尧图精选

构建你的第一个 Watermill 应用:用 Go 与 Kafka 实现消息消费、转换与再发布

🕒 发布时间:2026/9/15 23:46:16 📁 来源:尧图网络
构建你的第一个 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 应用容器与 KafkaRedpanda容器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.Publisher、message.Subscriber、Router 等与具体消息中间件的对接全部通过独立的 pub/sub 适配包完成。环境要求运行该示例需要本机安装 Docker 与 docker-composeKafka 与整个编译运行环境都由容器提供无需在宿主机安装 Go 工具链。Docker Compose 中定义了两个服务见 docker-compose.ymlserver基于golang:1.26镜像将当前目录挂载到/app启动时先执行go install github.com/ThreeDotsLabs/watermill/tools/milllatest安装命令行工具mill再执行go run main.gokafka使用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可简写为-tdocker-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:9092logger通过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 与 Subscriberfunc 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 并注册 Handlerrouter, 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 处理这么久。接着挂载了两类扩展点插件Pluginplugin.SignalsHandler在 Router 启动时注册对SIGINT/SIGTERM的监听见 message/router/plugin/signals.go收到信号后调用router.Close()优雅关闭保证 CtrlC 退出时正在处理的消息能正常收尾中间件Middlewaremiddleware.Recoverer包装 HandlerFunc把 handler 中抛出的 panic 捕获并转换为RecoveredPanicError附带完整堆栈返回见 message/router/middleware/recoverer.go避免单个消息的处理崩溃拖垮整个进程。随后是示例最重要的部分——注册业务 Handlerrouter.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...)一行即可完成消息发布。启动 Routerif 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/milllatest再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),仅供参考
上一篇/下一篇内容由系统自动关联 返回资讯列表 →