Watermill 入门指南:用 Go 以最简单的方式构建事件驱动应用

发布时间:2026/9/15 20:23:32

Watermill 入门指南:用 Go 以最简单的方式构建事件驱动应用 Watermill 入门指南用 Go 以最简单的方式构建事件驱动应用【免费下载链接】watermillBuilding event-driven applications the easy way in Go.项目地址: https://gitcode.com/GitHub_Trending/wa/watermillWatermill 是一个自带电池batteries included的 Go 消息处理库它把 Kafka、RabbitMQ、PostgreSQL、Google Cloud Pub/Sub 等异构 Pub/Sub 的复杂度统一隐藏在一组简洁接口背后。本文基于仓库中的 getting-started 文档 展开从底层Publisher/Subscriber接口讲到高层的Router组件并深入源码与官方示例读完你将能够独立完成从安装、发布订阅到路由处理、日志接入的完整 Watermill 开发流程。Watermill 是什么Watermill 是一个用 Go 构建以简单方式处理消息的库。你可以用它构建消息驱动message-driven和事件驱动event-driven应用底层对接 Kafka、RabbitMQ、PostgreSQL 等 Pub/Sub 系统而这些系统的接入由社区维护的独立包如watermill-kafka、watermill-amqp完成核心库本身不绑定任何具体消息队列。Watermill 自带电池它为每个消息驱动应用都会用到的能力消息模型、路由、中间件、插件、日志抽象等提供了开箱即用的工具而不是让你从零搭建。为什么使用 Watermill当你运行一个 HTTP 服务器时你并不直接操作 TCP 套接字、手动解析 HTTP 请求或管理连接——而是使用net/http这样的高层库由它替你处理所有复杂性。Watermill 之于消息正如net/http之于 HTTP。它为基于事件或其他异步模式构建应用提供了所需的一切。市面上存在大量消息队列各自拥有不同的特性、客户端库和 API。Watermill 将这些复杂性隐藏在一个易于使用和理解、且对所有消息队列统一的 API 之后——应用代码只面向抽象接口编程切换底层消息队列时业务代码几乎无需改动。需要特别强调Watermill 不是框架而是一个轻量级库可以非常容易地从项目中接入或移除。这一点可以从核心库的模块结构得到印证消息模型、路由、中间件、Pub/Sub 实现彼此解耦你的业务代码只依赖message.Publisher、message.Subscriber这样的接口而非某个具体实现。安装在 Go 项目中安装 Watermill 核心库go get -u github.com/ThreeDotsLabs/watermill如果还需要对接具体的消息队列则需额外安装对应适配包。例如仓库中的 Kafka 示例 main.go 使用了github.com/ThreeDotsLabs/watermill-kafka/v3/pkg/kafka和github.com/IBM/saramaAMQPRabbitMQ示例 main.go 使用了github.com/ThreeDotsLabs/watermill-amqp/v3/pkg/amqpNATS Streaming 示例 main.go 使用了github.com/ThreeDotsLabs/watermill-nats/pkg/nats。建议优先阅读官方示例目录 _examples 中各 Pub/Sub 子项目的go.mod与go.sum以确认与当前核心库版本兼容的适配包版本。一分钟背景事件驱动的基本模型事件驱动应用背后的思想始终如一一部分发布消息publish另一部分订阅消息subscribe。Watermill 为多种 发布者与订阅者 实现支持这一行为。三层 APIWatermill 提供了三套处理消息的 API它们层层叠加每一层都在上一层之上提供更高层的抽象自底向上依次是Publisher Subscriber最底层、最基础的消息收发接口对应 pub-sub.md 文档Router高层路由组件自动处理订阅、并发、Ack/Nack、优雅关闭等对应 messages-router.md 文档CQRS面向 Command/Event 的通用高层 API对应 cqrs.md 文档。本文将自底向上展开。即使你打算直接使用高层 API理解底层原理也很有价值——例如 Router 内部的 Ack/Nack 机制正是建立在底层消息模型的语义之上。第一层Publisher 与 Subscriber大多数 Pub/Sub 库都带有复杂的特性。Watermill 将这一复杂性隐藏在两个接口背后定义于 message/pubsub.gotype Publisher interface { Publish(topic string, messages ...*Message) error Close() error } type Subscriber interface { Subscribe(ctx context.Context, topic string) (-chan *Message, error) Close() error }从接口定义可以看出几个关键契约Publish可同步也可异步取决于具体实现Publish不接收 Context而是使用每条消息自身的 Context见Message.Context()Publish必须保证线程安全Subscribe返回一个接收消息的 channel该 channel 在Close()后关闭要接收下一条消息必须先对收到的消息调用Ack()如果处理失败并希望消息被重新投递应调用Nack()替代源码注释见 message/pubsub.go当传入的ctx被取消时订阅者会关闭订阅与输出 channel。创建消息Watermill 的核心是 Message 结构体——它之于 Watermill正如http.Request之于net/http包。大多数 Watermill 特性都基于这个结构体工作。Watermill 不强制任何消息格式。NewMessage期望一个字节切片作为 payload。你可以使用字符串、JSON、protobuf、Avro、gob或任何能序列化为[]byte的格式。消息 UUID 是可选的但强烈建议提供便于调试。msg : message.NewMessage(watermill.NewUUID(), []byte(Hello, world!))在源码 message/message.go 中可以看到Message的完整结构除UUID、Metadata类似 HTTP 请求头随消息一起持久化到 Pub/Sub、Payload外还内置了ack/noAck两个 channel 及Ack()/Nack()方法。Ack()与Nack()都是非阻塞且幂等的如果先调用了Nack()再调用Ack()会返回false反之亦然见 message/message.go。你还可以通过Acked()/Nacked()channel 在select中等待确认信号。发布消息Publish期望一个 topic 和一个或多个Messageerr : publisher.Publish(example.topic, msg) if err ! nil { panic(err) }仓库为每种受支持的 Pub/Sub 提供了可运行的最小示例均位于 _examples/pubsubs 目录。以 Go Channel内存版 Pub/Sub无外部依赖为例main.go 中完整的发布逻辑如下func publishMessages(publisher message.Publisher) { for { msg : message.NewMessage(watermill.NewUUID(), []byte(Hello, world!)) if err : publisher.Publish(example.topic, msg); err ! nil { panic(err) } time.Sleep(time.Second) } }其余示例的发布代码几乎一模一样差异只在发布者的构造方式上Kafka通过kafka.NewPublisher(kafka.PublisherConfig{Brokers: ..., Marshaler: kafka.DefaultMarshaler{}}, logger)创建见 _examples/pubsubs/kafka/main.goNATS Streaming通过nats.NewStreamingPublisher创建需配置ClusterID、ClientID与StanOptions见 _examples/pubsubs/nats-streaming/main.goGoogle Cloud Pub/Sub见 _examples/pubsubs/googlecloud/main.goRabbitMQ (AMQP)通过amqp.NewPublisher(amqp.NewDurableQueueConfig(amqpURI), logger)创建见 _examples/pubsubs/amqp/main.goSQL见 _examples/pubsubs/sql/main.goAWS SQS / SNS分别见 _examples/pubsubs/aws-sqs/main.go 与 _examples/pubsubs/aws-sns/main.go。订阅消息Subscribe期望一个 topic 名并返回接收消息的 channel。topic 的确切含义取决于 Pub/Sub 实现通常它需要与发布者使用的 topic 名保持一致。消息处理完成后必须调用Ack()确认否则消息会被重新投递在 go-channel/main.go 与 kafka/main.go 的注释中都强调了这一点。messages, err : subscriber.Subscribe(ctx, example.topic) if err ! nil { panic(err) } for msg : range messages { fmt.Printf(received message: %s, payload: %s\n, msg.UUID, string(msg.Payload)) msg.Ack() }典型的订阅处理函数封装为独立的process函数func process(messages -chan *message.Message) { for msg : range messages { log.Printf(received message: %s, payload: %s, msg.UUID, string(msg.Payload)) // we need to Acknowledge that we received and processed the message, // otherwise, it will be resent over and over again. msg.Ack() } }用 Docker 本地运行 Kafka 示例仓库为所有外部依赖型示例提供了docker-compose.yml。以 Kafka 示例的 docker-compose.yml 为例services: server: image: golang:1.25 restart: unless-stopped depends_on: - kafka volumes: - .:/app - $GOPATH/pkg/mod:/go/pkg/mod working_dir: /app command: go run main.go kafka: image: redpandadata/redpanda:v26.1.7 restart: unless-stopped logging: driver: none command: - redpanda - start - --kafka-addr internal://0.0.0.0:9092,external://0.0.0.0:19092 - --advertise-kafka-addr internal://kafka:9092,external://localhost:19092 - --mode dev-container - --smp 1 - --default-log-levelwarn将示例源码放到main.go后执行docker-compose up即可一键启动 Go 编译环境与 KafkaRedpanda服务。NATS Streaming、Google Cloud Pub/Sub官方模拟器、AMQP、SQL、AWS SQS/SNS 等示例的docker-compose.yml同样位于各自的示例目录中。第二层RouterPublisher 与 Subscriber 是 Watermill 的底层组件。对大多数场景你更想要的是高层 API——Router详见 messages-router.md。它自动处理消息订阅、并发分发、Ack/Nack 与优雅关闭。Router 的关键实现细节在 message/router.go 中HandlerFunc是消息到达时被调用的函数msg.Ack()在HandlerFunc不返回错误时被自动调用返回错误时则自动调用msg.Nack()见 message/router.goHandlerMiddleware则允许以装饰器模式包裹HandlerFunc在处理器前后执行逻辑见 message/router.go。配置 Router首先创建 Router 并添加插件与中间件router, err : message.NewRouter(message.RouterConfig{}, logger) if err ! nil { panic(err) } // SignalsHandler will gracefully shutdown Router when SIGTERM is received. // You can also close the router by just calling r.Close(). router.AddPlugin(plugin.SignalsHandler) // Router level middleware are executed for every message sent to the router router.AddMiddleware( // CorrelationID will copy the correlation id from the incoming messages metadata to the produced messages middleware.CorrelationID, // The handler function is retried if it returns an error. // After MaxRetries, the message is Nacked and its up to the PubSub to resend it. middleware.Retry{ MaxRetries: 3, InitialInterval: time.Millisecond * 100, Logger: logger, }.Middleware, // Recoverer handles panics from handlers. // In this case, it passes them as errors to the Retry middleware. middleware.Recoverer, )中间件middleware是作用于每条进入 Router 的消息的函数。你可以直接使用现成的中间件如 correlation关联 ID 透传、metrics指标、poison queue死信队列、retrying重试、throttling限流等完整清单见 messages-router.md 以及 message/router/middleware 目录circuit_breaker.go、deduplicator.go、delay_on_error.go、instant_ack.go、recoverer.go、timeout.go、ignore_errors.go等也可以编写自己的中间件。插件plugin在 Router 启动时执行。plugin.SignalsHandler会在收到 SIGTERM 信号时优雅关闭 Router。RouterConfig目前只有一个配置项CloseTimeoutRouter 关闭时等待 handler 完成的最长时间默认值为 30 秒见 message/router.go。处理器Handlers接下来为 Router 注册处理器。每个 handler 独立处理收到的消息handler 从给定的 subscriber 和 topic 读取消息handler 函数返回的任何消息都会被发布到给定的 publisher 和 topic。// AddHandler returns a handler which can be used to add handler level middleware // or to stop handler. handler : router.AddHandler( struct_handler, // handler name, must be unique incoming_messages_topic, // topic from which we will read events pubSub, outgoing_messages_topic, // topic to which we will publish events pubSub, structHandler{}.Handler, )注意上面的示例对 subscriber 和 publisher 使用了同一个pubSub参数因为我们使用的是GoChannel实现——一个简单的内存版 Pub/Sub。如果 handler 内部不打算发布消息可以使用更简单的AddConsumerHandler// just for debug, we are printing all messages received on incoming_messages_topic router.AddConsumerHandler( print_incoming_messages, incoming_messages_topic, pubSub, printMessages, )你可以使用两种类型的handler 函数无依赖的函数func(msg *message.Message) ([]*message.Message, error)结构体方法func (c structHandler) Handler(msg *message.Message) ([]*message.Message, error)如果你要写的 handler 没有任何依赖用第一种即可当 handler 需要数据库句柄、logger 等依赖时第二种更合适。例如 3-router/main.go 中的示例func printMessages(msg *message.Message) error { fmt.Printf( \n Received message: %s\n %s\n metadata: %v\n\n, msg.UUID, string(msg.Payload), msg.Metadata, ) return nil } type structHandler struct { // we can add some dependencies here } func (s structHandler) Handler(msg *message.Message) ([]*message.Message, error) { log.Println(structHandler received message, msg.UUID) msg message.NewMessage(watermill.NewUUID(), []byte(message produced by structHandler)) return message.Messages{msg}, nil }此外Router 还支持 handler 级中间件——只对特定 handler 生效添加方式与 router 级中间件相同见 3-router/main.go 中handler.AddMiddleware(...)的用法。最后运行 Router。Run在 Router 运行期间是阻塞的// Now that all handlers are registered, were running the Router. // Run is blocking while the router is running. ctx : context.Background() if err : router.Run(ctx); err ! nil { panic(err) }完整的 Router 示例源码见 _examples/basic/3-router/main.go。另一个展示 Router 与真实消息队列组合的入口是 _examples/basic/1-your-first-app/main.go它使用 Kafka Publisher/Subscriber注册了一个消费eventstopic、反序列化 JSON 事件、处理后发布到events-processedtopic 的 handler并演示了 handler 返回错误时默认触发 Nack、消息将被重新处理的语义可通过Retry、PoisonQueue等中间件改变该行为。日志要看到 Watermill 的日志只需传入任何实现了LoggerAdapter接口的 logger。该接口定义于 log.go包含Error、Info、Debug、Trace、With五个方法type LoggerAdapter interface { Error(msg string, err error, fields LogFields) Info(msg string, fields LogFields) Debug(msg string, fields LogFields) Trace(msg string, fields LogFields) With(fields LogFields) LoggerAdapter }Watermill 自带几种现成实现NewStdLogger(debug, trace bool)适用于实验性开发将日志输出到 stderr。两个布尔参数控制是否启用 Debug 与 Trace 级别的输出见 log.goNewSlogLogger(logger *slog.Logger)标准库log/slog的即用适配器logger传nil时自动使用slog.Default()NewSlogLoggerWithLevelMapping(logger, watermillLevelToSlog)可额外提供一个映射表把 Watermill 的日志级别映射到slog级别——例如将 Watermill 的 Info 日志降级为 slog 的 Debug见 slog.go。另外还有NopLogger丢弃所有日志与NewStdLoggerWithOut指定输出 io.Writer都在 log.go 中。下一步学什么查看 CQRS 组件了解第三层通用高层 API查看 文档主题 获取更多细节Outbox 模式 是事件驱动应用中需要掌握的关键模式对应组件见 components/forwarder参考 _examples 目录下的示例看 Watermill 在实践中的工作方式。推荐的示例入口仓库的 _examples 目录展示了如何上手使用 Watermill推荐起点Your first Watermill application——其docker-compose.yml中包含了包括 Go 和 Kafka 在内的完整环境一条命令即可运行接着看 Realtime feed——使用了更多中间件并包含两个 handler想了解不同的订阅者实现HTTP可以看 receiving-webhooks 示例——一个将 webhook 保存到 Kafka 的直白应用完整示例清单见 READMECQRS、Outbox 模式、SSE 等真实世界示例也在此列出。支持如果任何地方不清楚欢迎通过 support 页面 所列的渠道获取帮助。【免费下载链接】watermillBuilding event-driven applications the easy way in Go.项目地址: https://gitcode.com/GitHub_Trending/wa/watermill创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
延伸阅读

