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.whileInput 或 Schedule.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.tapInput 和 Schedule.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 当前的输出作为输入。这对于日志记录、调试或触发副作用都很有用。