拓冰建站拓冰建站
首页 / 资讯中心 / 正文

使用 Watermill 构建 SSE 实时推送应用:读写模型分离与事件驱动实战

使用 Watermill 构建 SSE 实时推送应用读写模型分离与事件驱动实战【免费下载链接】watermillBuilding event-driven applications the easy way in Go.项目地址: https://gitcode.com/GitHub_Trending/wa/watermill本文围绕 Watermill 仓库中的server-sent-events真实世界示例展开讲解如何用 Go 的 Watermill 消息路由配合 NATS Streaming 事件总线与 SSEServer-Sent Events协议实现一个创建文章 → 异步更新 Feed → 实时推送前端的类 Twitter 应用。读完本文你将掌握watermill-http的SSERouter与StreamAdapter用法、写模型MySQL与读模型MongoDB分离的落地方式以及前端EventSource消费实时流的完整链路。示例概览一个类 Twitter 的实时应用server-sent-events是 Watermill 仓库中的真实世界示例位于 _examples/real-world-examples/server-sent-events。它是一个类 Twitter 的 Web 应用使用 SSE 实现实时刷新你可以添加自己的帖子也可以点击按钮生成随机帖子支持#hashtag标签无论哪种方式Feed 列表与 Feed 内的帖子都会保持实时更新。示例的前端为一个单页应用Vue.js Bootstrap后端基于 Go 的 Watermill 消息路由构建底层事件总线选用 NATS Streaming。示例结构如下server/main.go应用入口负责组装存储、消息路由与 HTTP 路由server/http.goHTTP 路由、SSE 流适配器与 REST API 处理器server/event_handlers.goWatermill 消息路由器、事件处理器与 NATS Streaming 发布/订阅配置server/posts_storage.goMySQL 写模型存储server/feeds_storage.goMongoDB 读模型存储server/models.goPost、Feed及PostCreated、PostUpdated、FeedUpdated事件模型server/public/index.htmlVue.js 前端包含EventSource实时消费逻辑。运行与体验整个应用通过 Docker Compose 一键启动配置文件见 _examples/real-world-examples/server-sent-events/docker-compose.ymldocker-compose up启动后打开 http://localhost:8080 即可使用。Compose 中定义了四个服务服务镜像作用servergolang:1.25以go run .运行 Go 应用挂载./server目录暴露 8080 端口mysqlmysql:8.0写模型数据库启动时自动执行 schema.sql 初始化posts表mongomongo:3.6读模型数据库用于存储按标签组织的 Feednats-streamingnats-streaming:0.11.2事件总线承载post-created、post-updated、feed-updated三个主题的消息体验实时刷新的方式是在一个浏览器窗口添加帖子然后在第二个浏览器窗口观察 Feed 页面与帖子页面可以看到数据无需手动刷新即自动更新——这正是 SSE 长连接推送的效果。schema.sql中写模型表结构非常精简CREATE TABLE example.posts ( id VARCHAR(36) NOT NULL PRIMARY KEY, title VARCHAR(255) NOT NULL DEFAULT , content TEXT NOT NULL, author VARCHAR(255) NOT NULL DEFAULT );帖子 ID 由watermill.NewUUID()生成见server/http.go中CreatePost处理器标签并不直接存储而是在构造Post时通过正则从正文中的#hashtag提取出来见 models.go 的NewPost使用#([a-zA-Z0-9])匹配并统一转为小写、去重。整体架构写模型与读模型分离示例的核心架构思想是读写模型分离CQRS 风格的写路径 物化读模型写模型MySQL所有帖子以原子方式写入保证数据的一致性与完整性读模型MongoDB每个标签tag对应一个 Feed包含该标签下所有帖子Feed 由事件处理器异步维护数据开箱即用可直接返回给前端。原文档 README.md 中有一张架构图完整呈现了这一数据流流程可概括为用户提交帖子 →Create Post Handler写入 MySQL 并发布PostCreated事件 →Update Feeds Handler订阅事件、更新 MongoDB 中的 Feed 并发布FeedUpdated事件 →Get Feed Handler监听FeedUpdated通过 SSE 把最新 Feed 推送给前端。前端首次通过 HTTP GET 拉取初始数据之后每次 Feed 更新都会收到推送Fetch once, and every time the feed is updated。为什么要使用两套数据库原文档坦诚地指出对本示例应用而言使用双数据库引擎polyglot persistence显然是大材小用这样做是为了展示这种技术以及用 Watermill 落地它是多么容易。读写模型分离对高读/写比的应用很有价值所有写入原子地落到写模型此处为 MySQL事件处理器异步更新读模型此处为 Mongo读模型中的数据已经是可直接对外服务的形态读模型可以独立于写模型进行扩展。同时原文档也给出了务实的提醒使用该模式的前提是你的应用可以接受最终一致性eventual consistency大多数场景并不需要它请保持务实Be pragmatic!。事件流三类事件驱动全链路整个示例的所有异步通信包括 SSE 推送都由 Watermill 承担。原文档将事件定义为三类各自的职责如下事件触发时机处理动作PostCreated新帖子创建把帖子加入所有包含其标签的 FeedFeedUpdatedFeed 发生变化向所有正在浏览该 Feed 页面的客户端推送更新PostUpdated帖子被更新向正在浏览帖子页面的客户端推送更新同步更新包含其标签的所有 Feed其中PostUpdated对 Feed 的更新分为三种情况对应源码 feeds_storage.go 中UpdatePost的三步操作a)标签已存在更新该标签 Feed 中的帖子内容updatePostIfPresent用$set定位并替换posts.$b)新增了标签把帖子追加到新标签的 FeedappendPostIfNotPresent用$push$position: 0插入到 Feed 头部并通过$ne保证不重复插入c)删除了标签把帖子从旧标签的 Feed 中移除removePostIfNotInFeed用$pull按帖子 ID 删除。三个主题的常量定义在 event_handlers.goconst ( PostCreatedTopic post-created PostUpdatedTopic post-updated FeedUpdatedTopic feed-updated )事件模型见 models.go中PostUpdated同时携带OriginalPost与NewPost这正是推送过滤与 Feed 差分更新的依据type PostUpdated struct { OriginalPost Post json:original_post NewPost Post json:new_post OccurredAt time.Time json:occurred_at }核心组件一SSERouterSSERouter来自watermill-http包示例 go.mod 中为github.com/ThreeDotsLabs/watermill-http/v2 v2.3.1见 server/go.mod。创建 SSERouter 时传入一个上游订阅者upstream subscriber——来自该订阅者的消息将触发通过 HTTP 向客户端推送更新。示例使用 NATS Streaming 作为 Pub/Sub但按原文档说明这里可以是 Watermill 支持的任何 Pub/SubKafka、AMQP、GoChannel 等SSERouter只关心上游message.Subscriber接口。创建方式见 http.gosseRouter, err : watermillHTTP.NewSSERouter( watermillHTTP.SSERouterConfig{ UpstreamSubscriber: router.Subscriber, ErrorHandler: watermillHTTP.DefaultErrorHandler, }, router.Logger, ) if err ! nil { panic(err) }配置要点UpstreamSubscriber上游消息订阅者SSERouter从它消费事件消息决定推送给哪些连接ErrorHandler处理推送过程中的错误此处使用watermillHTTP.DefaultErrorHandler。启动 SSERouter 前需要先通过AddHandler把主题 流适配器注册进去。AddHandler返回一个标准http.Handler可以被任何路由库直接使用——这正是它和标准库无缝衔接的地方。示例中先注册处理器再挂载到 chi 路由postHandler : sseRouter.AddHandler(PostUpdatedTopic, postStream) feedHandler : sseRouter.AddHandler(FeedUpdatedTopic, feedStream) allFeedsHandler : sseRouter.AddHandler(FeedUpdatedTopic, allFeedsStream) r.Route(/api, func(r chi.Router) { r.Get(/posts/{id}, postHandler) r.Post(/posts, router.CreatePost) r.Post(/generate/post, router.GeneratePost) r.Patch(/posts/{id}, router.UpdatePost) r.Get(/feeds/{name}, feedHandler) r.Get(/feeds, allFeedsHandler) })注意这里/feeds/{name}与/feeds都注册在FeedUpdatedTopic上但使用了不同的流适配器实现单个 Feed 页与全部 Feed 列表两种粒度的实时推送。随后sseRouter.Run(context.Background())在 goroutine 中启动并通过-sseRouter.Running()等待就绪见 http.go。核心组件二StreamAdapter要与SSERouter协作需要实现一个StreamAdapter。原文档给出了接口的经典形态type StreamAdapter interface { // GetResponse returns the response to be sent back to client. // Any errors that occur should be handled and written to w, returning false as ok. GetResponse(w http.ResponseWriter, r *http.Request) (response interface{}, ok bool) // Validate validates if the incoming message should be handled by this handler. // Typically this involves checking some kind of model ID. Validate(r *http.Request, msg *message.Message) (ok bool) }从仓库实际代码看watermill-httpv2 的适配器将这两个职责细化成了四个方法可对照 http.go 中三个适配器的实现InitialStreamResponse(w, r) (response, ok)连接建立时返回初始响应相当于拉一次若出错则写错误并返回okfalseNextStreamResponse(r, msg) (response, ok)对每条上游消息决定是否推送新响应——对应原文档中Validate的职责返回okfalse表示该消息不适用于此连接。可见获取响应与校验消息是否该推送这两个核心职责在 v2 中被拆分为更明确的钩子但思路与原文档一致。按模型 ID 过滤的 Validate原文档给出的示例是校验消息是否属于 HTTP 请求中同一个帖子 ID只有匹配才推送func (p postStreamAdapter) Validate(r *http.Request, msg *message.Message) (ok bool) { postUpdated : PostUpdated{} err : json.Unmarshal(msg.Payload, postUpdated) if err ! nil { return false } postID : chi.URLParam(r, id) return postUpdated.OriginalPost.ID postID }仓库中的实现完全对应见 http.goNextStreamResponse先反序列化PostUpdated用chi.URLParam(r, id)取出路径参数与postUpdated.OriginalPost.ID比较不匹配就返回okfalse不向该连接推送匹配时再从 MySQL 查询最新帖子返回给客户端。feedStreamAdapter同理只是比较的是FeedUpdated.Name与chi.URLParam(r, name)。对每条消息都推送如果希望每条消息都触发更新Validate直接返回true即可func (f allFeedsStreamAdapter) Validate(r *http.Request, msg *message.Message) (ok bool) { return true }仓库中的allFeedsStreamAdapter正是如此NextStreamResponse不做任何过滤每次收到FeedUpdated消息都重新查询全部 Feed 并推送给/api/feeds的客户端见 http.go用于实时刷新导航栏中的 Feed 列表含每个 Feed 的帖子数。消息路由Watermill 路由器与 NATS Streaming事件处理的核心在 event_handlers.go 的SetupMessageRouter中。它创建标准 Watermill 路由器、注册Recoverer中间件、初始化 NATS Streaming 的发布者与订阅者然后注册两个处理器。NATS Streaming 的初始化使用watermill-nats包natsURL : stan.NatsURL(nats://nats-streaming:4222) pub, err : nats.NewStreamingPublisher(nats.StreamingPublisherConfig{ ClusterID: test-cluster, ClientID: publisher, StanOptions: []stan.Option{natsURL}, Marshaler: nats.GobMarshaler{}, }, logger) sub, err : nats.NewStreamingSubscriber(nats.StreamingSubscriberConfig{ ClusterID: test-cluster, ClientID: subscriber, StanOptions: []stan.Option{natsURL}, Unmarshaler: nats.GobMarshaler{}, }, logger)注意Marshaler/Unmarshaler统一使用nats.GobMarshaler保证消息序列化对称。随后通过router.AddHandler注册处理器将订阅主题 → 处理 → 发布主题串起来router.AddHandler( update-feeds-on-post-created, // 处理器名 PostCreatedTopic, // 订阅主题 sub, // 订阅者 FeedUpdatedTopic, // 发布主题 pub, // 发布者 func(msg *message.Message) (messages []*message.Message, err error) { event : PostCreated{} err json.Unmarshal(msg.Payload, event) if err ! nil { return nil, err } if len(event.Post.Tags) 0 { for _, tag : range event.Post.Tags { err feedsStorage.Add(msg.Context(), tag) if err ! nil { return nil, err } } err feedsStorage.AppendPost(msg.Context(), event.Post) if err ! nil { return nil, err } } return createFeedUpdatedEvents(event.Post.Tags) }, )处理函数返回的[]*message.Message会被路由器自动发布到FeedUpdatedTopic。createFeedUpdatedEvents为每个标签构造一个FeedUpdated事件消息UUID 由watermill.NewUUID()生成见 event_handlers.go从而驱动 SSE 推送。第二个处理器update-feeds-on-post-updated订阅PostUpdatedTopic先把新帖子的标签全部加入feedsStorage.Add再调用UpdatePost内部完成更新已有标签内容 / 追加新标签 / 移除被删标签三步最后对新标签 ∪ 原标签的并集发布FeedUpdated事件——这正是三个 StreamAdapter 收到推送消息的来源。SetupMessageRouter最后在 goroutine 中启动路由器并等待就绪然后把发布者与订阅者返回给main组装 HTTP 层go func() { err router.Run(context.Background()) if err ! nil { panic(err) } }() -router.Running() return pub, sub, nilmain.go中把sub作为UpstreamSubscriber注入SSERouter把pub包装成Publisher供 REST 处理器发布事件从而形成完整的HTTP 写入 → 事件发布 → 异步更新读模型 → SSE 推送闭环见 main.go。Publisher封装将任意interface{}事件序列化为 JSON 消息并发布func (p Publisher) Publish(topic string, event interface{}) error { payload, err : json.Marshal(event) if err ! nil { return err } msg : message.NewMessage(watermill.NewUUID(), payload) return p.publisher.Publish(topic, msg) }前端用 EventSource 消费 SSE 流前端应用使用 Vue.js 与 Bootstrap 构建见 server/public/index.htmlVue 2.7.16 vue-router 3.6.5 vue-resource。最值得关注的是EventSource的使用this.es new EventSource(/api/feeds/ this.feed) this.es.addEventListener(data, event { let data JSON.parse(event.data); this.posts_stream data.posts; }, false);前端通过监听data事件拿到推送的 JSON 载荷即StreamAdapter返回的response直接替换本地状态驱动 Vue 视图更新。示例中三处使用模式feeds-list组件监听/api/feeds实时刷新导航栏中的 Feed 列表每个 Feed 显示名称与帖子数ShowFeed组件监听/api/feeds/{name}进入页面时先建立 SSE 连接拿到初始帖子列表之后每次FeedUpdated到达都整体替换posts_stream切换 Feed 或组件销毁时会先es.close()关闭旧连接见setupFeed与destroyed钩子EditPost组件监听/api/posts/{post_id}实时展示帖子最新内容并注册了error事件处理区分连接关闭与其他错误。原文档也坦诚地补充了一句作者并非前端开发者index.html中的代码可能不够地道idiomatic欢迎提交 PR——因此本文对该前端的定位是功能演示重点仍应放在后端的事件驱动与 SSE 推送链路上。小结通过这个真实世界示例可以看到用 Watermill 搭建实时推送系统只需四步写模型落库并发布事件REST 处理器原子写入 MySQL随即通过Publisher发布PostCreated/PostUpdated异步更新读模型message.Router订阅事件把 Feed 物化到 MongoDB并产生FeedUpdatedSSE 路由推送watermill-http的SSERouter消费FeedUpdated通过StreamAdapter决定哪些连接该收到什么响应前端消费浏览器EventSource监听data事件实时刷新界面。其中读写模型分离与最终一致性是值得带入真实项目评估的架构决策而SSERouterStreamAdapter的过滤机制按模型 ID 精准推送或全量推送则提供了灵活的推送粒度控制。想要进一步动手可以阅读该目录下的 docker-compose.yml、event_handlers.go 与 http.go对照本文逐步验证事件链路。【免费下载链接】watermillBuilding event-driven applications the easy way in Go.项目地址: https://gitcode.com/GitHub_Trending/wa/watermill创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门