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

寻求FS2 Stream的chunkBy算子:按谓词分组为块的实现方案

在FS2中实现chunkBy算子

FS2原生并没有提供直接的chunkBy算子,但你可以通过组合现有API来实现按指定谓词将流元素分组为连续块的功能。

实现代码

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

def chunkBy[F[_], A, K](stream: Stream[F, A])(keyFn: A => K): Stream[F, Chunk[A]] = {
  stream
    // 扫描流,为每个元素带上其分组键,保留前一个元素的键信息用于比较
    .scan(Option.empty[(K, A)]) {
      case (None, elem) => Some((keyFn(elem), elem))
      case (_, elem) => Some((keyFn(elem), elem))
    }
    .drop(1) // 移除初始的空值
    // 当前元素键与前一个不同时,分割流
    .splitWhen { case (currentKey, _) =>
      _.exists { case (prevKey, _) => prevKey != currentKey }
    }
    // 将每个分割后的子流转换为Chunk
    .map(grouped => grouped.map(_._2).compile.toChunk)
    .flatten // 将F[Chunk[A]]转换为Stream[F, Chunk[A]]
}

测试示例

针对你提到的1-8的IO流按元素对4取余分组的场景:

val numberStream = Stream.emits(1 to 8).covary[IO]
val groupedChunks = chunkBy(numberStream)(_ % 4).compile.toList.unsafeRunSync()

// 预期输出:List(Chunk(1,5), Chunk(2,6), Chunk(3,7), Chunk(4,8))

注意事项

  • 这个实现是连续分组语义:只有当元素的谓词结果与前一个元素不同时,才会开启新的块。如果需要将所有谓词结果相同的元素(无论是否连续)归为同一块,那属于groupBy的语义,和chunkBy的连续分组逻辑不同。
  • 代码通过scan跟踪每个元素的分组键,splitWhen检测分组键的变化点,最后将每个子流编译为Chunk并展开,完整实现了所需的分组功能。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 13:29:58