调度组合子
学习如何通过组合与定制 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,并以调度当前的输出作为输入。它可以用于打日志、调试,或触发副作用。