更多相关文章

2026/9/15 20:18:32

支付宝小程序后端认证:手写RSA2签名绕开pycrypto与SDK坑

做支付宝小程序后端的时候,我第一个周末就栽在两个老熟人手上:alipay-sdk-python 和 pycrypto。先说结果,SDK 是从 PyPI 直接拉下来的,pycrypto 装不上,编译错误刷了一整屏,后来我索性把用户认证流程改成自…

2026/9/15 20:18:32

TFT多变量时序预测实战:原理、PyTorch实现与经验

做过多变量时序预测的朋友,应该都经历过这样的阶段:拿到一堆特征,不管三七二十一先上个LSTM再说。训练半天,loss降了,结果一上测试集,要么滞后严重,要么变量稍微多一点就直接崩溃。我也一样&…

2026/9/15 20:58:35

家居投资集团跨域管理:五维解法打通扩张困局

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/15 20:58:35

Prototypical Networks小样本学习PyTorch实战

1. 这不是又一个“花里胡哨”的小众模型——Prototypical Networks 是少有的、真正把“人类学习逻辑”刻进代码里的方法你有没有想过,为什么人能一眼认出“没见过的猫”?比如第一次看到一只苏格兰折耳猫,你不会犹豫,直接说“这是猫…

2026/9/15 20:53:34

