You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.26 13:54:01