FS2未知鉴别符条目实现同组保序跨组并行处理的方法
实现同鉴别符顺序处理、不同鉴别符并行处理的方案
核心思路:借助 FS2 内置的 groupBy 算子即可实现动态鉴别符的需求,无需提前预知所有鉴别符取值,完整可运行实现如下:
package org.example import cats.effect.std.Random import cats.effect.{ExitCode, IO, IOApp, Temporal} import cats.syntax.all._ import cats.{Applicative, Monad} import fs2._ import scala.concurrent.duration._ object GitterQuestion extends IOApp { override def run(args: List[String]): IO[ExitCode] = Random.scalaUtilRandom[IO].flatMap { implicit random => val flat = Stream( ("a", 1), ("a", 2), ("a", 3), ("b", 1), ("b", 2), ("b", 3), ("c", 1), ("c", 2), ("c", 3) ).covary[IO] // 核心逻辑替换原有硬编码部分 flat .groupBy(_._1) // 按鉴别符动态分组 .map { case (_, keyStream) => keyStream.through(rndDelay) // 每个同key子流串行处理,保证同key顺序 } .parJoin(100) // 所有子流并行合并,不同key互不影响 .printlns .compile.drain.as(ExitCode.Success) } def rndDelay[F[_]: Monad: Random: Temporal, A]: Pipe[F, A, A] = in => in.evalMap { v => (Random[F].nextDouble.map(_.seconds) >>= Temporal[F].sleep) >> Applicative[F].pure(v) } }
逻辑说明
groupBy会自动按流中元素的鉴别符动态创建子流,同一个鉴别符的所有元素会按原流中的输入顺序进入对应子流,不会乱序- 每个子流内部的处理是串行的,所以同鉴别符的条目会严格按输入顺序执行,不会出现重排,完全满足
(a,2)永远不会在(a,1)之前打印的要求 - 所有子流通过
parJoin合并并行执行,不同鉴别符的处理互不影响,和硬编码的效果完全一致,输出示例如下:
(c,1) (a,1) (c,2) (a,2) (c,3) (b,1) (a,3) (b,2) (b,3)
注意事项
如果你的业务场景中鉴别符数量较多,可以调整 groupBy 的 maxOpen 参数,默认值为256,超过这个数量的分组会触发背压,你可以根据实际情况调大,比如:groupBy(_._1, maxOpen = 2000);同时建议parJoin的并行数不小于groupBy的maxOpen值,避免子流处理排队。
内容的提问来源于stack exchange,提问作者Tomasz Bekas
相关产品推荐
相关产品推荐

