基础并发
通过并发、中断与竞速来管理和控制 effect 的执行。
并发选项
Effect 提供了一些选项来管理 effect 的执行方式,尤其侧重控制有多少 effect 并发运行。
type Options = {
readonly concurrency?: Concurrency
}
concurrency 选项用于确定并发级别,取值如下:
type Concurrency = number | "unbounded" | "inherit"
下面我们详细探讨每一种配置。
这里的示例使用的是 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.DurationInput) =>
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])
Effect.runPromise(sequential)
/*
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.DurationInput) =>
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,
})
Effect.runPromise(numbered)
/*
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.DurationInput) =>
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",
})
Effect.runPromise(unbounded)
/*
Output:
start task1
start task2
start task3
start task4
start task5
task2 done
task4 done
task5 done
task1 done
task3 done
*/
继承并发
当使用 concurrency: "inherit" 时,并发级别会从周围的上下文中继承。这个上下文可以用 Effect.withConcurrency(number | "unbounded") 设置。如果没有提供上下文,默认值为 "unbounded"。
示例(从上下文继承并发)
import { Effect, Duration } from "effect"
// Helper function to simulate a task with a delay
const makeTask = (n: number, delay: Duration.DurationInput) =>
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")
// Running all tasks with concurrency: "inherit",
// which defaults to "unbounded"
const inherit = Effect.all([task1, task2, task3, task4, task5], {
concurrency: "inherit",
})
Effect.runPromise(inherit)
/*
Output:
start task1
start task2
start task3
start task4
start task5
task2 done
task4 done
task5 done
task1 done
task3 done
*/
如果你使用 Effect.withConcurrency,并发配置就会调整为指定的选项。
示例(设置并发选项)
import { Effect, Duration } from "effect"
// Helper function to simulate a task with a delay
const makeTask = (n: number, delay: Duration.DurationInput) =>
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")
// Running tasks with concurrency: "inherit",
// which will inherit the surrounding context
const inherit = Effect.all([task1, task2, task3, task4, task5], {
concurrency: "inherit",
})
// Setting a concurrency limit of 2
const withConcurrency = inherit.pipe(Effect.withConcurrency(2))
Effect.runPromise(withConcurrency)
/*
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
*/
中断
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 } from "effect"
const program = Effect.gen(function* () {
console.log("start")
yield* Effect.sleep("2 seconds")
console.log("done")
return "some result"
})
Effect.runPromiseExit(program).then(console.log)
/*
Output:
start
done
{ _id: 'Exit', _tag: 'Success', value: 'some result' }
*/
示例(发生中断)
这里,Fiber 在打印日志 "start" 之后、打印 "done" 之前被中断。Effect.interrupt 会停止该 Fiber,因此它永远不会走到最后那行日志。
import { Effect } from "effect"
const program = Effect.gen(function* () {
console.log("start")
yield* Effect.sleep("2 seconds")
yield* Effect.interrupt
console.log("done")
return "some result"
})
Effect.runPromiseExit(program).then(console.log)
/*
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 } 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,
)
Effect.runFork(success)
/*
Output:
Task completed
*/
const failure = Console.log("Task failed").pipe(
Effect.andThen(Effect.fail("some error")),
handler,
)
Effect.runFork(failure)
/*
Output:
Task failed
*/
const interruption = Console.log("Task interrupted").pipe(
Effect.andThen(Effect.interrupt),
handler,
)
Effect.runFork(interruption)
/*
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) {
yield* Effect.interrupt
}
console.log(`done #${n}`)
}).pipe(Effect.onInterrupt(() => Console.log(`interrupted #${n}`))),
{ concurrency: "unbounded" },
)
Effect.runPromiseExit(program).then((exit) =>
console.log(JSON.stringify(exit, null, 2)),
)
/*
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)
Effect.runFork(program)
/*
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)
Effect.runFork(program)
/*
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)
Effect.runPromiseExit(program).then(console.log)
/*
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.either 函数。这个函数会把结果包装为 Either 类型,让你可以看出结果是成功(Right)还是失败(Left):
示例(用 Either 处理成功或失败)
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")),
)
// Run both tasks concurrently, wrapping the result
// in Either to capture success or failure
const program = Effect.race(Effect.either(task1), Effect.either(task2))
Effect.runPromise(program).then(console.log)
/*
Output:
task2 interrupted
{ _id: 'Either', _tag: 'Left', left: '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])
Effect.runFork(program)
/*
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])
Effect.runFork(program)
/*
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])
Effect.runPromiseExit(program).then(console.log)
/*
Output:
{
_id: 'Exit',
_tag: 'Failure',
cause: { _id: 'Cause', _tag: 'Fail', failure: 'task2' }
}
*/
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).pipe(
Effect.tap(Console.log("more work...")),
)
Effect.runPromiseExit(program).then(console.log)
/*
Output:
task1 done
task2 interrupted
more work...
{ _id: 'Exit', _tag: 'Success', value: 'task1' }
*/
示例(一个任务失败,一个任务成功)
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").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...")),
)
Effect.runPromiseExit(program).then(console.log)
/*
Output:
task2 interrupted
{
_id: 'Exit',
_tag: 'Failure',
cause: { _id: 'Cause', _tag: 'Fail', failure: 'task1' }
}
*/
断开 effect
Effect.raceFirst 函数会在另一方完成后安全地中断“输掉”的那个 effect,但直到输掉的一方被干净地终止之前,它都不会返回。
如果你希望更快地返回,可以为两个 effect 断开中断信号。与其像这样调用:
Effect.raceFirst(task1, task2)
可以改用:
Effect.raceFirst(Effect.disconnect(task1), Effect.disconnect(task2))
这样两个 effect 就能各自独立地完成,同时输掉的那个 effect 仍会在后台被终止。
示例(用 Effect.disconnect 更快返回)
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")),
),
)
// Race the two tasks with disconnect to allow quicker return
const program = Effect.raceFirst(
Effect.disconnect(task1),
Effect.disconnect(task2),
).pipe(Effect.tap(Console.log("more work...")))
Effect.runPromiseExit(program).then(console.log)
/*
Output:
task1 done
more work...
{ _id: 'Exit', _tag: 'Success', value: 'task1' }
task2 interrupted
*/
raceWith
这个函数会并发运行两个 effect,并在其中某个 effect 完成时调用指定的“finisher”(收尾)函数,无论它是成功还是失败。
每个 effect 各自的 finisher 函数让你可以在它们一完成时就去处理各自的结果。
该函数接收两个 finisher 回调,每个 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").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.raceWith(task1, task2, {
onSelfDone: (exit) => Console.log(`task1 exited with ${exit}`),
onOtherDone: (exit) => Console.log(`task2 exited with ${exit}`),
})
Effect.runFork(program)
/*
Output:
task1 done
task1 exited with {
"_id": "Exit",
"_tag": "Success",
"value": "task1"
}
task2 interrupted
*/