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

Queue

了解如何使用 Effect 的 Queue,以内置背压实现轻量、类型安全且异步的工作流。

Queue 是一个轻量级的内存队列,内置背压(back pressure),能够以异步、纯函数式且类型安全的方式处理数据。

基本操作

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

API说明
Queue.offer向队列中添加一个类型为 A 的值。
Queue.take移除并返回队列中最旧的值。

示例(添加并取出一个元素)

import { Effect, Queue } from "effect"

const program = Effect.gen(function* () {
  // Creates a bounded queue with capacity 100
  const queue = yield* Queue.bounded<number>(100)
  // Adds 1 to the queue
  yield* Queue.offer(queue, 1)
  // Retrieves and removes the oldest value
  const value = yield* Queue.take(queue)
  return value
})

await Effect.runPromise(program) // => 1

创建 Queue

Queue 可以是有界的(带有指定容量),也可以是无界的(没有上限)。不同类型的队列在达到容量上限时,对新值的处理方式各不相同。

有界 Queue

有界队列在已满时会施加背压,也就是说,任何 Queue.offer 操作都会挂起,直到有可用空间为止。

示例(创建有界 Queue)

import { Effect, Queue } from "effect"

// Creating a bounded queue with a capacity of 100
const boundedQueue = Queue.bounded<number>(100)

;(await Effect.runPromise(boundedQueue)).capacity // => 100

丢弃式 Queue

丢弃式队列在队列已满时会丢弃新值。

示例(创建丢弃式 Queue)

import { Effect, Queue } from "effect"

// Creating a dropping queue with a capacity of 100
const droppingQueue = Queue.dropping<number>(100)

;(await Effect.runPromise(droppingQueue)).capacity // => 100

滑动式 Queue

滑动式队列在达到容量上限时会移除旧值,为新值腾出空间。

示例(创建滑动式 Queue)

import { Effect, Queue } from "effect"

// Creating a sliding queue with a capacity of 100
const slidingQueue = Queue.sliding<number>(100)

;(await Effect.runPromise(slidingQueue)).capacity // => 100

无界 Queue

无界队列没有容量限制,因此可以不受约束地添加新值。

示例(创建无界 Queue)

import { Effect, Queue } from "effect"

// Creates an unbounded queue without a capacity limit
const unboundedQueue = Queue.unbounded<number>()

;(await Effect.runPromise(unboundedQueue)).capacity // => Infinity

向 Queue 添加元素

offer

使用 Queue.offer 向队列中添加值。

示例(添加单个元素)

import { Effect, Queue } from "effect"

const program = Effect.gen(function* () {
  const queue = yield* Queue.bounded<number>(100)
  // Adds 1 to the queue
  yield* Queue.offer(queue, 1)
  return yield* Queue.size(queue)
})

await Effect.runPromise(program) // => 1

使用带背压的队列时,如果队列已满,Queue.offer 会挂起。为了避免阻塞主 Fiber,你可以把 Queue.offer 操作 fork 出去。

示例(用 Effect.forkChild 处理已满的队列)

import { Effect, Queue, Fiber } from "effect"

const program = Effect.gen(function* () {
  const queue = yield* Queue.bounded<number>(1)
  // Fill the queue with one item
  yield* Queue.offer(queue, 1)
  // Attempting to add a second item will suspend as the queue is full
  const fiber = yield* Effect.forkChild(Queue.offer(queue, 2))
  // Empties the queue to make space
  yield* Queue.take(queue)
  // Joins the fiber, completing the suspended offer
  yield* Fiber.join(fiber)
  // Returns the size of the queue after additions
  return yield* Queue.size(queue)
})

await Effect.runPromise(program) // => 1

offerAll

你也可以用 Queue.offerAll 一次性添加多个元素。

示例(添加多个元素)

import { Effect, Queue, Array } from "effect"

const program = Effect.gen(function* () {
  const queue = yield* Queue.bounded<number>(100)
  const items = Array.range(1, 10)
  // Adds all items to the queue at once
  yield* Queue.offerAll(queue, items)
  // Returns the size of the queue after additions
  return yield* Queue.size(queue)
})

await Effect.runPromise(program) // => 10

从 Queue 消费元素

take

Queue.take 操作会从队列中移除并返回最旧的元素。如果队列为空,Queue.take 会挂起,直到有元素被添加时才恢复。为避免阻塞,你可以把 Queue.take 操作 fork 到一个新的 Fiber 中。

