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

如何用Scala的ZIO、ZStream定时流式输出JSON到API端点?

问题诊断与修复方案

核心问题分析

你的代码存在三个关键问题导致无输出:

  1. 流结构错误:用ZStream.fromZIO将定时任务的整体结果包装为单个流元素,而非将每个SomePayload拆分为独立输出单元,客户端无法接收分块数据。
  2. ContentLength配置错误:固定设置为100L与实际数据长度不匹配,服务器或客户端会因长度校验失败卡住。
  3. 定时逻辑位置错误:在ZIO层面使用repeat,而非ZStream层面,导致流仅输出一次结果后停止,无法定时重复触发任务。

修复后的完整代码

1. 先添加JSON序列化逻辑(以Circe为例)

import io.circe.generic.auto._
import io.circe.syntax._
import java.nio.charset.StandardCharsets

// 将Payload序列化为UTF-8编码的JSON字节
def payloadToBytes(payload: SomePayload): Task[Array[Byte]] =
  ZIO.succeed(payload.asJson.noSpaces.getBytes(StandardCharsets.UTF_8))

2. 修改端点定义(移除固定ContentLength)

val streamingEndpoint: PublicEndpoint[Unit, Unit, Stream[Throwable, Byte], ZioStreams] =
  endpoint.get
    .in("receive")
    .out(streamTextBody(ZioStreams)(CodecFormat.TextPlain(), Some(StandardCharsets.UTF_8)))

3. 修正Server逻辑(调整流结构与定时逻辑)

val streamingServerEndpoint: ZServerEndpoint[Any, ZioStreams] = streamingEndpoint.zServerLogic { _ =>
  val stream = ZStream
    // 每小时30分触发一次全量任务执行,首次执行立即触发
    .repeatSchedule(
      ZIO.collectAll(requests.map(service.processData)),
      Schedule.once ++ Schedule.minuteOfHour(30)
    )
    // 将每个请求返回的Seq[SomePayload]展开为单个元素流
    .flatMap(ZStream.fromIterable)
    // 序列化每个Payload为JSON字节
    .mapZIO(payloadToBytes)
    // 每个JSON后追加换行符,方便客户端逐行解析
    .concat(ZStream.succeed("\n".getBytes(StandardCharsets.UTF_8)))

  ZIO.succeed(stream)
}

关键修改说明

  • 流结构调整:用flatMap(ZStream.fromIterable)将批量结果拆分为单个Payload元素,确保流式输出的颗粒度正确。
  • 定时逻辑迁移:将repeat从ZIO层面移到ZStream的repeatSchedule,实现流的定时重复触发。
  • 移除固定ContentLength:流式输出无法提前预知总长度,框架会自动使用Transfer-Encoding: chunked分块传输。
  • 添加换行符:便于客户端逐行解析JSON数据,避免单一大块数据的解析压力。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 11:23:21