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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 21:55:34