示例(在 Fiber 中等待一个元素)

import { Effect, Queue, Fiber } from "effect"

const program = Effect.gen(function* () {
  const queue = yield* Queue.bounded<string>(100)
  // This take operation will suspend because the queue is empty
  const fiber = yield* Effect.forkChild(Queue.take(queue))
  // Adds an item to the queue
  yield* Queue.offer(queue, "something")
  // Joins the fiber to get the result of the take operation
  const value = yield* Fiber.join(fiber)
  return value
})

await Effect.runPromise(program) // => "something"

poll

若想在不挂起的情况下取出队列的第一个元素,请使用 Queue.poll。如果队列为空,Queue.poll 返回 None;如果队列中有元素,它会将该元素包装在 Some 中。

示例(轮询一个元素)

import { Effect, Queue, Option } from "effect"

const program = Effect.gen(function* () {
  const queue = yield* Queue.bounded<number>(100)
  // Adds items to the queue
  yield* Queue.offer(queue, 10)
  yield* Queue.offer(queue, 20)
  // Retrieves the first item if available
  const head = yield* Queue.poll(queue)
  return head
})

await Effect.runPromise(program) // => Option.some(10)

takeUpTo

要取出多个元素,请使用 Queue.takeBetween,它会返回最多达到指定数量的元素。 如果元素数量不足,它会返回所有当前可用的元素,而不会继续等待。

当不需要精确数量的元素时,这个函数对批处理特别有用。它能确保程序利用当前可用的数据继续工作。

如果你需要等待精确数量的元素再继续,可以考虑使用 takeN

示例(最多取出 N 个元素)

import { Effect, Queue } from "effect"

const program = Effect.gen(function* () {
  const queue = yield* Queue.bounded<number>(100)

  // Adds items to the queue
  yield* Queue.offer(queue, 1)
  yield* Queue.offer(queue, 2)
  yield* Queue.offer(queue, 3)

  // Retrieves up to 2 items (min must be at least 1, otherwise
  // takeBetween short-circuits and returns an empty array)
  const items = yield* Queue.takeBetween(queue, 1, 2)
  console.log(items)

  return "some result"
})

Effect.runPromise(program).then(console.log)
/*
Output:
[ 1, 2 ]
some result
*/

takeN

从队列中取出指定数量的元素。如果队列中的元素不足,该操作会挂起,直到所需数量的元素可用为止。

在每次处理都需要精确数量元素的场景中,这个函数很有用:它能确保在该批次凑齐之前,操作不会继续。

示例(取出固定数量的元素)

import { Effect, Queue, Fiber } from "effect"

const program = Effect.gen(function* () {
  // Create a queue that can hold up to 100 elements
  const queue = yield* Queue.bounded<number>(100)

  // Fork a fiber that attempts to take 3 items from the queue
  const fiber = yield* Effect.forkChild(
    Effect.gen(function* () {
      console.log("Attempting to take 3 items from the queue...")
      const chunk = yield* Queue.takeN(queue, 3)
      console.log(`Successfully took 3 items: ${chunk}`)
    }),
  )

  // Offer only 2 items initially
  yield* Queue.offer(queue, 1)
  yield* Queue.offer(queue, 2)
  console.log("Offered 2 items. The fiber is now waiting for the 3rd item...")

  // Simulate some delay
  yield* Effect.sleep("2 seconds")

  // Offer the 3rd item, which will unblock the takeN call
  yield* Queue.offer(queue, 3)
  console.log("Offered the 3rd item, which should unblock the fiber.")

  // Wait for the fiber to finish
  yield* Fiber.join(fiber)
  return "some result"
})

await Effect.runPromise(program) // => "some result"
/*
Output:
Offered 2 items. The fiber is now waiting for the 3rd item...
Attempting to take 3 items from the queue...
Offered the 3rd item, which should unblock the fiber.
Successfully took 3 items: 1,2,3
*/

takeAll

要一次性取出队列中的所有元素,请使用 Queue.takeAll。该操作会立即完成:如果队列为空,则返回空集合。

示例(取出所有元素)

import { Effect, Queue } from "effect"

const program = Effect.gen(function* () {
  const queue = yield* Queue.bounded<number>(100)
  // Adds items to the queue
  yield* Queue.offer(queue, 10)
  yield* Queue.offer(queue, 20)
  yield* Queue.offer(queue, 30)
  // Retrieves all items from the queue
  const chunk = yield* Queue.takeAll(queue)
  return chunk
})