你的问卷没被退回来,不是因为好,是因为没人认真看

毕夏AI官网 www.bixiaai.com 毕夏AI写作官网 www.bixiaai.com 毕夏官网 www.bixiaai.com 毕夏智能写作官网 www.bixiaai.com 说一个反直觉的判断:本科毕业论文问卷设计的真实淘汰线,不是“好不好”,是“有没有硬伤”。 抽检评议要素写…

2026/9/15 4:54:30

拯救者Y7000黑屏故障排查与维修实战指南

1. 项目概述:一台黑屏的拯救者Y7000,到底卡在哪一步? 联想拯救者Y7000系列笔记本,从2018年第一代搭载i5-8300H开始,到后来的i7-9750H、i7-10750H、i5-11400H,再到2023年款的R7-7840HS,它始终是学…

2026/9/15 0:01:16

AI英语单词APP开发:自适应学习算法与移动端优化实践

1. 项目概述 作为一名在移动应用开发领域摸爬滚打多年的老手,我最近完成了一个AI英语单词APP的开发项目。这个项目将传统单词记忆方法与现代AI技术相结合,打造了一款能够智能适应不同用户学习习惯的英语学习工具。 市面上大多数单词APP都存在一个通病&a…

2026/9/15 0:01:16

Flutter与OpenHarmony结合开发手语学习APP实战

1. 项目背景与核心价值作为一名同时接触过Flutter和OpenHarmony的开发者,最近我完成了一个基于Flutter for OpenHarmony的手语学习APP实战项目。这个项目最大的特点在于实现了跨平台框架与国产操作系统深度结合的创新实践——用Flutter开发的应用能完美运行在OpenHa…

2026/9/15 0:01:16

六个月成为机器人工程师:从ROS2到SLAM的实战路径

1. 六个月的紧迫感从哪来:先搞清楚你要成为哪种机器人工程师说实话,六个月的期限并不是一个宽松的时间线。市面上任何一本正经的机器人学教材都超过五百页,ROS2的官方文档可以翻到你怀疑人生,再加上ABB、KUKA这些工业机器人厂家动…

2026/9/15 14:22:53

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行,Type-C接口算是典型的“看着简单,做起来全坑”的东西。光引脚就24个,高低速信号、电源、控制线全部塞在一个小小的连接器里,如果PCB布局不做规划,打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/14 13:53:59

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/15 11:42:23

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
咨询二维码