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

使用fs2预读Doobie查询结果捕获错误遇ResultSet关闭问题

解决Doobie流与http4s流式响应的异常捕获问题

你的核心矛盾是:流式输出Doobie查询结果时,查询格式错误等异常只会在流消费阶段抛出,但此时http4s已经发送了200响应头,无法再返回错误状态码。而你写的peekFirstElement函数虽然试图提前捕获异常,但破坏了Doobie流的资源生命周期,导致结果集提前关闭,后续读取失败。

问题根源

peekFirstElement的实现通过pull.uncons拆分原流,再将第一个chunk和剩余流重新组合后返回。但Doobie的流与数据库结果集强绑定,拆分操作会触发Doobie内部的资源状态变更,当你重新消费组合后的流时,原结果集已经被提前关闭,自然会抛出“结果集已关闭”的异常。

解决思路

不要将流的拆分与重组拆分为独立的IO操作,而是在同一个流上下文内完成“头部异常捕获”和“完整流式输出”,确保Doobie的资源生命周期不被中断。具体实现是通过pull操作一次性获取第一个chunk和剩余流,捕获异常后直接返回完整的安全流,或者返回错误信息。

修改后的代码

1. 替换peekFirstElement为checkStreamHead

/** 检查流的头部,捕获首次读取时的异常,返回安全的流或错误信息 */
def checkStreamHead[O](stream: Stream[IO, O]): IO[Either[Throwable, Stream[IO, O]]] = {
  // 使用pull操作获取第一个chunk和剩余流,避免重复消费流
  stream.pull.uncons
    .attempt // 捕获pull操作中的异常
    .map {
      case Left(error) => Left(error)
      case Right(Some((firstChunk, remainingStream))) =>
        // 将第一个chunk和剩余流重新组合成完整的流
        Right(Stream.chunk(firstChunk) ++ remainingStream)
      case Right(None) =>
        // 流为空,返回空流
        Right(Stream.empty)
    }
    .stream
    .compile
    .onlyOrError // 提取Pull操作的结果
}

2. 调整路由中的响应逻辑

def routes(conns: List[(Regex, Transactor[IO])]): HttpRoutes[IO] = {
  val config = FileService.Config[IO](absolutePath)
  val pathCollector: FileService.Fs2PathCollector[IO] = (f, cfg, request) => {
    conns.find(_._1.matches(request.pathInfo.toString))
      .map { case (_, transactor) =>
        for {
          content <- config.fs2PathCollector(f, cfg, request)
          queryBytes <- OptionT.liftF(content.body.compile.toList)
          query = new String(queryBytes.toArray)
          _ <- OptionT.liftF(logger[IO].debug(s"executing query: \"$query\""))
          outerQuery = sql"select to_json(querydata) from (${Fragment.const(query)}) as querydata".query[Json]
          queryResultStream = outerQuery.stream.transact(transactor)
          // 先检查流的头部,捕获异常
          checkedStreamResult <- OptionT.liftF(checkStreamHead(queryResultStream))
          response <- OptionT.liftF(
            checkedStreamResult match {
              case Left(error) =>
                // 异常发生在响应头发送前,返回500错误
                InternalServerError(s"Query execution failed: ${error.getMessage}")
              case Right(safeStream) =>
                // 构建流式JSON响应
                val jsonStream = Stream.emit(Token.StartArray) ++ safeStream.through(tokenize) ++ Stream(Token.EndArray)
                Ok(jsonStream.through(fs2.data.json.render.compact).through(fs2.text.utf8.encode))
            }
          )
        } yield response
      }.getOrElse(config.fs2PathCollector(f, cfg, request))
  }

  fileService[IO](FileService.Config[IO](absolutePath, pathCollector, "", NoopCacheStrategy[IO], 32 * 1024))
}

关键说明

  • checkStreamHead通过pull.uncons在同一个操作中获取第一个chunk和剩余流,确保Doobie的结果集资源不会被提前关闭。
  • 异常捕获发生在构建http4s响应之前,此时还未发送200响应头,因此可以正常返回500错误状态码。
  • 重组后的流保持了原Doobie流的完整性,流式输出时不会出现资源已关闭的问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 19:39:52