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
})
Effect.runPromise(program).then(console.log)
// Output: 1
创建 Queue
Queue 可以是有界的(带有指定容量),也可以是无界的(没有上限)。不同类型的队列在达到容量上限时,对新值的处理方式各不相同。
有界 Queue
有界队列在已满时会施加背压,也就是说,任何 Queue.offer 操作都会挂起,直到有可用空间为止。
示例(创建有界 Queue)
import { Queue } from "effect"
// Creating a bounded queue with a capacity of 100
const boundedQueue = Queue.bounded<number>(100)
丢弃式 Queue
丢弃式队列在队列已满时会丢弃新值。
示例(创建丢弃式 Queue)
import { Queue } from "effect"
// Creating a dropping queue with a capacity of 100
const droppingQueue = Queue.dropping<number>(100)
滑动式 Queue
滑动式队列在达到容量上限时会移除旧值,为新值腾出空间。
示例(创建滑动式 Queue)
import { Queue } from "effect"
// Creating a sliding queue with a capacity of 100
const slidingQueue = Queue.sliding<number>(100)
无界 Queue
无界队列没有容量限制,因此可以不受约束地添加新值。
示例(创建无界 Queue)
import { Queue } from "effect"
// Creates an unbounded queue without a capacity limit
const unboundedQueue = Queue.unbounded<number>()
向 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)
})
使用带背压的队列时,如果队列已满,Queue.offer 会挂起。为了避免阻塞主 Fiber,你可以把 Queue.offer 操作 fork 出去。
示例(用 Effect.fork 处理已满的队列)
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.fork(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)
})
Effect.runPromise(program).then(console.log)
// Output: 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)
})
Effect.runPromise(program).then(console.log)
// Output: 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.fork(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
})
Effect.runPromise(program).then(console.log)
// Output: something
poll
若想在不挂起的情况下取出队列的第一个元素,请使用 Queue.poll。如果队列为空,Queue.poll 返回 None;如果队列中有元素,它会将该元素包装在 Some 中。
示例(轮询一个元素)
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)
// Retrieves the first item if available
const head = yield* Queue.poll(queue)
return head
})
Effect.runPromise(program).then(console.log)
/*
Output:
{
_id: "Option",
_tag: "Some",
value: 10
}
*/
takeUpTo
要取出多个元素,请使用 Queue.takeUpTo,它会返回最多达到指定数量的元素。
如果元素数量不足,它会返回所有当前可用的元素,而不会继续等待。
当不需要精确数量的元素时,这个函数对批处理特别有用。它能确保程序利用当前可用的数据继续工作。
如果你需要等待精确数量的元素再继续,可以考虑使用 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
const chunk = yield* Queue.takeUpTo(queue, 2)
console.log(chunk)
return "some result"
})
Effect.runPromise(program).then(console.log)
/*
Output:
{ _id: 'Chunk', values: [ 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.fork(
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"
})
Effect.runPromise(program).then(console.log)
/*
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: {
"_id": "Chunk",
"values": [
1,
2,
3
]
}
some result
*/
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
})
Effect.runPromise(program).then(console.log)
/*
Output:
{
_id: "Chunk",
values: [ 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.fork(Queue.take(queue))
// Shuts down the queue, interrupting the fiber
yield* Queue.shutdown(queue)
// Joins the interrupted fiber
yield* Fiber.join(fiber)
})
awaitShutdown
Queue.awaitShutdown 操作可用于在队列关闭时运行一个 effect。它会等待直到队列被关闭;如果队列已经关闭,则会立即恢复。
示例(等待 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
const fiber = yield* Effect.fork(
Queue.awaitShutdown(queue).pipe(
Effect.andThen(Console.log("shutting down")),
),
)
// Shuts down the queue, triggering the await in the fiber
yield* Queue.shutdown(queue)
yield* Fiber.join(fiber)
})
Effect.runPromise(program)
// 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
console.log(yield* receive(queue))
console.log(yield* receive(queue))
})
Effect.runFork(program)
/*
Output:
1
2
*/