如何让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
相关产品推荐
相关产品推荐

