已发布 上游基线 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, Queue } 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
    console.log("Subscriber 1: " + (yield* Queue.take(dequeue1)))
    console.log("Subscriber 2: " + (yield* Queue.take(dequeue2)))
  }),
)

Effect.runFork(program)
/*
Output:
Subscriber 1: Hello from a PubSub!
Subscriber 2: Hello from a PubSub!
*/
Subscribe Before Publishing

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

创建 PubSub

有界 PubSub

有界 PubSub 在达到容量上限时会对发布者施加背压(back pressure),暂停后续发布,直到有可用空间为止。

背压能确保所有订阅者在订阅期间都能收到全部消息。不过,如果某个订阅者速度较慢,消息投递也会随之变慢。

示例(创建有界 PubSub)

import { PubSub } from "effect"

// Creates a bounded PubSub with a capacity of 2
const boundedPubSub = PubSub.bounded<string>(2)

丢弃式 PubSub

丢弃式 PubSub 在容量已满时会丢弃新值。如果消息被丢弃,PubSub.publish 操作会返回 false

在丢弃式 pubsub 中,发布者可以继续发布新值,但不保证订阅者能收到所有消息。

示例(创建丢弃式 PubSub)

import { PubSub } from "effect"

// Creates a dropping PubSub with a capacity of 2
const droppingPubSub = PubSub.dropping<string>(2)

滑动式 PubSub

滑动式 PubSub 会移除最早的消息,为新消息腾出空间,从而确保发布永不阻塞。

滑动式 pubsub 能避免慢订阅者影响消息投递速率。不过,慢订阅者仍有漏掉部分消息的风险。

示例(创建滑动式 PubSub)

import { PubSub } from "effect"

// Creates a sliding PubSub with a capacity of 2
const slidingPubSub = PubSub.sliding<string>(2)

无界 PubSub

无界 PubSub 没有容量限制,因此发布总是立即成功。

无界 pubsub 保证所有订阅者都能收到全部消息,且不会拖慢消息投递。不过,如果消息的发布速度快于消费速度,它可以无限增长。

一般来说,除非你有特定的使用场景需要无界 pubsub,否则建议使用有界、丢弃式或滑动式 pubsub。

示例

import { PubSub } from "effect"

// Creates an unbounded PubSub with unlimited capacity
const unboundedPubSub = PubSub.unbounded<string>()

PubSub 上的操作符

publishAll

PubSub.publishAll 函数让你可以一次性向 pubsub 发布多个值。

示例(发布多条消息)

import { Effect, PubSub, Queue } 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"])
    console.log(yield* Queue.takeAll(dequeue))
  }),
)

Effect.runFork(program)
/*
Output:
{ _id: 'Chunk', values: [ 'Message 1', 'Message 2' ] }
*/

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)}`)
  console.log(`size: ${yield* PubSub.size(pubsub)}`)
})

Effect.runFork(program)
/*
Output:
capacity: 2
size: 0
*/

关闭 PubSub

要关闭 pubsub,请使用 PubSub.shutdown。你也可以用 PubSub.isShutdown 检查它是否已关闭,或用 PubSub.awaitShutdown 等待关闭完成。关闭 pubsub 还会终止所有关联的队列,确保关闭信号被有效传达。

PubSub 作为 Enqueue

PubSub 的操作符与 Queue 类似,主要区别在于用 PubSub.publishPubSub.subscribe 代替了 Queue.offerQueue.take。如果你已经熟悉 Queue 的用法,那么 PubSub 对你来说会很容易上手。

本质上,PubSub 可以被看作一个只允许写入的 Enqueue

import type { Queue } from "effect"

interface PubSub<A> extends Queue.Enqueue<A> {}

这里的 Enqueue 类型指的是只接受入队操作(enqueue,即写入)的队列。任何在这里入队的值都会被发布到 pubsub,而 shutdown 之类的操作也会影响该 pubsub。

这种设计让 PubSub 非常灵活,你可以在任何需要一个只接受已发布值的 Enqueue 的地方使用它。