已发布 上游基线 bf46254 原文 ↗ 在 GitHub 编辑

PubSub

在 Effect 中使用 PubSub,轻松实现消息广播与异步通信。

PubSub 是一个异步消息中枢,发布者发送的消息可以被当前所有订阅者接收。

Queue 不同——在 Queue 中每个值只会投递给一个消费者——PubSub 会把每条已发布的消息广播给所有订阅者。因此,在需要消息广播而非负载分发的场景中,PubSub 是理想之选。

基本操作

PubSub<A> 存储类型为 A 的消息,并提供两个基础操作:

API说明
PubSub.publishPubSub 发送一条类型为 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
Subscribe Before Publishing

订阅者只会收到它处于活跃订阅状态期间发布的消息。若要确保某个订阅者收到某条特定消息, 请在发布该消息之前先建立订阅。

创建 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.capacityPubSub.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 还会终止所有关联的队列,确保关闭信号被有效传达。