await Effect.runPromise(program) // => [10, 20, 30]

关闭 Queue

shutdown

Queue.shutdown 操作允许你中断当前所有挂起在 offer*take* 操作上的 Fiber。该操作还会清空队列,并使之后任何 offer*take* 调用立即终止。

示例(关闭 Queue 时中断 Fiber)

import { Effect, Queue, Fiber } from "effect"

const program = Effect.gen(function* () {
  const queue = yield* Queue.bounded<number>(3)
  // Forks a fiber that waits to take an item from the queue
  const fiber = yield* Effect.forkChild(Queue.take(queue))
  // Shuts down the queue, interrupting the fiber
  yield* Queue.shutdown(queue)
  // Joins the interrupted fiber
  yield* Fiber.join(fiber)
})

const exit = await Effect.runPromiseExit(program)
exit._tag // => "Failure"

await

Queue.await 操作会一直等待,直到队列进入 Done 状态。由于 Queue.shutdown 会以一个中断原因(interrupt cause)完成队列,正在等待关闭的 effect 自身也会被中断,而不会作为正常成功来结束。若想无论结果如何都执行清理逻辑,请使用 Effect.onExit 而不是 Effect.andThen,并使用 Fiber.await 而不是 Fiber.join,这样中断就不会传播到执行 join 的 Fiber。

示例(等待 Queue 关闭)

import { Effect, Queue, Fiber, Console } from "effect"

const program = Effect.gen(function* () {
  const queue = yield* Queue.bounded<number>(3)
  // Forks a fiber to await queue shutdown and log a message,
  // regardless of whether the wait completes or is interrupted
  const fiber = yield* Effect.forkChild(
    Queue.await(queue).pipe(Effect.onExit(() => Console.log("shutting down"))),
  )
  // Shuts down the queue, triggering (and interrupting) the await in the fiber
  yield* Queue.shutdown(queue)
  yield* Fiber.await(fiber)
})

await Effect.runPromise(program) // => undefined
// Output: shutting down

只允许 offer / 只允许 take 的 Queue

有时,你可能希望代码的某些部分只能向队列添加值(Enqueue),或者只能从队列取出值(Dequeue)。Effect 提供了接口来强制约束这些特定的能力。

Enqueue

所有向队列添加值的方法都由 Enqueue 接口定义。这样就把队列限制为只能执行 offer 操作。

示例(把 Queue 限制为只能执行 offer 操作)

import { Queue } from "effect"

const send = (offerOnlyQueue: Queue.Enqueue<number>, value: number) => {
  // This queue is restricted to offer operations only

  // Error: cannot use take on an offer-only queue
  // @errors: 2345
  Queue.take(offerOnlyQueue)

  // Valid offer operation
  return Queue.offer(offerOnlyQueue, value)
}

Dequeue

类似地,所有从队列取出值的方法都由 Dequeue 接口定义,它把队列限制为只能执行 take 操作。

示例(把 Queue 限制为只能执行 take 操作)

import { Queue } from "effect"

const receive = (takeOnlyQueue: Queue.Dequeue<number>) => {
  // This queue is restricted to take operations only

  // Error: cannot use offer on a take-only queue
  // @errors: 2345
  Queue.offer(takeOnlyQueue, 1)

  // Valid take operation
  return Queue.take(takeOnlyQueue)
}

Queue 类型同时组合了 EnqueueDequeue,因此你可以轻松地把它传给代码的不同部分,并按需只暴露 EnqueueDequeue 的行为。

示例(同时使用只 offer 和只 take 的 Queue)

import { Effect, Queue } from "effect"

const send = (offerOnlyQueue: Queue.Enqueue<number>, value: number) => {
  return Queue.offer(offerOnlyQueue, value)
}

const receive = (takeOnlyQueue: Queue.Dequeue<number>) => {
  return Queue.take(takeOnlyQueue)
}

const program = Effect.gen(function* () {
  const queue = yield* Queue.unbounded<number>()

  // Add values to the queue
  yield* send(queue, 1)
  yield* send(queue, 2)

  // Retrieve values from the queue
  const first = yield* receive(queue)
  const second = yield* receive(queue)
  console.log(first)
  console.log(second)
  first // => 1
  second // => 2
})

await Effect.runPromise(program) // => undefined