UnixSocketPubSub
UnixSocketPubSub 是一种 PubSub 实现,使用 Unix domain socket 在单个主机上的进程间投递事件。它选举一个进程作为 broker,其他进程作为 client 连接。如果 broker 退出,剩余 client 会自动选举新的 broker。
当多个本地进程需要共享流时使用它,例如在 Mastra Code 终端界面中协调 thread stream。对于单进程投递,请使用 EventEmitterPubSub。对于跨主机的分布式投递,请使用 RedisStreamsPubSub 或 GoogleCloudPubSub。
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 状态。