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 类型同时组合了 Enqueue 和 Dequeue,因此你可以轻松地把它传给代码的不同部分,并按需只暴露 Enqueue 或 Dequeue 的行为。
示例(同时使用只 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