如何用Scala的ZIO、ZStream定时流式输出JSON到API端点?
问题诊断与修复方案
核心问题分析
你的代码存在三个关键问题导致无输出:
- 流结构错误:用
ZStream.fromZIO将定时任务的整体结果包装为单个流元素,而非将每个SomePayload拆分为独立输出单元,客户端无法接收分块数据。 - ContentLength配置错误:固定设置为
100L与实际数据长度不匹配,服务器或客户端会因长度校验失败卡住。 - 定时逻辑位置错误:在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
相关产品推荐
相关产品推荐

