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

如何在Ktor中将ByteReadChannel转为Flow并过滤JSON数组?

在Ktor中流式处理JSON数组并过滤对象

前置准备

首先确保你的Ktor项目已经配置了ContentNegotiation和Kotlinx序列化插件,用于JSON处理:

install(ContentNegotiation) {
    Json(Json {
        ignoreUnknownKeys = true
        prettyPrint = false
    })
}

同时,定义目标数据类并添加序列化注解:

@Serializable
data class MyJsonObject(
    val id: Int,
    val name: String,
    // 其他业务属性
)

1. 将ByteReadChannel转为流式Flow

通过Kotlinx序列化的JsonReader实现流式解析,避免一次性加载整个JSON数组到内存:

import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.flow.flowOn
import kotlinx.serialization.json.Json
import kotlinx.serialization.json.JsonReader
import kotlinx.serialization.decodeFromJsonReader
import io.ktor.utils.io.ByteReadChannel
import kotlinx.coroutines.Dispatchers

fun ByteReadChannel.toJsonObjectFlow(): Flow<MyJsonObject> = flow {
    val json = Json { ignoreUnknownKeys = true }
    // 将ByteReadChannel转为输入流,流式读取数据
    this@toJsonObjectFlow.inputStream().use { inputStream ->
        JsonReader(inputStream).use { reader ->
            // 开始遍历JSON数组结构
            reader.beginArray()
            while (reader.hasNext()) {
                // 逐个解析数组中的JSON对象
                val obj = json.decodeFromJsonReader(MyJsonObject.serializer(), reader)
                emit(obj)
            }
            reader.endArray()
        }
    }
}.flowOn(Dispatchers.IO) // IO操作切换到IO调度器执行

2. 过滤Flow中的对象

直接使用Flow的filter操作符,根据业务属性筛选符合条件的对象:

val filteredFlow = byteReadChannel.toJsonObjectFlow()
    .filter { obj ->
        // 自定义过滤逻辑,例如筛选name字段为"target"的对象
        obj.name == "target"
    }

3. 将过滤结果流式响应给客户端

为了避免再次加载全部结果到内存,使用respondTextWriter逐个写入过滤后的JSON对象,组成合法的JSON数组响应:

import io.ktor.server.application.call
import io.ktor.server.response.respondTextWriter
import io.ktor.http.ContentType

call.respondTextWriter(contentType = ContentType.Application.Json) {
    write("[")
    var isFirstElement = true
    filteredFlow.collect { obj ->
        if (!isFirstElement) write(",")
        isFirstElement = false
        // 逐个序列化并写入对象
        Json.encodeToWriter(this, MyJsonObject.serializer(), obj)
    }
    write("]")
}

关键说明

  • 全程流式处理:从读取响应体、解析JSON对象、过滤到响应客户端,所有环节都不会一次性加载全部数据到内存,适合处理包含数千个对象的超大JSON数组。
  • 异常处理:可以在flow块中添加try-catch捕获解析或IO异常(如JSON格式错误、通道中断等),根据业务需求返回对应的错误响应。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 18:44:52