跳到主要内容

GoogleCloudPubSub

GoogleCloudPubSub 是由 Google Cloud Pub/Sub 支持的 PubSub 实现。它使用 Google Cloud topic 和 subscription 跨进程及主机投递事件,并提供有序投递和消息确认。

将它用于 Google Cloud 上的分布式部署。对于单进程投递,请使用 EventEmitterPubSub。对于 Redis,请使用 RedisStreamsPubSub

每个 topic 对应一个 Google Cloud topic。带 group 的订阅共享一个 subscription,因此成员会竞争事件。没有 group 的订阅会创建每实例 subscription,因此每个实例都会收到每个事件。

安装
安装的直接链接

npm install @mastra/google-cloud-pubsub

使用示例
使用示例的直接链接

传入 Google Cloud client 配置,例如项目 ID。

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { GoogleCloudPubSub } from '@mastra/google-cloud-pubsub'

export const mastra = new Mastra({
pubsub: new GoogleCloudPubSub({
projectId: 'my-project',
}),
})

构造函数参数
构造函数参数的直接链接

config:

ClientConfig
Google Cloud Pub/Sub client 的配置,包括凭据和项目 ID。有关所有字段,请参阅 @google-cloud/pubsub client 文档。

方法
方法的直接链接

GoogleCloudPubSub 实现 PubSub contract。以下方法是此实现特有的。

init(topicName, group?)
inittopicname-group的直接链接

如果 topic 和 subscription 尚不存在,则创建它们,并返回 subscription。subscribe 会在内部调用此方法,因此很少需要直接调用。

await pubsub.init('workflow.events')

subscribe(topic, cb, options?)
subscribetopic-cb-options的直接链接

订阅一个 topic。使用 options.group 时,组成员共享 subscription 并竞争事件。没有 group 时,该实例会通过自己的 subscription 收到每个事件。

await pubsub.subscribe('workflow.events', (event, ack, nack) => {
console.log(event)
})

flush()
flush的直接链接

等待待处理确认完成。

await pubsub.flush()

destroy(topicName)
destroytopicname的直接链接

移除一个 topic 名称对应的 subscription 和 topic。使用此方法清理 Google Cloud 资源。

await pubsub.destroy('workflow.events')

确认
确认的直接链接

每个已投递事件均包含 acknack 函数。成功处理后调用 ack,以从 subscription 移除事件。当两者都未调用时,Google Cloud 会在其确认截止时间过期后重新投递该事件。