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

FS2中如何按终止类型为输入流添加对应终止元素

解决方案

核心实现思路

利用Deferred捕获流的终止状态(正常完成、取消、错误),将原始ByteBuffer流与收尾元素流合并,确保收尾元素能在流终止时正确输出,同时避免多余的Data元素发送。

具体代码实现

import fs2._
import cats.effect._
import java.nio.ByteBuffer

sealed trait Input
case class Data(x: ByteBuffer) extends Input
case object EndOfInput extends Input
case object Cancel extends Input

def byteBufferToInput[F[_]: Concurrent](source: Stream[F, ByteBuffer]): Stream[F, Input] = {
  // 创建Deferred用于传递流的终止信号
  Stream.eval(Deferred[F, ExitCase[Throwable]]).flatMap { exitSignal =>
    // 原始ByteBuffer流映射为Data元素,终止时将状态写入Deferred
    val dataStream = source
      .map(Data(_))
      .onFinalizeCase(exitCase => exitSignal.complete(exitCase).void)

    // 根据终止状态生成收尾元素流
    val finalizerStream = Stream
      .eval(exitSignal.get)
      .flatMap {
        case ExitCase.Completed => Stream.emit(EndOfInput)
        case ExitCase.Canceled => Stream.emit(Cancel)
        case ExitCase.Error(_) => 
          // 此处可根据需求调整:如需特定错误下发送Cancel,可添加错误类型判断
          Stream.empty
      }

    // 合并两个流,任一终止则停止另一流,保证收尾元素仅在原始流终止后输出
    dataStream.mergeHaltBoth(finalizerStream)
  }
}

自定义扩展方法(实现理想中的latelyCase)

如果需要复用这种"根据终止状态追加元素"的逻辑,可以封装成流的扩展方法:

implicit class StreamLatelyCaseOps[F[_], O](private val stream: Stream[F, O]) extends AnyVal {
  def latelyCase[F2[x] >: F[x]: Concurrent](f: ExitCase[Throwable] => Stream[F2, O]): Stream[F2, O] = {
    Stream.eval(Deferred[F2, ExitCase[Throwable]]).flatMap { exitSignal =>
      val main = stream.onFinalizeCase(exitSignal.complete(_).void)
      val finalStream = Stream.eval(exitSignal.get).flatMap(f)
      main.mergeHaltBoth(finalStream)
    }
  }
}

使用时可以简化为:

val inputStream: Stream[IO, Input] = byteBufferSource
  .map(Data(_))
  .latelyCase {
    case ExitCase.Completed => Stream.emit(EndOfInput)
    case ExitCase.Canceled => Stream.emit(Cancel)
    case ExitCase.Error(e) => 
      // 示例:仅在自定义错误时发送Cancel
      if (e.isInstanceOf[MyCustomError]) Stream.emit(Cancel) else Stream.empty
  }

关键细节说明

  1. onFinalizeCase的作用:会在流终止时触发,无论终止原因是正常完成、取消还是错误,确保Deferred能正确捕获终止状态,解决之前Deferred在complete前被取消的问题。
  2. mergeHaltBoth的作用:保证原始流一旦终止,收尾流立即输出对应元素,同时不会有后续的Data元素被发送。
  3. 错误处理灵活性:在ExitCase.Error分支可根据错误类型自定义逻辑,避免第三方服务返回的错误误发Cancel。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 00:00:58