PubSub
在 Effect 中使用 PubSub,轻松实现消息广播与异步通信。
PubSub 是一个异步消息中枢,发布者发送的消息可以被当前所有订阅者接收。
与 Queue 不同——在 Queue 中每个值只会投递给一个消费者——PubSub 会把每条已发布的消息广播给所有订阅者。因此,在需要消息广播而非负载分发的场景中,PubSub 是理想之选。
基本操作
PubSub<A> 存储类型为 A 的消息,并提供两个基础操作:
| API | 说明 |
|---|---|
PubSub.publish | 向 PubSub 发送一条类型为 A 的消息,返回一个 effect,指示该消息是否发布成功。 |
PubSub.subscribe | 创建一个 scoped effect,用于订阅该 PubSub,并在作用域结束时自动取消订阅。订阅者通过 Dequeue 接收消息,Dequeue 中保存着已发布的消息。 |
示例(向多个订阅者发布消息)
import { Effect, PubSub } from "effect"
const program = Effect.scoped(
Effect.gen(function* () {
const pubsub = yield* PubSub.bounded<string>(2)
// Two subscribers
const dequeue1 = yield* PubSub.subscribe(pubsub)
const dequeue2 = yield* PubSub.subscribe(pubsub)
// Publish a message to the pubsub
yield* PubSub.publish(pubsub, "Hello from a PubSub!")
// Each subscriber receives the message
const message1 = yield* PubSub.take(dequeue1)
const message2 = yield* PubSub.take(dequeue2)
console.log("Subscriber 1: " + message1)
console.log("Subscriber 2: " + message2)
message1 // => "Hello from a PubSub!"
message2 // => "Hello from a PubSub!"
}),
)
await Effect.runPromise(program) // => undefined
订阅者只会收到它处于活跃订阅状态期间发布的消息。若要确保某个订阅者收到某条特定消息, 请在发布该消息之前先建立订阅。
创建 PubSub
有界 PubSub
有界 PubSub 在达到容量上限时会对发布者施加背压(back pressure),暂停后续发布,直到有可用空间为止。
背压能确保所有订阅者在订阅期间都能收到全部消息。不过,如果某个订阅者速度较慢,消息投递也会随之变慢。
示例(创建有界 PubSub)
import { Effect, PubSub } from "effect"
// Creates a bounded PubSub with a capacity of 2
const boundedPubSub = PubSub.bounded<string>(2)
PubSub.capacity(await Effect.runPromise(boundedPubSub)) // => 2
丢弃式 PubSub
丢弃式 PubSub 在容量已满时会丢弃新值。如果消息被丢弃,PubSub.publish 操作会返回 false。
在丢弃式 pubsub 中,发布者可以继续发布新值,但不保证订阅者能收到所有消息。
示例(创建丢弃式 PubSub)
import { Effect, PubSub } from "effect"
// Creates a dropping PubSub with a capacity of 2
const droppingPubSub = PubSub.dropping<string>(2)
PubSub.capacity(await Effect.runPromise(droppingPubSub)) // => 2
滑动式 PubSub
滑动式 PubSub 会移除最早的消息,为新消息腾出空间,从而确保发布永不阻塞。
滑动式 pubsub 能避免慢订阅者影响消息投递速率。不过,慢订阅者仍有漏掉部分消息的风险。
示例(创建滑动式 PubSub)
import { Effect, PubSub } from "effect"
// Creates a sliding PubSub with a capacity of 2
const slidingPubSub = PubSub.sliding<string>(2)
PubSub.capacity(await Effect.runPromise(slidingPubSub)) // => 2
无界 PubSub
无界 PubSub 没有容量限制,因此发布总是立即成功。
无界 pubsub 保证所有订阅者都能收到全部消息,且不会拖慢消息投递。不过,如果消息的发布速度快于消费速度,它可以无限增长。
一般来说,除非你有特定的使用场景需要无界 pubsub,否则建议使用有界、丢弃式或滑动式 pubsub。
示例
import { Effect, PubSub } from "effect"
// Creates an unbounded PubSub with unlimited capacity
const unboundedPubSub = PubSub.unbounded<string>()
PubSub.capacity(await Effect.runPromise(unboundedPubSub)) // => Number.MAX_SAFE_INTEGER
PubSub 上的操作符
publishAll
PubSub.publishAll 函数让你可以一次性向 pubsub 发布多个值。
示例(发布多条消息)
import { Effect, PubSub } from "effect"
const program = Effect.scoped(
Effect.gen(function* () {
const pubsub = yield* PubSub.bounded<string>(2)
const dequeue = yield* PubSub.subscribe(pubsub)
yield* PubSub.publishAll(pubsub, ["Message 1", "Message 2"])
const messages = yield* PubSub.takeAll(dequeue)
console.log(messages)
messages // => ["Message 1", "Message 2"]
}),
)
await Effect.runPromise(program) // => undefined
capacity / size
你可以分别用 PubSub.capacity 和 PubSub.size 查看 pubsub 的容量与当前大小。
注意,PubSub.capacity 返回一个 number,因为容量在 pubsub 创建时就已设定,之后不会再改变。
相比之下,由于 pubsub 中消息的数量会随时间变化,PubSub.size 返回一个 effect,用于获取 pubsub 的当前大小。
示例(获取 PubSub 的容量与大小)
import { Effect, PubSub } from "effect"
const program = Effect.gen(function* () {
const pubsub = yield* PubSub.bounded<number>(2)
console.log(`capacity: ${PubSub.capacity(pubsub)}`)
const capacityMessage = `capacity: ${PubSub.capacity(pubsub)}` // => "capacity: 2"
console.log(`size: ${yield* PubSub.size(pubsub)}`)
const sizeMessage = `size: ${yield* PubSub.size(pubsub)}` // => "size: 0"
})
await Effect.runPromise(program) // => undefined
关闭 PubSub
要关闭 pubsub,请使用 PubSub.shutdown。你也可以用 PubSub.isShutdown 检查它是否已关闭,或用 PubSub.awaitShutdown 等待关闭完成。关闭 pubsub 还会终止所有关联的队列,确保关闭信号被有效传达。