内置调度
了解 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]
按固定间隔重复
你可以定义控制”两次执行之间间隔多久”的调度。spaced 与 fixed 的区别在于间隔如何计量:
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]