跳至主要內容

GoogleCloudPubSub

GoogleCloudPubSub 是由 Google Cloud Pub/Sub 支援的 PubSub 實作。它使用 Google Cloud 主題和訂閱,在不同進程及主機之間傳送事件,並支援按序傳送和訊息確認。

適用於 Google Cloud 上的分散式部署。單一進程傳送請使用 EventEmitterPubSub;Redis 則請使用 RedisStreamsPubSub

每個主題都會映射至一個 Google Cloud 主題。設有群組的訂閱會共用同一訂閱,因此成員會競爭接收事件。沒有群組的訂閱則會建立個別執行個體專用的訂閱,讓每個執行個體都收到所有事件。

安裝
安裝 的直接連結

npm install @mastra/google-cloud-pubsub

使用範例
使用範例 的直接連結

傳入 Google Cloud 用戶端設定,例如項目 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 用戶端的設定,包括憑證和項目 ID。所有欄位請參閱 @google-cloud/pubsub 用戶端文檔。

方法
方法 的直接連結

GoogleCloudPubSub 實作 PubSub 合約。以下方法為此實作所特有。

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

如果主題及訂閱尚未存在,便建立它們,然後傳回訂閱。subscribe 會在內部呼叫此方法,因此你很少需要直接呼叫它。

await pubsub.init('workflow.events')

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

訂閱主題。使用 options.group 時,群組成員會共用訂閱並競爭接收事件。沒有群組時,此執行個體會透過自己的訂閱接收所有事件。

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

flush()
flush 的直接連結

等待尚未完成的確認操作結束。

await pubsub.flush()

destroy(topicName)
destroytopicname 的直接連結

移除指定主題名稱的訂閱和主題。使用此方法清理 Google Cloud 資源。

await pubsub.destroy('workflow.events')

確認
確認 的直接連結

每個已傳送事件都包含 acknack 函數。處理成功後呼叫 ack,即可從訂閱移除事件。兩者皆未呼叫時,Google Cloud 會在確認期限屆滿後重新傳送事件。