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

Schedule 组合子

学习如何在 Effect 中组合与定制 Schedule,构造复杂的重复模式,包括并集、交集、顺序执行等。

Schedule 定义的是有状态、且可能带 effect 的事件重复计划,它们能以多种方式组合。组合子(combinator)让我们可以把多个 Schedule 组合到一起,从而得到新的 Schedule。

为了演示不同 Schedule 的功能,我们将使用下面这个辅助函数,它会记录每一次重复以及对应的延迟毫秒数,格式为:

#<repetition>: <delay in ms>

辅助函数(记录执行延迟)

import { Array, Chunk, Duration, Effect, Schedule } from "effect"

const log = (
  schedule: Schedule.Schedule<unknown>,
  delay: Duration.DurationInput = 0,
): void => {
  const maxRecurs = 10 // Limit the number of executions
  const delays = Chunk.toArray(
    Effect.runSync(
      Schedule.run(
        Schedule.delays(Schedule.addDelay(schedule, () => delay)),
        Date.now(),
        Array.range(0, maxRecurs),
      ),
    ),
  )
  delays.forEach((duration, i) => {
    console.log(
      i === maxRecurs
        ? "..." // Indicate truncation if there are more executions
        : i === delays.length - 1
          ? "(end)" // Mark the last execution
          : `#${i + 1}: ${Duration.toMillis(duration)}ms`,
    )
  })
}

组合方式

Schedule 可以通过不同的方式组合:

模式说明
并集组合两个 Schedule,只要任一 Schedule 想继续就继续,并使用较短的延迟。
交集组合两个 Schedule,仅当两个 Schedule 都想继续时才继续,并使用较长的延迟。
顺序执行组合两个 Schedule,先完整运行第一个,然后切换到第二个。

并集

组合两个 Schedule,只要任一 Schedule 想继续就继续,并使用较短的延迟。

示例(组合指数退避与固定间隔)

import { Array, Chunk, Duration, Effect, Schedule } from "effect"

const log = (
  schedule: Schedule.Schedule<unknown>,
  delay: Duration.DurationInput = 0,
): void => {
  const maxRecurs = 10
  const delays = Chunk.toArray(
    Effect.runSync(
      Schedule.run(
        Schedule.delays(Schedule.addDelay(schedule, () => delay)),
        Date.now(),
        Array.range(0, maxRecurs),
      ),
    ),
  )
  delays.forEach((duration, i) => {
    console.log(
      i === maxRecurs
        ? "..."
        : i === delays.length - 1
          ? "(end)"
          : `#${i + 1}: ${Duration.toMillis(duration)}ms`,
    )
  })
}

const schedule = Schedule.union(
  Schedule.exponential("100 millis"),
  Schedule.spaced("1 second"),
)

log(schedule)
/*
Output:
#1: 100ms  < exponential
#2: 200ms
#3: 400ms
#4: 800ms
#5: 1000ms < spaced
#6: 1000ms
#7: 1000ms
#8: 1000ms
#9: 1000ms
#10: 1000ms
...
*/

Schedule.union 操作符在每一步都挑选最短的延迟,因此在把指数退避的 Schedule 与固定间隔组合时,最初的几次重复会遵循指数退避,而当延迟超过该值后,就会稳定在固定间隔上。

交集

组合两个 Schedule,仅当两个 Schedule 都想继续时才继续,并使用较长的延迟。

示例(用固定重试次数限制指数退避)

import { Array, Chunk, Duration, Effect, Schedule } from "effect"

const log = (
  schedule: Schedule.Schedule<unknown>,
  delay: Duration.DurationInput = 0,
): void => {
  const maxRecurs = 10
  const delays = Chunk.toArray(
    Effect.runSync(
      Schedule.run(
        Schedule.delays(Schedule.addDelay(schedule, () => delay)),
        Date.now(),
        Array.range(0, maxRecurs),
      ),
    ),
  )
  delays.forEach((duration, i) => {
    console.log(
      i === maxRecurs
        ? "..."
        : i === delays.length - 1
          ? "(end)"
          : `#${i + 1}: ${Duration.toMillis(duration)}ms`,
    )
  })
}

