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

调度组合子

学习如何通过组合与定制 Effect 中的调度,构建复杂的重复模式,包括并集、交集、顺序连接等。

调度(Schedule)定义的是有状态、可能带副作用的事件重复计划,并且可以以多种方式组合。组合子(combinator)让我们可以把多个调度组合在一起,得到新的调度。

为了演示不同调度的行为,我们会用到下面这个辅助函数:它把每一次重复连同对应的延迟(毫秒)一起打印出来,格式为:

#<repetition>: <delay in ms>

辅助函数(打印执行延迟)

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

const log = (
  schedule: Schedule.Schedule<unknown>,
  delay: Duration.Input = 0,
): void => {
  const maxRecurs = 10 // Limit the number of executions
  const withDelay = Schedule.addDelay(schedule, () => Effect.succeed(delay))
  const delays = Effect.runSync(
    Effect.gen(function* () {
      const step = yield* Schedule.toStep(withDelay)
      const out: Array<Duration.Duration> = []
      let now = Date.now()
      for (const input of Array.range(0, maxRecurs)) {
        const duration = yield* Pull.matchEffect(step(now, input), {
          onSuccess: ([, duration]) => Effect.succeed(duration),
          onFailure: Effect.failCause,
          onDone: () => Effect.succeed(undefined),
        })
        if (duration === undefined) break
        out.push(duration)
        now += Duration.toMillis(duration)
      }
      return out
    }),
  )
  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`,
    )
  })
}

typeof log // => "function"

组合

调度可以通过不同方式组合:

模式说明
并集(Union)组合两个调度,只要其中一个还想继续就重复,并取较短的延迟。
交集(Intersection)组合两个调度,只有两个都想继续才重复,并取较长的延迟。
顺序连接(Sequencing)先完整跑完第一个调度,再切换到第二个。

并集(Union)

组合两个调度,只要其中一个还想继续就重复,并取较短的延迟。

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

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

const log = (
  schedule: Schedule.Schedule<unknown>,
  delay: Duration.Input = 0,
): Array<Duration.Duration> => {
  const maxRecurs = 10
  const withDelay = Schedule.addDelay(schedule, () => Effect.succeed(delay))
  const delays = Effect.runSync(
    Effect.gen(function* () {
      const step = yield* Schedule.toStep(withDelay)
      const out: Array<Duration.Duration> = []
      let now = Date.now()
      for (const input of Array.range(0, maxRecurs)) {
        const duration = yield* Pull.matchEffect(step(now, input), {
          onSuccess: ([, duration]) => Effect.succeed(duration),
          onFailure: Effect.failCause,
          onDone: () => Effect.succeed(undefined),
        })
        if (duration === undefined) break
        out.push(duration)
        now += Duration.toMillis(duration)
      }
      return out
    }),
  )
  delays.forEach((duration, i) => {
    console.log(
      i === maxRecurs
        ? "..."
        : i === delays.length - 1
          ? "(end)"
          : `#${i + 1}: ${Duration.toMillis(duration)}ms`,
    )
  })
  return delays
}

