Sink 并发
了解如何通过并发的 sink 操作提升性能,例如合并结果,或竞速以捕获最先完成者。
本节介绍并发操作,它们允许多个 sink 同时运行。当你确实需要并发执行时,这些操作对于提升任务性能很有价值。
通过并发 zip 合并结果
要并发运行两个 sink 并合并它们的结果,可以使用 Sink.zip。该操作会并发执行两个 sink,并把它们的结果合并成一个元组。
示例(并发运行两个 Sink 并合并结果)
import { Sink, Console, Stream, Schedule, Effect } from "effect"
const stream = Stream.make("1", "2", "3", "4", "5").pipe(
Stream.schedule(Schedule.spaced("10 millis")),
)
const sink1 = Sink.forEach((s: string) => Console.log(`sink 1: ${s}`)).pipe(
Sink.as(1),
)
const sink2 = Sink.forEach((s: string) => Console.log(`sink 2: ${s}`)).pipe(
Sink.as(2),
)
// Combine the two sinks to run concurrently and collect results in a tuple
const sink = Sink.zip(sink1, sink2, { concurrent: true })
Effect.runPromise(Stream.run(stream, sink)).then(console.log)
/*
Output:
sink 1: 1
sink 2: 1
sink 1: 2
sink 2: 2
sink 1: 3
sink 2: 3
sink 1: 4
sink 2: 4
sink 1: 5
sink 2: 5
[ 1, 2 ]
*/
竞速 Sink:最先完成者胜出
Sink.race 操作允许多个 sink 竞争完成。最先完成的那个 sink 提供结果。
示例(让两个 Sink 竞速以捕获最先产生的结果)
import { Sink, Console, Stream, Schedule, Effect } from "effect"
const stream = Stream.make("1", "2", "3", "4", "5").pipe(
Stream.schedule(Schedule.spaced("10 millis")),
)
const sink1 = Sink.forEach((s: string) => Console.log(`sink 1: ${s}`)).pipe(
Sink.as(1),
)
const sink2 = Sink.forEach((s: string) => Console.log(`sink 2: ${s}`)).pipe(
Sink.as(2),
)
// Race the two sinks, the result will be from the first to complete
const sink = Sink.race(sink1, sink2)
Effect.runPromise(Stream.run(stream, sink)).then(console.log)
/*
Output:
sink 1: 1
sink 2: 1
sink 1: 2
sink 2: 2
sink 1: 3
sink 2: 3
sink 1: 4
sink 2: 4
sink 1: 5
sink 2: 5
1
*/