const schedule = Schedule.intersect(
  Schedule.exponential("10 millis"),
  Schedule.recurs(5),
)

log(schedule)
/*
Output:
#1: 10ms  < exponential
#2: 20ms
#3: 40ms
#4: 80ms
#5: 160ms
(end)     < recurs
*/

Schedule.intersect 操作符会同时施加两个 Schedule 的约束。在这个示例中,Schedule 遵循指数退避,但由于 Schedule.recurs(5) 的限制,它在 5 次重复之后就会停止。

顺序执行

组合两个 Schedule,先完整运行第一个,然后切换到第二个。

示例(从固定重试次数切换到周期性执行)

import { Array, Chunk, Duration, Effect, Schedule } from "effect"

const log = (
  schedule: Schedule.Schedule<unknown>,
  delay: Duration.DurationInput = 0,
): void => {
  const maxRecurs = 10
  const delays = Chunk.toArray(
    Effect.runSync(
      Schedule.run(
        Schedule.delays(Schedule.addDelay(schedule, () => delay)),
        Date.now(),
        Array.range(0, maxRecurs),
      ),
    ),
  )
  delays.forEach((duration, i) => {
    console.log(
      i === maxRecurs
        ? "..."
        : i === delays.length - 1
          ? "(end)"
          : `#${i + 1}: ${Duration.toMillis(duration)}ms`,
    )
  })
}

const schedule = Schedule.andThen(
  Schedule.recurs(5),
  Schedule.spaced("1 second"),
)

log(schedule)
/*
Output:
#1: 0ms    < recurs
#2: 0ms
#3: 0ms
#4: 0ms
#5: 0ms
#6: 1000ms < spaced
#7: 1000ms
#8: 1000ms
#9: 1000ms
#10: 1000ms
...
*/

第一个 Schedule 会一直运行到结束,之后由第二个 Schedule 接管。在这个示例中,effect 最初会毫无延迟地执行 5 次,然后每隔 1 秒继续执行。

为重试延迟加入随机性

Schedule.jittered 组合子会通过在指定范围内施加随机延迟来修改一个 Schedule。

当资源因过载或争用而不可用时,重试与退避并不能帮到我们。如果所有失败的 API 调用都在同一时间点退避,它们会造成新一轮的过载或争用。抖动(jitter)为 Schedule 的延迟加入了一定程度的随机性,这有助于避免意外地让它们同步在一起,从而意外拖垮服务。

研究表明,Schedule.jittered(0.0, 1.0) 是在重试中引入随机性的一种有效方式。

示例(加入抖动的指数退避)

import { Array, Chunk, Duration, Effect, Schedule } from "effect"

const log = (
  schedule: Schedule.Schedule<unknown>,
  delay: Duration.DurationInput = 0,
): void => {
  const maxRecurs = 10
  const delays = Chunk.toArray(
    Effect.runSync(
      Schedule.run(
        Schedule.delays(Schedule.addDelay(schedule, () => delay)),
        Date.now(),
        Array.range(0, maxRecurs),
      ),
    ),
  )
  delays.forEach((duration, i) => {
    console.log(
      i === maxRecurs
        ? "..."
        : i === delays.length - 1
          ? "(end)"
          : `#${i + 1}: ${Duration.toMillis(duration)}ms`,
    )
  })
}

const schedule = Schedule.jittered(Schedule.exponential("10 millis"))

log(schedule)
/*
Output:
#1: 10.448486ms
#2: 21.134521ms
#3: 47.245117ms
#4: 88.263184ms
#5: 163.651367ms
#6: 335.818848ms
#7: 719.126709ms
#8: 1266.18457ms
#9: 2931.252441ms
#10: 6121.593018ms
...
*/

Schedule.jittered 组合子在某个范围内为延迟引入随机性。例如,对指数退避施加抖动可以确保每次重试都发生在略微不同的时间点,从而降低压垮系统的风险。

用过滤器控制重复次数

你可以使用 Schedule.whileInputSchedule.whileOutput,根据施加在 Schedule 输入或输出上的条件来限制它持续多久。

示例(根据输出停止)

