Watermill 入门指南:用 Go 以最简单的方式构建事件驱动应用
【免费下载链接】watermillBuilding event-driven applications the easy way in Go.项目地址: https://gitcode.com/GitHub_Trending/wa/watermill
Watermill 是一个"自带电池"(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/sarama,AMQP(RabbitMQ)示例 main.go 使用了github.com/ThreeDotsLabs/watermill-amqp/v3/pkg/amqp,NATS Streaming 示例 main.go 使用了github.com/ThreeDotsLabs/watermill-nats/pkg/nats。建议优先阅读官方示例目录 _examples 中各 Pub/Sub 子项目的go.mod与go.sum,以确认与当前核心库版本兼容的适配包版本。
一分钟背景:事件驱动的基本模型
事件驱动应用背后的思想始终如一:一部分发布消息(publish),另一部分订阅消息(subscribe)。Watermill 为多种 发布者与订阅者 实现支持这一行为。
三层 API
Watermill 提供了三套处理消息的 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.go:
type 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 和一个或多个Message:
err := 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.go; - NATS Streaming:通过
nats.NewStreamingPublisher创建,需配置ClusterID、ClientID与StanOptions,见 _examples/pubsubs/nats-streaming/main.go; - Google Cloud Pub/Sub:见 _examples/pubsubs/googlecloud/main.go;
- RabbitMQ (AMQP):通过
amqp.NewPublisher(amqp.NewDurableQueueConfig(amqpURI), logger)创建,见 _examples/pubsubs/amqp/main.go; - SQL:见 _examples/pubsubs/sql/main.go;
- AWS 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-level=warn将示例源码放到main.go后,执行docker-compose up即可一键启动 Go 编译环境与 Kafka(Redpanda)服务。NATS Streaming、Google Cloud Pub/Sub(官方模拟器)、AMQP、SQL、AWS SQS/SNS 等示例的docker-compose.yml同样位于各自的示例目录中。
第二层:Router
Publisher 与 Subscriber 是 Watermill 的底层组件。对大多数场景,你更想要的是高层 API——Router(详见 messages-router.md)。它自动处理消息订阅、并发分发、Ack/Nack 与优雅关闭。
Router 的关键实现细节在 message/router.go 中:HandlerFunc是消息到达时被调用的函数,msg.Ack()在HandlerFunc不返回错误时被自动调用,返回错误时则自动调用msg.Nack()(见 message/router.go);HandlerMiddleware则允许以装饰器模式包裹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 message's metadata to the produced messages middleware.CorrelationID, // The handler function is retried if it returns an error. // After MaxRetries, the message is Nacked and it's 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目前只有一个配置项CloseTimeout(Router 关闭时等待 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, we're 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.go;NewSlogLogger(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 的直白应用;
- 完整示例清单见 README(CQRS、Outbox 模式、SSE 等真实世界示例也在此列出)。
支持
如果任何地方不清楚,欢迎通过 support 页面 所列的渠道获取帮助。
【免费下载链接】watermillBuilding event-driven applications the easy way in Go.项目地址: https://gitcode.com/GitHub_Trending/wa/watermill
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考