跳到主要内容

UnixSocketPubSub

UnixSocketPubSub 是一种 PubSub 实现,使用 Unix domain socket 在单个主机上的进程间投递事件。它选举一个进程作为 broker,其他进程作为 client 连接。如果 broker 退出,剩余 client 会自动选举新的 broker。

当多个本地进程需要共享流时使用它,例如在 Mastra Code 终端界面中协调 thread stream。对于单进程投递,请使用 EventEmitterPubSub。对于跨主机的分布式投递,请使用 RedisStreamsPubSubGoogleCloudPubSub

UnixSocketPubSub 是 push transport:事件无需读取循环即可到达,因此 Mastra 不会为其运行 pull Worker。

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

传入所有参与进程共享的 socket 路径。

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { UnixSocketPubSub } from '@mastra/core/events'

export const mastra = new Mastra({
pubsub: new UnixSocketPubSub('/tmp/mastra/events.sock'),
})

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

socketPath:

string
Unix domain socket 的路径。共享流的所有进程必须使用相同路径。

options?:

UnixSocketPubSubOptions
可选配置。
number

属性
属性的直接链接

socketPath:

string
传入构造函数的 socket 路径。

supportedModes:

ReadonlyArray<"pull" | "push">
返回 ["push"]

isBroker:

boolean
此实例当前是否充当 broker。

remoteClientCount:

number
连接到此 broker 的远程 client 数量。当该实例不是 broker 时始终为 0。

方法
方法的直接链接

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

close()
close的直接链接

关闭 socket 连接;当此实例是 broker 时,释放 broker 角色。在优雅关闭期间调用此方法。

await pubsub.close()

Broker 选举
Broker 选举的直接链接

第一个绑定 socket 的进程成为 broker,并在所有已连接 client 之间路由事件。其他进程作为 client 连接。当 broker 退出时,独占锁文件会串行化下一次选举。恰好一个 client 成为新 broker,其余 client 重新订阅它。这避免了两个进程同时作为 broker 的 split-brain 状态。