import { Array, Chunk, Duration, Effect, Schedule } from "effect"

const log = (
  schedule: Schedule.Schedule<unknown>,
  delay: Duration.DurationInput = 0,
): void => {
  const maxRecurs = 10
  const delays = Chunk.toArray(
    Effect.runSync(
      Schedule.run(
        Schedule.delays(Schedule.addDelay(schedule, () => delay)),
        Date.now(),
        Array.range(0, maxRecurs),
      ),
    ),
  )
  delays.forEach((duration, i) => {
    console.log(
      i === maxRecurs
        ? "..."
        : i === delays.length - 1
          ? "(end)"
          : `#${i + 1}: ${Duration.toMillis(duration)}ms`,
    )
  })
}

const schedule = Schedule.whileOutput(Schedule.recurs(5), (n) => n <= 2)

log(schedule)
/*
Output:
#1: 0ms < recurs
#2: 0ms
#3: 0ms
(end)   < whileOutput
*/

Schedule.whileOutput 会根据 Schedule 的输出来过滤重复次数。在这个示例中,一旦输出超过 2,Schedule 就会停止,尽管 Schedule.recurs(5) 允许最多 5 次重复。

根据输出调整延迟

Schedule.modifyDelay 组合子允许你根据重复次数或其他输出条件,动态地改变 Schedule 的延迟。

示例(在若干次重复之后缩短延迟)

import { Array, Chunk, Duration, Effect, Schedule } from "effect"

const log = (
  schedule: Schedule.Schedule<unknown>,
  delay: Duration.DurationInput = 0,
): void => {
  const maxRecurs = 10
  const delays = Chunk.toArray(
    Effect.runSync(
      Schedule.run(
        Schedule.delays(Schedule.addDelay(schedule, () => delay)),
        Date.now(),
        Array.range(0, maxRecurs),
      ),
    ),
  )
  delays.forEach((duration, i) => {
    console.log(
      i === maxRecurs
        ? "..."
        : i === delays.length - 1
          ? "(end)"
          : `#${i + 1}: ${Duration.toMillis(duration)}ms`,
    )
  })
}

const schedule = Schedule.modifyDelay(
  Schedule.spaced("1 second"),
  (out, duration) => (out > 2 ? "100 millis" : duration),
)

log(schedule)
/*
Output:
#1: 1000ms
#2: 1000ms
#3: 1000ms
#4: 100ms  < modifyDelay
#5: 100ms
#6: 100ms
#7: 100ms
#8: 100ms
#9: 100ms
#10: 100ms
...
*/

延迟的修改会在执行过程中动态生效。在这个示例中,前 3 次重复遵循原本的 1 秒 间隔;此后延迟降到 100 毫秒,使得后续的重复更加频繁。

探查

Schedule.tapInputSchedule.tapOutput 允许你在不改变 Schedule 行为的前提下,对它的输入或输出执行额外的带 effect 操作。

示例(记录 Schedule 的输出)

import { Array, Chunk, Duration, Effect, Schedule, Console } from "effect"

const log = (
  schedule: Schedule.Schedule<unknown>,
  delay: Duration.DurationInput = 0,
): void => {
  const maxRecurs = 10
  const delays = Chunk.toArray(
    Effect.runSync(
      Schedule.run(
        Schedule.delays(Schedule.addDelay(schedule, () => delay)),
        Date.now(),
        Array.range(0, maxRecurs),
      ),
    ),
  )
  delays.forEach((duration, i) => {
    console.log(
      i === maxRecurs
        ? "..."
        : i === delays.length - 1
          ? "(end)"
          : `#${i + 1}: ${Duration.toMillis(duration)}ms`,
    )
  })
}

const schedule = Schedule.tapOutput(Schedule.recurs(2), (n) =>
  Console.log(`Schedule Output: ${n}`),
)

log(schedule)
/*
Output:
Schedule Output: 0
Schedule Output: 1
Schedule Output: 2
#1: 0ms
#2: 0ms
(end)
*/

Schedule.tapOutput 会在每次重复之前运行一个 effect,并以 Schedule 当前的输出作为输入。这对于日志记录、调试或触发副作用都很有用。