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

基础并发

通过并发、中断与竞速来管理和控制 effect 的执行。

并发选项

Effect 提供了一些选项来管理 effect 的执行方式,尤其侧重控制有多少 effect 并发运行。

type Options = {
  readonly concurrency?: Concurrency
}

concurrency 选项用于确定并发级别,取值如下:

type Concurrency = number | "unbounded" | "inherit"

下面我们详细探讨每一种配置。

Applicability of Concurrency Options

这里的示例使用的是 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
*/