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

FS2(Functional Streams for Scala)如何实现group(n)?求类Akka Streams的grouped方法

在FS2中实现类似Akka Streams grouped的分组功能

我懂你想要的效果——就像Akka Streams里的grouped(n)那样,把流中的元素按固定数量n打包成组,哪怕最后一组元素不足n也会完整输出。你试过的chunkLimit、segmentLimit这些方法要么受限于chunk的原有结构,要么终止条件不符合预期,确实容易踩坑。咱们来搞定这个需求:

官方推荐方案:使用chunkN方法

FS2其实已经内置了满足这个需求的方法——chunkN,你只需要注意它的第二个参数allowFewer:

import cats.effect.IO
import fs2.Stream

// 创建测试流
val numberStream = Stream(1, 2, 3, 4, 5, 6, 7)

// 按每组3个元素分组,允许最后一组不足3个
val groupedStream = numberStream.chunkN(3, allowFewer = true)
  .map(_.toList) // 把Chunk转换成List(可选,根据你的需求调整)

// 运行流并输出结果
groupedStream.compile.toList.flatMap { groups =>
  IO(groups.foreach(println))
}.unsafeRunSync()

这段代码的输出会和你预期的一致:

List(1, 2, 3)
List(4, 5, 6)
List(7)

为什么之前的方法不适用?

  • chunkLimit(n):只会在输入的单个chunk大小达到n时才输出,不会合并多个小chunk来凑够n个元素,所以如果流中的元素本来就是小chunk,会直接输出而不分组。
  • segmentLimit(n):和chunkLimit类似,基于segment的边界来处理,同样不会跨segment合并元素。
  • groupAdjacentBy:是按自定义条件对相邻元素分组,不是按固定数量分组,完全不符合你的需求。

自定义实现(如果需要更灵活的控制)

如果你需要对分组逻辑做额外定制,可以自己实现一个Pipe:

import cats.effect.IO
import fs2.{Pipe, Stream}

def groupedN[F[_], A](n: Int, allowFewer: Boolean = true): Pipe[F, A, List[A]] =
  stream => {
    // 先把整个流合并成一个大chunk,再按n分组
    stream.chunkAll.flatMap { fullChunk =>
      val groups = fullChunk.grouped(n).map(_.toList).toList
      if (allowFewer) Stream.emits(groups)
      else {
        // 如果不允许不足n的组,就过滤掉最后一组(如果它的大小不够)
        val validGroups = if (groups.lastOption.exists(_.size < n)) groups.init else groups
        Stream.emits(validGroups)
      }
    }
  }

// 使用自定义Pipe
val customGrouped = numberStream.through(groupedN(3)).compile.toList.unsafeRunSync()
println(customGrouped) // 输出:List(List(1,2,3), List(4,5,6), List(7))

这个自定义实现先把整个流收集成一个完整的chunk,再按指定数量拆分,适合需要一次性处理所有元素的场景。如果你的流是无限流,这种方式就不适用了,这时候还是用chunkN更合适,因为它是增量处理的,不会积压元素。

内容的提问来源于stack exchange,提问作者iuriisusuk

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 08:48:22