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
相关产品推荐
相关产品推荐

