如何在Ktor中编写SSE事件?下载进度推送失败求助
SSE 推送下载进度失效问题排查与修复
我需要在收到触发下载的POST请求后,通过Server Sent Events(SSE)将下载进度推送给客户端,但当前代码无法正常运行。以下是我的客户端和服务端代码:
客户端JS代码
const result = document.getElementById("result"); const es = new EventSource("/download/process"); es.onmessage = function (event) { const data = event.data; console.log("Displaying:" + data); if (result.innerText.length === 0) { result.innerText = data; } else { result.innerText = result.innerText + "\n" + data; } result.scrollTop = result.scrollHeight; }
服务端Ktor代码
val flo = MutableSharedFlow<String>() routing { route("download") { get("process") { call.response.cacheControl(CacheControl.NoCache(null)) call.respondBytesWriter(ContentType.Text.EventStream) { flo.collect { writeStringUtf8(it) flush() } } } post("start") { // Get parameters here runBlocking { // Perform download and send to client event throughout flo.tryEmit(it) } call.respondText("Download finished") } } }
问题分析与修复方案
- SSE消息格式不规范:EventSource仅能识别以
data:开头、\n\n结尾的消息。当前直接写入字符串,客户端无法解析为有效事件,需按照规范格式化消息内容。 - SharedFlow配置不合理:默认
MutableSharedFlow的replay=0,无法保留历史事件,导致晚连接的客户端或下载过程中连接的客户端错过进度。需设置replay=1来保留最新进度。 - emit参数错误:POST请求中
flo.tryEmit(it)的it是ApplicationCall对象,并非实际下载进度内容,需替换为真实的进度信息(如百分比、已下载大小)。 - 协程使用不当:用
runBlocking会阻塞Ktor线程池,应使用Ktor自带的协程上下文启动下载任务,避免阻塞。
修正后的代码
客户端代码(增加错误处理)
const result = document.getElementById("result"); const es = new EventSource("/download/process"); es.onmessage = function (event) { const data = event.data; console.log("Displaying:" + data); // 简化文本拼接逻辑 result.innerText = result.innerText ? `${result.innerText}\n${data}` : data; result.scrollTop = result.scrollHeight; }; // 新增错误处理,便于排查连接问题 es.onerror = function(error) { console.error("EventSource 连接错误:", error); es.close(); };
服务端代码
// 配置SharedFlow,保留最新1条进度,确保新连接能获取当前状态 val progressFlow = MutableSharedFlow<String>(replay = 1) routing { route("download") { get("process") { call.response.cacheControl(CacheControl.NoCache(null)) // 设置SSE必需的响应头 call.response.headers.append("Connection", "keep-alive") call.respondBytesWriter(ContentType.Text.EventStream) { progressFlow.collect { progress -> // 按照SSE规范写入消息 writeStringUtf8("data: $progress\n\n") flush() } } } post("start") { // 获取下载参数(示例:从查询参数取下载链接) val downloadUrl = call.request.queryParameters["url"] ?: run { call.respondBadRequest("缺少下载链接参数") return@post } // 使用Ktor协程上下文启动下载任务,避免阻塞线程 launch(call.application.coroutineContext) { try { // 模拟下载过程,实际替换为真实下载逻辑 for (progress in 0..100 step 10) { delay(500) // 模拟下载耗时 progressFlow.emit("当前下载进度:$progress%") } progressFlow.emit("下载完成") } catch (e: Exception) { progressFlow.emit("下载失败:${e.message ?: "未知错误"}") } } call.respondText("下载已启动") } } }
额外注意事项
- 需确保客户端先建立SSE连接(页面加载时初始化EventSource),再发送POST请求触发下载,避免错过初始进度事件。
- 如果需支持多用户独立进度,不能使用全局SharedFlow,需结合会话(如
call.sessions)为每个用户创建独立的进度流。
内容的提问来源于stack exchange,提问作者DVT
相关产品推荐
相关产品推荐