const schedule = Schedule.min([
  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
...
*/

const result = log(schedule)
result.map(Duration.toMillis) // => [100, 200, 400, 800, 1000, 1000, 1000, 1000, 1000, 1000, 1000]

Schedule.min 运算符在每一步都取最短的延迟,因此把指数退避和固定间隔组合时,初始的重复会走指数退避,等到延迟超过那个固定值后就稳定成固定间隔。

交集(Intersection)

组合两个调度,只有两个都想继续才重复,并取较长的延迟。

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

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

const log = (
  schedule: Schedule.Schedule<unknown>,
  delay: Duration.Input = 0,
): Array<Duration.Duration> => {
  const maxRecurs = 10
  const withDelay = Schedule.addDelay(schedule, () => Effect.succeed(delay))
  const delays = Effect.runSync(
    Effect.gen(function* () {
      const step = yield* Schedule.toStep(withDelay)
      const out: Array<Duration.Duration> = []
      let now = Date.now()
      for (const input of Array.range(0, maxRecurs)) {
        const duration = yield* Pull.matchEffect(step(now, input), {
          onSuccess: ([, duration]) => Effect.succeed(duration),
          onFailure: Effect.failCause,
          onDone: () => Effect.succeed(undefined),
        })
        if (duration === undefined) break
        out.push(duration)
        now += Duration.toMillis(duration)
      }
      return out
    }),
  )
  delays.forEach((duration, i) => {
    console.log(
      i === maxRecurs
        ? "..."
        : i === delays.length - 1
          ? "(end)"
          : `#${i + 1}: ${Duration.toMillis(duration)}ms`,
    )
  })
  return delays
}

const schedule = Schedule.max([
  Schedule.exponential("10 millis"),
  Schedule.recurs(5),
])

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

const result = log(schedule)
result.map(Duration.toMillis) // => [10, 20, 40, 80, 160]

Schedule.max 运算符会同时受两个调度的约束。在这个例子中,调度走的是指数退避,但因为 Schedule.recurs(5) 的限制,在第 5 次重复后停止。

顺序连接(Sequencing)

组合两个调度,先完整跑完第一个,再切换到第二个。

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

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

const log = (
  schedule: Schedule.Schedule<unknown>,
  delay: Duration.Input = 0,
): Array<Duration.Duration> => {
  const maxRecurs = 10
  const withDelay = Schedule.addDelay(schedule, () => Effect.succeed(delay))
  const delays = Effect.runSync(
    Effect.gen(function* () {
      const step = yield* Schedule.toStep(withDelay)
      const out: Array<Duration.Duration> = []
      let now = Date.now()
      for (const input of Array.range(0, maxRecurs)) {
        const duration = yield* Pull.matchEffect(step(now, input), {
          onSuccess: ([, duration]) => Effect.succeed(duration),
          onFailure: Effect.failCause,
          onDone: () => Effect.succeed(undefined),
        })
        if (duration === undefined) break
        out.push(duration)
        now += Duration.toMillis(duration)
      }
      return out
    }),
  )
  delays.forEach((duration, i) => {
    console.log(
      i === maxRecurs
        ? "..."
        : i === delays.length - 1
          ? "(end)"
          : `#${i + 1}: ${Duration.toMillis(duration)}ms`,
    )
  })
  return delays
}

const schedule = Schedule.concat(
  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
...
*/

const result = log(schedule)
result.map(Duration.toMillis) // => [0, 0, 0, 0, 0, 1000, 1000, 1000, 1000, 1000, 1000]

第一个调度先跑到结束,之后由第二个调度接管。在这个例子中,effect 一开始连续执行 5 次(无延迟),然后每 1 秒执行一次。

给重试延迟加入随机性

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

当某个资源因为过载或竞争而不可用时,重试加退避并不能帮上忙。如果所有失败的 API 调用都被退避到同一个时间点,它们会再次造成过载或竞争。Jitter(抖动)会给调度的延迟加入一定量的随机性,这样我们就不会在无意中把请求同步化、从而意外地把服务压垮。

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

示例(带抖动的指数退避)

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

const log = (
  schedule: Schedule.Schedule<unknown>,
  delay: Duration.Input = 0,
): Array<Duration.Duration> => {
  const maxRecurs = 10
  const withDelay = Schedule.addDelay(schedule, () => Effect.succeed(delay))
  const delays = Effect.runSync(
    Effect.gen(function* () {
      const step = yield* Schedule.toStep(withDelay)
      const out: Array<Duration.Duration> = []
      let now = Date.now()
      for (const input of Array.range(0, maxRecurs)) {
        const duration = yield* Pull.matchEffect(step(now, input), {
          onSuccess: ([, duration]) => Effect.succeed(duration),
          onFailure: Effect.failCause,
          onDone: () => Effect.succeed(undefined),
        })
        if (duration === undefined) break
        out.push(duration)
        now += Duration.toMillis(duration)
      }
      return out
    }),
  )
  delays.forEach((duration, i) => {
    console.log(
      i === maxRecurs
        ? "..."
        : i === delays.length - 1
          ? "(end)"
          : `#${i + 1}: ${Duration.toMillis(duration)}ms`,
    )
  })
  return delays
}

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

// The exact delays are randomized by jitter, but the schedule still
// produces exactly 11 values within maxRecurs (10 shown + 1 truncated)
const result = log(schedule)
result.length // => 11

Schedule.jittered 组合子会在指定范围内给延迟加入随机性。例如,对指数退避施加抖动,能保证每次重试发生在略微不同的时刻,从而降低压垮系统的风险。

用过滤器控制重复次数

使用 Schedule.while 可以限制调度持续多久。它的谓词会收到一段元数据,里面包含调度的输入和输出。

示例(基于输出停止)

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

const log = (
  schedule: Schedule.Schedule<unknown>,
  delay: Duration.Input = 0,
): Array<Duration.Duration> => {
  const maxRecurs = 10
  const withDelay = Schedule.addDelay(schedule, () => Effect.succeed(delay))
  const delays = Effect.runSync(
    Effect.gen(function* () {
      const step = yield* Schedule.toStep(withDelay)
      const out: Array<Duration.Duration> = []
      let now = Date.now()
      for (const input of Array.range(0, maxRecurs)) {
        const duration = yield* Pull.matchEffect(step(now, input), {
          onSuccess: ([, duration]) => Effect.succeed(duration),
          onFailure: Effect.failCause,
          onDone: () => Effect.succeed(undefined),
        })
        if (duration === undefined) break
        out.push(duration)
        now += Duration.toMillis(duration)
      }
      return out
    }),
  )
  delays.forEach((duration, i) => {
    console.log(
      i === maxRecurs
        ? "..."
        : i === delays.length - 1
          ? "(end)"
          : `#${i + 1}: ${Duration.toMillis(duration)}ms`,
    )
  })
  return delays
}

const schedule = Schedule.while(Schedule.recurs(5), ({ output }) => output <= 2)

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

const result = log(schedule)
result.map(Duration.toMillis) // => [0, 0, 0]

Schedule.while 根据调度的输出过滤重复。在这个例子中,即便 Schedule.recurs(5) 允许最多 5 次重复,一旦输出超过 2,调度就会停止。

根据输出调整延迟

Schedule.modifyDelay 组合子让你能基于重复次数或其它输出条件,动态地改变调度的延迟。

示例(在达到一定重复次数后缩短延迟)

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

const log = (
  schedule: Schedule.Schedule<unknown>,
  delay: Duration.Input = 0,
): Array<Duration.Duration> => {
  const maxRecurs = 10
  const withDelay = Schedule.addDelay(schedule, () => Effect.succeed(delay))
  const delays = Effect.runSync(
    Effect.gen(function* () {
      const step = yield* Schedule.toStep(withDelay)
      const out: Array<Duration.Duration> = []
      let now = Date.now()
      for (const input of Array.range(0, maxRecurs)) {
        const duration = yield* Pull.matchEffect(step(now, input), {
          onSuccess: ([, duration]) => Effect.succeed(duration),
          onFailure: Effect.failCause,
          onDone: () => Effect.succeed(undefined),
        })
        if (duration === undefined) break
        out.push(duration)
        now += Duration.toMillis(duration)
      }
      return out
    }),
  )
  delays.forEach((duration, i) => {
    console.log(
      i === maxRecurs
        ? "..."
        : i === delays.length - 1
          ? "(end)"
          : `#${i + 1}: ${Duration.toMillis(duration)}ms`,
    )
  })
  return delays
}

const schedule = Schedule.modifyDelay(
  Schedule.spaced("1 second"),
  ({ output, duration }) =>
    Effect.succeed(output > 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
...
*/

const result = log(schedule)
result.map(Duration.toMillis) // => [1000, 1000, 1000, 100, 100, 100, 100, 100, 100, 100, 100]

延迟的修改在运行时动态生效。在这个例子中,前三次重复沿用原始的 1 秒 间隔;之后延迟降到 100 毫秒,使后续重复更加频繁。

旁路(Tapping)

Schedule.tap 会执行一个额外的、带副作用的操作,而不改变调度的行为。它的回调会收到一段元数据,里面包含调度的输入和输出。

示例(记录调度输出)

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

const log = (
  schedule: Schedule.Schedule<unknown>,
  delay: Duration.Input = 0,
): Array<Duration.Duration> => {
  const maxRecurs = 10
  const withDelay = Schedule.addDelay(schedule, () => Effect.succeed(delay))
  const delays = Effect.runSync(
    Effect.gen(function* () {
      const step = yield* Schedule.toStep(withDelay)
      const out: Array<Duration.Duration> = []
      let now = Date.now()
      for (const input of Array.range(0, maxRecurs)) {
        const duration = yield* Pull.matchEffect(step(now, input), {
          onSuccess: ([, duration]) => Effect.succeed(duration),
          onFailure: Effect.failCause,
          onDone: () => Effect.succeed(undefined),
        })
        if (duration === undefined) break
        out.push(duration)
        now += Duration.toMillis(duration)
      }
      return out
    }),
  )
  delays.forEach((duration, i) => {
    console.log(
      i === maxRecurs
        ? "..."
        : i === delays.length - 1
          ? "(end)"
          : `#${i + 1}: ${Duration.toMillis(duration)}ms`,
    )
  })
  return delays
}

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

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

const result = log(schedule)
result.map(Duration.toMillis) // => [0, 0]

Schedule.tap 会在每次重复之前运行一个 effect,并以调度当前的输出作为输入。它可以用于打日志、调试,或触发副作用。