
Watermill 实战用 SQL Publisher 将 Google Cloud Pub/Sub 事件持久化为 MySQL 事件日志【免费下载链接】watermillBuilding event-driven applications the easy way in Go.项目地址: https://gitcode.com/GitHub_Trending/wa/watermill本篇文章基于仓库中的 persistent-event-log 真实示例讲解如何用 Watermill 的Router将一对发布者/订阅者代理起来把 Google Cloud Pub/Sub 上无存储能力的事件流落地到 MySQL 持久化表中形成可审计、可重放的事件日志。读完本文你将掌握 SQL Publisher 的接入方式、Router.AddHandler的桥接原理以及如何在缺少消息存储的 Pub/Sub 之上自建事件日志。背景当 Pub/Sub 不提供消息存储时一些消息系统如 Kafka天然支持将已处理的消息长时间甚至永久保存在 Broker 中这对审计audit或未来按需重放replay特定消息非常有价值。但并非所有 PubSub 都提供存储能力——例如 Google Cloud Pub/Sub 默认只做短期投递不会为订阅者保留完整的历史事件流。那么问题来了如果项目使用的是一个不提供存储的 PubSub却又需要一份持久的、可查询的事件日志该怎么办答案就是本示例展示的方案用 Watermill 的 Router 将无存储 PubSub 的订阅者与SQL 数据库的发布者桥接起来把每条消费到的事件写入数据库表从而让数据库充当持久化事件日志。本示例用 Google Cloud Pub/Sub 作为事件源、MySQL 作为存储端Google Cloud Pub/Sub 仅作为演示任何其他 Watermill 支持的 Subscriber 都可以替换它。解决方案Router 作为 PubSub 之间的代理Watermill 的 Router 允许注册一个订阅 处理 发布的处理器从订阅者收到消息经过处理函数再把可能转换后的消息交给发布者。因此把 Google Cloud Pub/Sub 的 Subscriber 接到 SQL 的 Publisher 上就完成了一次从内存/短期队列到持久化存储的搬运。整体数据流如下simulateEvents (Google Cloud Publisher模拟产生事件) │ Publish(events, msg) ▼ Google Cloud Pub/Sub 主题 events │ Router.AddHandler(googlecloud-to-mysql) 订阅消费 ▼ 处理函数json.Unmarshal 解析 payload → log 打印 → 原样返回消息 │ SQL Publisher.Publish(events 表) ▼ MySQL 表 watermill.eventswatermill_events默认 MySQL schema示例中的核心桥接代码位于 main.gorouter.AddHandler( googlecloud-to-mysql, // handler 名称必须唯一 googleCloudTopic, // 订阅主题events subscriber, // Google Cloud Pub/Sub 订阅者 mysqlTable, // 发布目标MySQL 表名events publisher, // SQL 发布者 func(msg *message.Message) ([]*message.Message, error) { consumedEvent : event{} err : json.Unmarshal(msg.Payload, consumedEvent) if err ! nil { return nil, err } log.Printf(received event %v with UUID %s, consumedEvent, msg.UUID) return []*message.Message{msg}, nil // 原样透传给 SQL Publisher }, )从 message/router.go 的AddHandler签名可以看到每个 handler 都绑定了一对(subscriber, publisher)处理函数返回的[]*message.Message会被自动发布到publishTopic。这里的处理函数没有修改消息只是把消息原样透传这与 Router 内置的 PassthroughHandler 行为一致var PassthroughHandler HandlerFunc func(msg *Message) ([]*Message, error) { return []*Message{msg}, nil }也就是说你完全可以用message.PassthroughHandler直接替代这个匿名函数——示例保留它是为了顺带演示如何在搬运过程中对 payload 做解析、校验或转换。消息顺序用 OccurredAt 而非投递顺序Watermill 对消息顺序的保证取决于底层 PubSub。Google Cloud Pub/Sub 这类系统不保证消息的严格有序投递因此示例在事件 payload 中内置了OccurredAt发生时间字段供后续消费方自行排序type event struct { Name string json:name OccurredAt string json:occurred_at }模拟器simulateEvents每秒钟发布一条UserSignedUp事件OccurredAt取 UTC 时间并按 RFC3339 格式化见 main.goe : event{ Name: UserSignedUp, OccurredAt: time.Now().UTC().Format(time.RFC3339), }这样写入数据库后即使各条记录的created_at入库时间与事件实际发生时间存在细微偏差下游依然可以依据occurred_at还原业务上的真实顺序。运行环境与依赖示例依赖 Docker 与 docker-compose通过 docker-compose.yml 一键拉起三个服务服务镜像作用servergolang:1.25挂载当前目录并执行go run main.gomysqlmysql:8.0存储事件日志预建数据库watermill空密码googlecloudgoogle/cloud-sdk:228.0.0启动本地 Pub/Sub 模拟器emulator监听0.0.0.0:8085关键点server通过环境变量PUBSUB_EMULATOR_HOST: googlecloud:8085把 Google Cloud Pub/Sub 客户端指向本地模拟器因此无需真实 GCP 账号即可运行mysql暴露3306:3306方便在宿主机上用客户端连接检查数据googlecloud服务内部通过gcloud beta emulators pubsub start --host-port0.0.0.0:8085启动模拟器。依赖版本记录在 go.mod 中核心包括github.com/ThreeDotsLabs/watermill v1.5.1消息、Router、中间件github.com/ThreeDotsLabs/watermill-googlecloud/v2 v2.0.0Google Cloud Pub/Sub PubSubgithub.com/ThreeDotsLabs/watermill-sql/v4 v4.1.3SQL Publisher/Subscribergithub.com/go-sql-driver/mysql v1.10.0MySQL 驱动启动与验证在 persistent-event-log 目录下执行docker-compose up数秒后模拟器开始产生事件Router 消费并写入 MySQL。在另一个终端验证表内容docker-compose exec mysql mysql -e select * from watermill.watermill_events;预期输出如下README 原始记录-------------------------------------------------------------------------------------- | offset | uuid | created_at | payload | metadata | -------------------------------------------------------------------------------------- | 1 | 2faf6a14-f52a-4d6c-a4be-7355db428be1 | 2019-08-17 12:23:35 | {...} | {} | | 2 | cccfe73c-1968-4e20-b8b7-3763f68dc60b | 2019-08-17 12:23:35 | {...} | {} | | 3 | e8585f50-5e38-4569-bd93-fe4f6e960e61 | 2019-08-17 12:23:36 | {...} | {} | | 4 | 2d364b7e-fc4d-459c-972a-8859c8f1a655 | 2019-08-17 12:23:37 | {...} | {} | | 5 | 3b9da717-aad8-4e4b-a6e2-2d7040454015 | 2019-08-17 12:23:38 | {...} | {} | | 6 | 5c07a2e7-464e-4ffb-8ada-0e2f02e48111 | 2019-08-17 12:23:39 | {...} | {} | | 7 | 60a30b9e-6a40-4f41-94f9-8e7c8a38a998 | 2019-08-17 12:23:40 | {...} | {} | | 8 | 3d28a15a-7448-4535-9b79-27111579e341 | 2019-08-17 12:23:41 | {...} | {} | | 9 | 3c448aff-6bdd-4fc4-9b56-8bacab0b2746 | 2019-08-17 12:23:42 | {...} | {} | | 10 | 9b56ca67-4c47-4bcd-931f-86f9af62775d | 2019-08-17 12:23:43 | {...} | {} | --------------------------------------------------------------------------------------可以看到表由offset自增主键兼作日志序号、uuid消息 UUID、created_at入库时间、payload事件 JSON和metadata消息元数据组成。offset单调递增正适合作为事件日志的全局序号由于 Google Cloud Pub/Sub 不保证有序业务排序请优先使用 payload 中的occurred_at。代码逐段剖析1. 创建 MySQL 连接createDBmain.go 使用go-sql-driver/mysql构造 DSNconf : driver.NewConfig() conf.Net tcp conf.User root conf.Addr mysql // docker-compose 中的服务名 conf.DBName watermill // 与 docker-compose 中 MYSQL_DATABASE 一致 db, err : stdSQL.Open(mysql, conf.FormatDSN()) // db.Ping() 确认连接可用注意Addr指向的是 compose 网络中的服务名mysql而不是localhost——这是容器间通信的关键。2. 创建 Google Cloud Pub/Sub 订阅者createSubscribersub, err : googlecloud.NewSubscriber( googlecloud.SubscriberConfig{ ProjectID: example, }, logger, )ProjectID在本地模拟器场景下任意指定即可。对比 pubsubs/googlecloud 示例 可以看到SubscriberConfig还支持GenerateSubscriptionName自定义订阅名本示例使用默认命名Router 会自动为 handler 建立订阅。3. 创建 SQL PublishercreatePublisherpub, err : sql.NewPublisher( sql.BeginnerFromStdSQL(db), sql.PublisherConfig{ SchemaAdapter: sql.DefaultMySQLSchema{}, AutoInitializeSchema: true, }, logger, )这是整个方案的核心两个配置项各自承担重要职责BeginnerFromStdSQL(db)把标准库*sql.DB包装成sql.BeginnerSQL Publisher 通过它在插入事件时开启事务SchemaAdapter: sql.DefaultMySQLSchema{}使用 Watermill SQL 为 MySQL 预置的默认表结构即上面看到的watermill_events表含 offset/uuid/created_at/payload/metadata 列AutoInitializeSchema: true启动时自动执行建表 SQL无需手工CREATE TABLE。4. 模拟事件生产simulateEventspub, err : googlecloud.NewPublisher(googlecloud.PublisherConfig{ ProjectID: example, }, logger) for { e : event{Name: UserSignedUp, OccurredAt: time.Now().UTC().Format(time.RFC3339)} payload, _ : json.Marshal(e) err pub.Publish(googleCloudTopic, message.NewMessage(watermill.NewUUID(), payload)) time.Sleep(time.Second) }每秒钟向events主题发布一条新消息。message.NewMessage(watermill.NewUUID(), payload)创建的消息包含一个全局唯一 UUID——这正是写入数据库uuid列的值。关于消息结构可见 message/message.goMessage由UUID、Metadata类似 HTTP 头可存无需解析 payload 的附加数据、Payload构成并内置Ack()/Nack()确认机制。Router 桥接的底层原理消息确认与失败重投AddHandler注册的 handler 在收到消息后执行处理函数Router 会根据处理结果自动确认消息源码见 message/router.go 的handleMessage处理函数返回nil错误 → 发布产出消息成功后调用msg.Ack()处理函数返回错误 → 调用msg.Nack()由底层 PubSub 决定重投发布产出消息失败 → 同样msg.Nack()。这意味着本示例天然具备**至少一次at-least-once**语义只要 MySQL 写入失败消息就不会被确认Google Cloud Pub/Sub 会重新投递避免事件丢失。但也要注意重复投递可能造成同一事件被写入多次生产环境应结合幂等如按uuid去重处理。优雅退出与异常兜底示例在启动 Router 时挂载了插件与中间件main.gorouter.AddPlugin(plugin.SignalsHandler) router.AddMiddleware(middleware.Recoverer)plugin.SignalsHandler监听SIGINT/SIGTERM收到信号后调用router.Close()优雅关闭等待CloseTimeout默认 30 秒见 router.go 的RouterConfig.setDefaults内正在处理的消息完成middleware.Recoverer捕获 handler 中的 panic附带完整堆栈转为错误返回从而触发 Nack 而非让进程崩溃。最后通过router.Run(context.Background())阻塞运行main.go。Run内部会先执行所有插件再调用RunHandlers为每个 handler 建立订阅并派发消息见 router.go。进阶自定义 Schema 与表结构示例默认使用DefaultMySQLSchema。若默认表结构不满足需求例如想改用BINARY(16)存 UUID、增减列、调整索引可以基于默认实现扩展SchemaAdapter重写SchemaInitializingQueries控制建表 SQL只有在配置InitializeSchema: true或AutoInitializeSchema: true时才会执行重写UnmarshalMessage/MarshalMessage控制数据库记录与message.Message之间的互转。相关的完整配置项可参考仓库中的 SQL Pub/Sub 文档。该文档还强调SQL Pub/Sub 的订阅端通过周期性SELECT轮询新记录并记录已处理位置如果使用默认的自增offset作为定位依据则无需额外适配。若你的日志表只增不改且需要消费后删除的队列语义文档还提供了不依赖消费者组、可自定义WHERE条件的 Queue schema 作为替代方案。总结本示例展示了一条极具实用价值的 Watermill 集成路径用 Router 把短生命周期 PubSub 的订阅者与SQL 持久化发布者桥接起来数据库即事件日志。其核心要点可归纳为Router.AddHandler天然支持异构 PubSub 之间的代理处理函数可原样透传也可解析/转换SQL Publisher 通过DefaultMySQLSchemaAutoInitializeSchema快速落库offset提供单调递增的日志序号无序 PubSub 下的顺序还原依赖业务字段如occurred_at不要依赖投递顺序Router 的 Ack/Nack 语义保证了写库失败即重投的至少一次交付配合幂等设计即可用于审计与重放场景。如果你需要把 Kafka、NATS 或 AMQP 等任何其他 Watermill 支持的 PubSub 事件持久化只需替换本示例中的 Subscriber 一行即可其余架构完全复用。【免费下载链接】watermillBuilding event-driven applications the easy way in Go.项目地址: https://gitcode.com/GitHub_Trending/wa/watermill创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考