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

基础并发

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

并发选项

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

type Options = {
  readonly concurrency?: Concurrency
}

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

type Concurrency = number | "unbounded"

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

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.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
*/