基础并发
通过并发、中断与竞速来管理和控制 effect 的执行。
并发选项
Effect 提供了一些选项来管理 effect 的执行方式,尤其侧重控制有多少 effect 并发运行。
type Options = {
readonly concurrency?: Concurrency
}
concurrency 选项用于确定并发级别,取值如下:
type Concurrency = number | "unbounded"
下面我们详细探讨每一种配置。
这里的示例使用的是 Effect.all 函数,但这些选项同样适用于许多其他
Effect API。
顺序执行(默认)
默认情况下,如果你不指定任何并发选项,effect 会顺序执行,一个接一个。这意味着每个 effect 只会在前一个 effect 完成之后才开始。
示例(顺序执行)
import { Effect, Duration } from "effect"
// Helper function to simulate a task with a delay
const makeTask = (n: number, delay: Duration.Input) =>
Effect.promise(
() =>
new Promise<void>((resolve) => {
console.log(`start task${n}`) // Logs when the task starts
setTimeout(() => {
console.log(`task${n} done`) // Logs when the task finishes
resolve()
}, Duration.toMillis(delay))
}),
)
const task1 = makeTask(1, "200 millis")
const task2 = makeTask(2, "100 millis")
const sequential = Effect.all([task1, task2])
await Effect.runPromise(sequential) // => [undefined, undefined]
/*
Output:
start task1
task1 done
start task2 <-- task2 starts only after task1 completes
task2 done
*/
数字并发
你可以通过为 concurrency 设置一个 number 来控制有多少 effect 并发运行。例如,concurrency: 2 允许最多两个 effect 同时运行。
示例(限制为 2 个并发任务)
import { Effect, Duration } from "effect"
// Helper function to simulate a task with a delay
const makeTask = (n: number, delay: Duration.Input) =>
Effect.promise(
() =>
new Promise<void>((resolve) => {
console.log(`start task${n}`) // Logs when the task starts
setTimeout(() => {
console.log(`task${n} done`) // Logs when the task finishes
resolve()
}, Duration.toMillis(delay))
}),
)
const task1 = makeTask(1, "200 millis")
const task2 = makeTask(2, "100 millis")
const task3 = makeTask(3, "210 millis")
const task4 = makeTask(4, "110 millis")
const task5 = makeTask(5, "150 millis")
const numbered = Effect.all([task1, task2, task3, task4, task5], {
concurrency: 2,
})
await Effect.runPromise(numbered) // => [undefined, undefined, undefined, undefined, undefined]
/*
Output:
start task1
start task2 <-- active tasks: task1, task2
task2 done
start task3 <-- active tasks: task1, task3
task1 done
start task4 <-- active tasks: task3, task4
task4 done
start task5 <-- active tasks: task3, task5
task3 done
task5 done
*/
无界并发
当使用 concurrency: "unbounded" 时,并发运行的 effect 数量没有上限。
示例(无界并发)
import { Effect, Duration } from "effect"
// Helper function to simulate a task with a delay
const makeTask = (n: number, delay: Duration.Input) =>
Effect.promise(
() =>
new Promise<void>((resolve) => {
console.log(`start task${n}`) // Logs when the task starts
setTimeout(() => {
console.log(`task${n} done`) // Logs when the task finishes
resolve()
}, Duration.toMillis(delay))
}),
)
const task1 = makeTask(1, "200 millis")
const task2 = makeTask(2, "100 millis")
const task3 = makeTask(3, "210 millis")
const task4 = makeTask(4, "110 millis")
const task5 = makeTask(5, "150 millis")
const unbounded = Effect.all([task1, task2, task3, task4, task5], {
concurrency: "unbounded",
})
await Effect.runPromise(unbounded) // => [undefined, undefined, undefined, undefined, undefined]
/*
Output:
start task1
start task2
start task3
start task4
start task5
task2 done
task4 done
task5 done
task1 done
task3 done
*/
中断
Effect 中的所有 effect 都由 Fiber 执行。如果你没有自己创建 Fiber,那么它要么是由你正在使用的某个操作创建的(如果该操作是并发的),要么是由 Effect 运行时 系统创建的。
每当一个 effect 被运行时,都会创建一个 Fiber。并发运行 effect 时,会为每个并发 effect 创建一个 Fiber。
总结如下:
Effect是更高层的概念,用于描述一段带副作用的计算。它是惰性且不可变的,这意味着它表示一段可能产生值、也可能失败的计算,但并不会立即执行。- 而 Fiber 表示
Effect正在运行的执行过程。它可以被中断,也可以被等待以获取其结果。可以把它看作一种控制和交互正在进行的计算的方式。
Fiber 可以通过多种方式被中断。下面我们来探讨其中一些场景,并看看在 Effect 中如何中断 Fiber 的示例。
interrupt
可以使用 Effect.interrupt effect 来中断指定的 Fiber。
这个 effect 模拟它所在的 Fiber 被显式中断的行为。 执行时,它会让该 Fiber 立即停止运行,并捕获中断的详细信息,例如该 Fiber 的 ID 和它的启动时间。 如果使用 runPromiseExit 这类函数运行 effect,就可以在 Exit 类型中观察到由此产生的中断。
示例(无中断)
在这个例子中,程序在没有任何中断的情况下运行,记录了任务的开始与完成。
import { Effect, Exit } from "effect"
const program = Effect.gen(function* () {
console.log("start")
yield* Effect.sleep("2 seconds")
console.log("done")
return "some result"
})
await Effect.runPromiseExit(program) // => Exit.succeed("some result")
/*
Output:
start
done
*/
示例(发生中断)
这里,Fiber 在打印日志 "start" 之后、打印 "done" 之前被中断。Effect.interrupt 会停止该 Fiber,因此它永远不会走到最后那行日志。
import { Effect } from "effect"
const program = Effect.gen(function* () {
console.log("start")
yield* Effect.sleep("2 seconds")
return yield* Effect.interrupt
})
const exit = await Effect.runPromiseExit(program)
exit._tag // => "Failure"
/*
Output:
start
{
_id: 'Exit',
_tag: 'Failure',
cause: {
_id: 'Cause',
_tag: 'Interrupt',
fiberId: {
_id: 'FiberId',
_tag: 'Runtime',
id: 0,
startTimeMillis: ...
}
}
}
*/
onInterrupt
注册一个清理 effect,在某个 effect 被中断时运行。
这个函数允许你指定一个 effect,在 Fiber 被中断时运行。该 effect 会在 Fiber 被中断时执行, 让你可以执行清理或其他操作。
示例(在中断时运行清理操作)
在这个示例中,我们设置了一个处理器,每当 Fiber 被中断时就记录 “Cleanup completed”。然后展示了三种情况:成功的 effect、失败的 effect 以及被中断的 effect,以此说明处理器如何根据 effect 结束方式的不同而被触发。
import { Console, Effect, Exit } from "effect"
// This handler is executed when the fiber is interrupted
const handler = Effect.onInterrupt((_fibers) =>
Console.log("Cleanup completed"),
)
const success = Console.log("Task completed").pipe(
Effect.as("some result"),
handler,
)
await Effect.runPromise(success) // => "some result"
/*
Output:
Task completed
*/
const failure = Console.log("Task failed").pipe(
Effect.andThen(Effect.fail("some error")),
handler,
)
await Effect.runPromiseExit(failure) // => Exit.fail("some error")
/*
Output:
Task failed
*/
const interruption = Console.log("Task interrupted").pipe(
Effect.andThen(Effect.interrupt),
handler,
)
const interruptionExit = await Effect.runPromiseExit(interruption)
interruptionExit._tag // => "Failure"
/*
Output:
Task interrupted
Cleanup completed
*/
并发 effect 的中断
当并发运行多个 effect 时(例如使用 Effect.forEach),如果其中一个 effect 被中断,就会导致所有并发 effect 也一并被中断。
由此得到的 cause 会包含哪些 Fiber 被中断的信息。
示例(中断并发 effect)
import { Effect, Console } from "effect"
const program = Effect.forEach(
[1, 2, 3],
(n) =>
Effect.gen(function* () {
console.log(`start #${n}`)
yield* Effect.sleep(`${n} seconds`)
if (n > 1) {
return yield* Effect.interrupt
}
console.log(`done #${n}`)
}).pipe(Effect.onInterrupt(() => Console.log(`interrupted #${n}`))),
{ concurrency: "unbounded" },
)
const exit = await Effect.runPromiseExit(program)
console.log(JSON.stringify(exit, null, 2))
exit._tag // => "Failure"
/*
Output:
start #1
start #2
start #3
done #1
interrupted #2
interrupted #3
{
"_id": "Exit",
"_tag": "Failure",
"cause": {
"_id": "Cause",
"_tag": "Parallel",
"left": {
"_id": "Cause",
"_tag": "Interrupt",
"fiberId": {
"_id": "FiberId",
"_tag": "Runtime",
"id": 3,
"startTimeMillis": ...
}
},
"right": {
"_id": "Cause",
"_tag": "Sequential",
"left": {
"_id": "Cause",
"_tag": "Empty"
},
"right": {
"_id": "Cause",
"_tag": "Interrupt",
"fiberId": {
"_id": "FiberId",
"_tag": "Runtime",
"id": 0,
"startTimeMillis": ...
}
}
}
}
}
*/
竞速
race
这个函数接收两个 effect 并并发运行它们。第一个成功完成的 effect 将决定这次竞速的结果,而另一个 effect 会被中断。
如果两个 effect 都没有成功,该函数会以一个包含所有错误的 cause 失败。
当你希望并发运行两个 effect、但只关心第一个成功的那个时,这很有用。它常用于超时、重试等场景,或者当你希望优化为更快得到响应、而不必顾虑另一个 effect 时。
示例(两个任务都成功)
import { Effect, Console } from "effect"
const task1 = Effect.succeed("task1").pipe(
Effect.delay("200 millis"),
Effect.tap(Console.log("task1 done")),
Effect.onInterrupt(() => Console.log("task1 interrupted")),
)
const task2 = Effect.succeed("task2").pipe(
Effect.delay("100 millis"),
Effect.tap(Console.log("task2 done")),
Effect.onInterrupt(() => Console.log("task2 interrupted")),
)
const program = Effect.race(task1, task2)
await Effect.runPromise(program) // => "task2"
/*
Output:
task2 done
task1 interrupted
*/
示例(一个任务失败,一个任务成功)
import { Effect, Console } from "effect"
const task1 = Effect.fail("task1").pipe(
Effect.delay("100 millis"),
Effect.tap(Console.log("task1 done")),
Effect.onInterrupt(() => Console.log("task1 interrupted")),
)
const task2 = Effect.succeed("task2").pipe(
Effect.delay("200 millis"),
Effect.tap(Console.log("task2 done")),
Effect.onInterrupt(() => Console.log("task2 interrupted")),
)
const program = Effect.race(task1, task2)
await Effect.runPromise(program) // => "task2"
/*
Output:
task2 done
*/
示例(两个任务都失败)
import { Effect, Console } from "effect"
const task1 = Effect.fail("task1").pipe(
Effect.delay("100 millis"),
Effect.tap(Console.log("task1 done")),
Effect.onInterrupt(() => Console.log("task1 interrupted")),
)
const task2 = Effect.fail("task2").pipe(
Effect.delay("200 millis"),
Effect.tap(Console.log("task2 done")),
Effect.onInterrupt(() => Console.log("task2 interrupted")),
)
const program = Effect.race(task1, task2)
const exit = await Effect.runPromiseExit(program)
console.log(exit)
exit._tag // => "Failure"
/*
Output:
{
_id: 'Exit',
_tag: 'Failure',
cause: {
_id: 'Cause',
_tag: 'Parallel',
left: { _id: 'Cause', _tag: 'Fail', failure: 'task1' },
right: { _id: 'Cause', _tag: 'Fail', failure: 'task2' }
}
}
*/
如果你想处理最先完成的任务的结果,无论它成功还是失败,都可以使用 Effect.result 函数。这个函数会把结果包装为 Result 类型,让你可以看出结果是成功(Success)还是失败(Failure):
示例(用 Result 处理成功或失败)
import { Effect, Console, Result } from "effect"
const task1 = Effect.fail("task1").pipe(
Effect.delay("100 millis"),
Effect.tap(Console.log("task1 done")),
Effect.onInterrupt(() => Console.log("task1 interrupted")),
)
const task2 = Effect.succeed("task2").pipe(
Effect.delay("200 millis"),
Effect.tap(Console.log("task2 done")),
Effect.onInterrupt(() => Console.log("task2 interrupted")),
)
// Run both tasks concurrently, wrapping the result
// in Result to capture success or failure
const program = Effect.race(Effect.result(task1), Effect.result(task2))
await Effect.runPromise(program) // => Result.fail("task1")
/*
Output:
task2 interrupted
{ _id: 'Result', _tag: 'Failure', failure: 'task1' }
*/
raceAll
该函数会并发运行多个 effect,并返回第一个成功的 effect 的结果。一旦某个 effect 成功,其余的都会被中断。
如果所有 effect 都没有成功,该函数会以最后遇到的错误失败。
当你想要让多个 effect 竞速、但只关心第一个成功的那个时,这很有用。 它常用于超时、重试之类的场景, 或者当你想优化出更快的响应、 而不必关心其余 effect 的时候。
示例(所有任务都成功)
import { Effect, Console } from "effect"
const task1 = Effect.succeed("task1").pipe(
Effect.delay("100 millis"),
Effect.tap(Console.log("task1 done")),
Effect.onInterrupt(() => Console.log("task1 interrupted")),
)
const task2 = Effect.succeed("task2").pipe(
Effect.delay("200 millis"),
Effect.tap(Console.log("task2 done")),
Effect.onInterrupt(() => Console.log("task2 interrupted")),
)
const task3 = Effect.succeed("task3").pipe(
Effect.delay("150 millis"),
Effect.tap(Console.log("task3 done")),
Effect.onInterrupt(() => Console.log("task3 interrupted")),
)
const program = Effect.raceAll([task1, task2, task3])
await Effect.runPromise(program) // => "task1"
/*
Output:
task1 done
task2 interrupted
task3 interrupted
*/
示例(一个任务失败,两个任务成功)
import { Effect, Console } from "effect"
const task1 = Effect.fail("task1").pipe(
Effect.delay("100 millis"),
Effect.tap(Console.log("task1 done")),
Effect.onInterrupt(() => Console.log("task1 interrupted")),
)
const task2 = Effect.succeed("task2").pipe(
Effect.delay("200 millis"),
Effect.tap(Console.log("task2 done")),
Effect.onInterrupt(() => Console.log("task2 interrupted")),
)
const task3 = Effect.succeed("task3").pipe(
Effect.delay("150 millis"),
Effect.tap(Console.log("task3 done")),
Effect.onInterrupt(() => Console.log("task3 interrupted")),
)
const program = Effect.raceAll([task1, task2, task3])
await Effect.runPromise(program) // => "task3"
/*
Output:
task3 done
task2 interrupted
*/
示例(所有任务都失败)
import { Effect, Console } from "effect"
const task1 = Effect.fail("task1").pipe(
Effect.delay("100 millis"),
Effect.tap(Console.log("task1 done")),
Effect.onInterrupt(() => Console.log("task1 interrupted")),
)
const task2 = Effect.fail("task2").pipe(
Effect.delay("200 millis"),
Effect.tap(Console.log("task2 done")),
Effect.onInterrupt(() => Console.log("task2 interrupted")),
)
const task3 = Effect.fail("task3").pipe(
Effect.delay("150 millis"),
Effect.tap(Console.log("task3 done")),
Effect.onInterrupt(() => Console.log("task3 interrupted")),
)
const program = Effect.raceAll([task1, task2, task3])
const exit = await Effect.runPromiseExit(program)
console.log(exit)
exit._tag // => "Failure"
/*
Output:
{
_id: 'Exit',
_tag: 'Failure',
cause: { _id: 'Cause', _tag: 'Fail', failure: 'task2' }
}
*/
raceFirst
该函数接收两个 effect 并发运行它们, 返回第一个完成的那个的结果, 无论它是成功还是失败。
当你想要让两个操作竞速, 并希望以先完成的那个继续执行(无论它是成功还是失败)时, 这个函数很有用。
示例(两个任务都成功)
import { Effect, Console, Exit } from "effect"
const task1 = Effect.succeed("task1").pipe(
Effect.delay("100 millis"),
Effect.tap(Console.log("task1 done")),
Effect.onInterrupt(() =>
Console.log("task1 interrupted").pipe(Effect.delay("100 millis")),
),
)
const task2 = Effect.succeed("task2").pipe(
Effect.delay("200 millis"),
Effect.tap(Console.log("task2 done")),
Effect.onInterrupt(() =>
Console.log("task2 interrupted").pipe(Effect.delay("100 millis")),
),
)
const program = Effect.raceFirst(task1, task2).pipe(
Effect.tap(Console.log("more work...")),
)
await Effect.runPromiseExit(program) // => Exit.succeed("task1")
/*
Output:
task1 done
task2 interrupted
more work...
*/
示例(一个任务失败,一个任务成功)
import { Effect, Console, Exit } from "effect"
const task1 = Effect.fail("task1").pipe(
Effect.delay("100 millis"),
Effect.tap(Console.log("task1 done")),
Effect.onInterrupt(() =>
Console.log("task1 interrupted").pipe(Effect.delay("100 millis")),
),
)
const task2 = Effect.succeed("task2").pipe(
Effect.delay("200 millis"),
Effect.tap(Console.log("task2 done")),
Effect.onInterrupt(() =>
Console.log("task2 interrupted").pipe(Effect.delay("100 millis")),
),
)
const program = Effect.raceFirst(task1, task2).pipe(
Effect.tap(Console.log("more work...")),
)
await Effect.runPromiseExit(program) // => Exit.fail("task1")
/*
Output:
task2 interrupted
*/
观察胜出者
可选的 onWinner 回调会收到胜出的 fiber 及其索引(第一个 effect 为 0,第二个为 1)。该回调只用于观察:raceFirst 的结果仍然由最先完成的那个 effect 决定。
示例(观察哪个任务先完成)
import { Effect, Console } from "effect"
const task1 = Effect.succeed("task1").pipe(
Effect.delay("100 millis"),
Effect.tap(Console.log("task1 done")),
Effect.onInterrupt(() =>
Console.log("task1 interrupted").pipe(Effect.delay("100 millis")),
),
)
const task2 = Effect.succeed("task2").pipe(
Effect.delay("200 millis"),
Effect.tap(Console.log("task2 done")),
Effect.onInterrupt(() =>
Console.log("task2 interrupted").pipe(Effect.delay("100 millis")),
),
)
const program = Effect.raceFirst(task1, task2, {
onWinner: ({ index }) => console.log(`task${index + 1} won`),
})
Effect.runFork(program)
/*
Output:
task1 done
task1 won
task2 interrupted
*/