Mastra Google Cloud Pub/Sub 集成指南:用 GCP Pub/Sub 为 AI 应用构建跨主机事件总线

Mastra Google Cloud Pub/Sub 集成指南:用 GCP Pub/Sub 为 AI 应用构建跨主机事件总线 Mastra Google Cloud Pub/Sub 集成指南用 GCP Pub/Sub 为 AI 应用构建跨主机事件总线【免费下载链接】mastraMastra is the modern TypeScript framework for AI-powered applications and agents.项目地址: https://gitcode.com/GitHub_Trending/ma/mastramastra/google-cloud-pubsub是 Mastra 官方的事件发布/订阅Pub/Sub后端它通过 Google Cloud Pub/Sub 的 topic 与 subscription 传递 Mastra 事件。当你的多进程 Mastra 应用需要跨主机、可持久化的事件投递而进程内 EventEmitter 已无法满足需求时就用它来替换默认的事件总线。读完本文你将掌握该包从安装、认证、接入 Mastra 到消费组group、本地仿真器测试等完整实战方法并理解其订阅命名、消息排序、ack/nack、localOnly 本地投递等底层实现原理。适用场景为什么需要替代进程内事件发射器Mastra 默认的事件机制是进程内 EventEmitter——事件只在同一个 Node.js 进程内流转。一旦你以多副本、多实例、Serverless 或分布式部署方式运行 Mastra例如多个 worker 共同处理同一个 agent 或 workflow就需要一个跨主机的事件通道让“生产事件”与“消费事件”发生在不同的机器/进程上。Google Cloud Pub/Sub 提供的正是这种能力持久化消息在 topic 上留存订阅者离线期间的消息可被补投跨主机生产者和消费者无需同进程天然支持多副本水平扩展可靠投递支持 ack/nack、redelivery重投与 deliveryAttempt投递尝试次数跟踪。GoogleCloudPubSub是mastra/core/events中抽象基类PubSub的一个实现见 packages/core/src/events/pubsub.ts因此它可以被无缝注入Mastra实例成为整个应用包括事件化 workflow 引擎、agent stream/observe 循环、durable agent 等的底层事件通道。安装npm install mastra/google-cloud-pubsub包的运行时依赖为google-cloud/pubsub^5.3.1并以mastra/core版本范围1.13.2-0 2.0.0-0作为 peer 依赖见 pubsub/google-cloud-pubsub/package.json。当前版本要求Node.js 22.13.0包本身为纯 ESM同时导出importdist/index.js与requiredist/index.cjs双入口。认证Application Default Credentials 优先在使用前你需要完成 Google Cloud 认证。官方推荐两种方式Application Default CredentialsADC在本地通过gcloud auth application-default login登录或在 GCE/GKE/Cloud Run 等环境由元数据服务器自动提供显式设置服务账号在启动 Mastra 前设置环境变量指向服务账号 JSONexport GOOGLE_APPLICATION_CREDENTIALS/path/to/service-account-key.json构造函数接收的配置直接透传给google-cloud/pubsub的PubSub客户端ClientConfig因此你也可以在构造参数中传入projectId、apiEndpoint用于连接仿真器等标准客户端选项见 src/index.ts。快速开始接入 Mastra在Mastra初始化时传入pubsub选项即可完成接入import { Mastra } from mastra/core/mastra; import { GoogleCloudPubSub } from mastra/google-cloud-pubsub; export const mastra new Mastra({ pubsub: new GoogleCloudPubSub({ projectId: process.env.GCP_PROJECT_ID!, }), });接入后事件化 workflow 的streamAsync、agent 的stream/observe、durable agent 的唤醒等机制都会经由 Google Cloud Pub/Sub 投递事件。在仓库的集成测试中可以看到完整链路pubsub/google-cloud-pubsub/src/index.test.ts 用new Mastra({ ..., pubsub: new GoogleCloudPubSub({ projectId: pubsub-test }) })配合mastra.startEventEngine()/mastra.stopEventEngine()验证了 workflow 流式执行、suspend/resume、agent 作为 step、sleep 等待等场景。事件模型理解底层契约要正确使用这个后端先了解它传递的事件结构。核心Event类型定义在 packages/core/src/events/types.ts字段类型说明typestring事件类型标识idstring事件唯一 ID接收端由 Pub/Sub message id 覆盖dataany事件负载runIdstring关联的 run 标识createdAtDate创建时间接收端由 message publishTime 覆盖indexnumber?顺序索引用于位置跟踪与断点恢复deliveryAttemptnumber?投递尝试次数从 1 开始每次 nack/重投递增后端不支持时默认为 1订阅回调签名EventCallback为(event, ack?, nack?) void | Promisevoidack()确认处理成功消息从队列移除nack()负确认消息延迟后重新入队重投都不调用消息保持 in-flight直到后端 ack 截止时间GCP 通常约 10 秒过期。需要注意在持久化后端上每个投递的事件都必须调用ack包括你过滤掉的事件否则消息会一直停留在订阅的 pending 集合里导致 topic 无限增长。回调返回 Promise 时reject 会被视为 nack但返回的 Promise 不会被等待需要顺序处理的订阅方必须自行串行化。源码实现剖析订阅命名、竞态合并与消息排序GoogleCloudPubSub的实现在 pubsub/google-cloud-pubsub/src/index.ts约 300 行下面拆解几个关键机制。订阅命名规则fan-out 与 groupgetSubscriptionName(topic: string, group?: string) { if (group) return ${topic}-${group}; return ${topic}-${this.instanceId}; }每个GoogleCloudPubSub实例在构造时都会生成一个crypto.randomUUID()作为instanceId不带 group 订阅订阅名为${topic}-${instanceId}每个实例进程拥有独立订阅实现fan-out——每个订阅者都能收到全部消息带 group 订阅订阅名为${topic}-${group}同一 group 下的所有订阅者共享同一个订阅消息只投递给 group 中的一个成员实现竞争消费competing consumers。group.test.ts中的测试验证了这一点fan-out 与 group 订阅在同一 topic 上会创建不同名称的底层订阅互不干扰、各自收到全部消息见 pubsub/google-cloud-pubsub/src/group.test.ts。init并发竞态合并与 ALREADY_EXISTS 恢复init()负责幂等地创建 topic 与 subscription并解决两个现实问题并发合并inFlightInit当生产者的agent.stream()与消费者的agent.observe()几乎同时订阅一个全新的 run topic 时两次init()会竞争创建同一个订阅。实现用inFlightInit记录进行中的创建 Promise后续调用直接复用避免重复创建见 src/index.tsALREADY_EXISTS 重连如果创建订阅失败且错误码为 gRPC code 6ALREADY_EXISTS——无论是并发竞争、其他进程通过 group 共享还是上次进程遗留——都会转而 attach 到已存在的订阅而不是抛出异常。这一逻辑不依赖 group 参数未分组场景同样生效这是 CHANGELOG 中记录的上游竞态修复见 pubsub/google-cloud-pubsub/CHANGELOG.md。创建订阅时还开启了两个关键选项await this.pubsub.topic(topicName).createSubscription(subscriptionName, { enableMessageOrdering: true, enableExactlyOnceDelivery: topicName workflows || !!group, });enableMessageOrdering: true配合 publish 时携带的orderingKey: workflows保证 workflow 事件按发布顺序投递enableExactlyOnceDelivery对workflowstopic 及带 group 的订阅启用精确一次投递避免重复处理。publishJSON 序列化、topic 归一化与自动建 topicawait topic.publishMessage({ data: Buffer.from(JSON.stringify(event)), orderingKey: workflows, });发布时事件被 JSON 序列化并附带orderingKey: workflows见 src/index.ts。此外 publish/subscribe/unsubscribe 都会对workflow.events.*前缀的 topic 做版本归一化若倒数第二段是v2则归一化为workflow.events.v2否则归一化为workflow.events.v1确保新旧版本的事件流收敛到同一个 topic。若 publish 遇到 code 5NOT_FOUNDtopic 不存在会先createTopic再重试发布。ack 缓冲与 flushackMessage使用message.ackWithResponse()并通过Promise.race加了 5 秒超时保护将进行中的 ack Promise 存入ackBufferflush()会Promise.all等待所有缓冲中的 ack 完成用于优雅关闭前排空见 src/index.ts 与 src/index.ts。unsubscribe按回调精确注销unsubscribe(topic, cb)会同时检查未分组与所有topic:group组合键仅当某个订阅上的回调清空时才移除 message listener、关闭并删除订阅状态见 src/index.ts。这意味着多次subscribe共享同一底层订阅时注销单个回调不会中断其他回调。localOnly不经过 Google Cloud 的进程内投递publish(topic, event, { localOnly: true })时事件完全不会进入 Google Cloud Pub/Sub而是通过localCallbacks按引用直接分发给同进程的订阅者见 src/index.ts。为什么要这么做事件化 agent 循环的workflows-finish运行结果中携带了MastraModelOutput类实例——如果把它 JSON 序列化进 Pub/Sub再在接收端JSON.parse回来实例的方法和原型如consumeStream会全部丢失。localOnly 按引用投递payload 上的活对象Date、Map、Error、MastraModelOutput等原型得以完整保留。group.test.ts中专门验证了localOnly 事件能拿到与原对象相等的引用且不会泄漏给共享同一 topic 的其他实例见 pubsub/google-cloud-pubsub/src/group.test.ts。高级用法group 消费组实现工作负载分担当多个 worker 进程需要分担同一批事件的处理时使用 group 选项// worker A 与 worker B 使用相同 groupPub/Sub 将每条消息只投递给其中一个 await pubsub.subscribe( jobs, async (event, ack, nack) { try { await processJob(event); await ack?.(); } catch (err) { console.error(err); await nack?.(); } }, { group: workers }, );同一 group 的所有订阅者共享一个订阅每条消息精确投递一次配合enableExactlyOnceDelivery不做任何处理、仅过滤部分事件时也必须对未处理事件调用ack真正的多进程竞争分发需要每个进程各建一个GoogleCloudPubSub实例单进程内多个实例共享 group 订阅时底层仍是同一个订阅消息不会被复制见 group.test.ts。本地开发与测试Pub/Sub Emulator集成测试依赖 Google Cloud Pub/Sub本地仿真器见 group.test.ts 的说明你可以用 Docker 快速拉起# 方式一docker compose docker compose -f .dev/docker-compose.yaml up -d pubsub-emulator # 方式二直接运行镜像 docker run -p 8085:8085 gcr.io/google.com/cloudsdktool/google-cloud-cli:emulators \ gcloud beta emulators pubsub start --host-port0.0.0.0:8085然后通过环境变量指定仿真器地址并设置apiEndpointconst ps new GoogleCloudPubSub({ projectId: pubsub-test, apiEndpoint: process.env.PUBSUB_EMULATOR_HOST ?? localhost:8085, });测试代码约定PUBSUB_EMULATOR_HOST默认值为localhost:8085。仿真器无需真实 GCP 账号即可完整验证 fan-out、group、localOnly、ALREADY_EXISTS 恢复等全部行为是本地迭代与 CI 的推荐方式。版本、限制与注意事项版本要求Node.js 22.13.0mastra/corepeer 依赖范围1.13.2-0 2.0.0-0当前 package 版本 1.1.3详见 package.json投递语义未分组订阅是 fan-out每条消息投递给所有订阅者分组订阅是竞争消费每条消息只给一个成员workflow 相关事件默认开启消息排序与精确一次投递必须 ack持久化后端上未 ack 的消息会滞留订阅导致 topic 无限增长回调抛错/返回 reject 会被视为 nack 触发重投可通过event.deliveryAttempt感知重投次数并实现指数退避或死信逻辑本地事件优先凡是携带活对象类实例、函数的进程内事件应使用localOnly发布避免序列化破坏原型跨主机场景则必须接受 JSON 序列化的纯数据负载生命周期应用关闭前建议调用flush()排空 ack 缓冲destroy(topicName)会关闭并删除订阅和 topic仅用于明确的清理场景。延伸阅读实现源码pubsub/google-cloud-pubsub/src/index.ts抽象基类与投递模式pull/push、clearTopic、getHistory、subscribeWithReplay 等packages/core/src/events/pubsub.ts事件类型与订阅选项group、startFrom、batchpackages/core/src/events/types.ts与 Mastra 事件化 workflow 的集成测试pubsub/google-cloud-pubsub/src/index.test.tsfan-out / group / localOnly 行为测试pubsub/google-cloud-pubsub/src/group.test.ts版本演进记录pubsub/google-cloud-pubsub/CHANGELOG.md【免费下载链接】mastraMastra is the modern TypeScript framework for AI-powered applications and agents.项目地址: https://gitcode.com/GitHub_Trending/ma/mastra创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考