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

内置调度

了解 Effect 中的内置调度模式,用于高效地实现定时重复与延迟。

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

#<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"

无限重复与固定次数重复

forever

一个无限重复的调度,每次运行时产出当前的重复次数。

示例(无限重复的调度)

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.forever

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

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

once

一个只重复一次的调度。

示例(只重复一次的调度)

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.duration(Duration.zero)

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

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

recurs

一个重复指定次数的调度,每次运行时产出当前的重复次数。

示例(固定重复次数)

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.recurs(5)

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

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

按固定间隔重复

你可以定义控制”两次执行之间间隔多久”的调度。spacedfixed 的区别在于间隔如何计量

  • spaced:从上一次执行结束开始计算下一次的延迟。
  • fixed:保证以固定节奏重复,不受执行耗时影响。

spaced

一个无限重复的调度,每次重复与上一次运行之间相隔指定的时长。它每次运行时返回重复次数。

示例(执行之间带延迟地重复)

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.spaced("200 millis")

//               ┌─── Simulating an effect that takes
//               │    100 milliseconds to complete
//               ▼
log(schedule, "100 millis")
/*
Output:
#1: 300ms < spaced
#2: 300ms
#3: 300ms
#4: 300ms
#5: 300ms
#6: 300ms
#7: 300ms
#8: 300ms
#9: 300ms
#10: 300ms
...
*/

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

第一次延迟大约是 100 毫秒,因为首次执行不受调度影响。之后的每次延迟大约相隔 200 毫秒,体现了 spaced 调度的效果。

fixed

一个按固定间隔重复的调度,每次运行时返回重复次数。 如果两次更新之间执行的动作耗时超过了间隔,那么该动作会被立即执行,但重复执行不会堆积

|-----interval-----|-----interval-----|-----interval-----|
|---------action--------|action-------|action------------|

示例(固定间隔重复)

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.fixed("200 millis")

//               ┌─── Simulating an effect that takes
//               │    100 milliseconds to complete
//               ▼
log(schedule, "100 millis")
/*
Output:
#1: 300ms < fixed
#2: 200ms
#3: 200ms
#4: 200ms
#5: 200ms
#6: 200ms
#7: 200ms
#8: 200ms
#9: 200ms
#10: 200ms
...
*/

const result = log(schedule, "100 millis")
result.map(Duration.toMillis) // => [300, 200, 200, 200, 200, 200, 200, 200, 200, 200, 200]

逐渐拉长执行间隔

exponential

一个使用指数退避重复的调度,每次延迟按指数增长。返回当前的重复间隔。

示例(指数退避调度)

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.exponential("10 millis")

log(schedule)
/*
Output:
#1: 10ms < exponential
#2: 20ms
#3: 40ms
#4: 80ms
#5: 160ms
#6: 320ms
#7: 640ms
#8: 1280ms
#9: 2560ms
#10: 5120ms
...
*/

const result = log(schedule)
result.map(Duration.toMillis) // => [10, 20, 40, 80, 160, 320, 640, 1280, 2560, 5120, 10240]

fibonacci

一个总是重复的调度,每次延迟等于前两次延迟之和(类似斐波那契数列)。返回当前的重复间隔。

示例(斐波那契延迟调度)

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.fibonacci("10 millis")

log(schedule)
/*
Output:
#1: 10ms
#2: 20ms
#3: 30ms
#4: 50ms
#5: 80ms
#6: 130ms
#7: 210ms
#8: 340ms
#9: 550ms
#10: 890ms
...
*/

const result = log(schedule)
result.map(Duration.toMillis) // => [10, 20, 30, 50, 80, 130, 210, 340, 550, 890, 1440]