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

如何让Ktor调用Groq API时实现SSE实时流式数据处理?

Ktor调用Groq API实现SSE实时流式传输问题

问题描述

在Ktor中通过Server-Sent Events(SSE)调用Groq API时,代码会等待完整响应返回后才处理并输出数据块,无法做到数据一到达就立即处理,不符合实时流式传输的预期。

现有实现代码

class ApiClient(
    private val apiKey: String,
    private val baseUrl: String = "https://api.groq.com/openai"
) {
    private val client = HttpClient(CIO) {
        install(ContentNegotiation) {
            json(Json {
                encodeDefaults = true
                ignoreUnknownKeys = true
                explicitNulls = false
            })
        }
        install(SSE)
    }

    suspend fun chatCompletionCreateWithStream(
        model: String,
        messages: List<Message>,
        stream: Boolean = true,
        tool_choice: String = "auto",
        tools: List<Tool>? = null,
        coroutineContext: CoroutineContext = Dispatchers.Default
    ): Flow<Result<ChatCompletion>> {
        val requestPayload = ChatRequest(
            model = model,
            messages = messages,
            stream = stream,
            tool_choice = tool_choice,
            tools = tools
        )

        return flow {
            val response = client.post("$baseUrl/v1/chat/completions") {
                contentType(ContentType.Application.Json)
                headers {
                    append("Accept", "text/event-stream")
                    append("Authorization", "Bearer $apiKey")
                }
                setBody(requestPayload)
            }

            response.bodyAsChannel().apply {
                while (!isClosedForRead) {
                    val line = readUTF8Line(Int.MAX_VALUE)
                        ?.removePrefix("data: ")
                        ?.takeIf { it.startsWith("{") }
                        ?: continue

                    try {
                        val chunk = Json.decodeFromString<ChatCompletion>(line)
                        emit(Result.success(chunk))  // Emitting as the data comes in
                    } catch (e: Exception) {
                        emit(Result.failure<ChatCompletion>(e))
                    }
                }
            }
        }.flowOn(coroutineContext)
         .catch { e ->
             emit(Result.failure<ChatCompletion>(e))
         }
    }
}

问题详情

期望数据块一到达就被实时处理并输出,无需等待全部响应接收完成,但当前代码仅在获取完整响应后才开始处理和输出数据块。已设置Accept: text/event-stream头启用流式传输,但未实现实时流式处理。

当前已尝试方案

  • 使用bodyAsChannel()以流的方式读取响应
  • 使用readUTF8Line()处理每行数据并通过Kotlin Flow输出数据块
  • 设置Accept: text/event-stream请求头用于流式传输

解决方案

1. 配置CIO客户端禁用响应缓冲

CIO客户端默认会缓冲响应数据,需要显式关闭缓冲并适配长连接参数,确保数据一到达就可读:

private val client = HttpClient(CIO) {
    install(ContentNegotiation) {
        json(Json {
            encodeDefaults = true
            ignoreUnknownKeys = true
            explicitNulls = false
        })
    }
    install(SSE)
    // 配置引擎参数,禁用缓冲与超时,适配流式长连接
    engine {
        requestTimeout = 0
        pipelining = false
        endpoint {
            connectTimeout = 10_000
            sslConnectTimeout = 10_000
            socketTimeout = 0
            keepAliveTime = 30_000
            keepAliveTimeout = 5_000
        }
    }
}

2. 使用Ktor SSE专用API处理响应

利用已安装的SSE插件提供的receiveSSE()方法,自动处理SSE格式(包括data:前缀、结束标记等),替代手动读取Channel:

return flow {
    client.post("$baseUrl/v1/chat/completions") {
        contentType(ContentType.Application.Json)
        headers {
            append("Accept", "text/event-stream")
            append("Authorization", "Bearer $apiKey")
            append("Cache-Control", "no-cache")
            append("Connection", "keep-alive")
        }
        setBody(requestPayload)
    }.receiveSSE {
        for (event in this) {
            // 过滤Groq API的流式结束标记
            if (event.data == "[DONE]") continue
            
            try {
                val chunk = Json.decodeFromString<ChatCompletion>(event.data)
                emit(Result.success(chunk))
            } catch (e: Exception) {
                emit(Result.failure<ChatCompletion>(e))
            }
        }
    }
}.flowOn(coroutineContext)
 .catch { e ->
     emit(Result.failure<ChatCompletion>(e))
 }

3. 确认请求参数有效性

确保请求Payload中的stream参数严格设置为true,Groq API仅在该参数为真时返回流式响应。

4. 避免Flow下游不必要缓冲

收集Flow时避免使用buffer()等缓冲操作,确保数据实时处理:

apiClient.chatCompletionCreateWithStream(...)
    .collect { result ->
        result.onSuccess { chunk ->
            // 实时处理每个数据块
            println(chunk.choices.firstOrNull()?.delta?.content)
        }.onFailure { error ->
            println("Error: ${error.message}")
        }
    }

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 10:43:12