1. 为什么同步阻塞的 Skill 调用撑不住长任务
流式 Skill 调用、MCP、长时间运行任务、异步调用、Server-Sent Events 这几个词放在一起,本质上是在解决同一个尴尬:Agent 发起一个 Skill 调用后,只能干等一个完整结果。查询订单、算个价格这类毫秒级操作没问题,但视频生成、日志扫描、跨系统数据同步动辄几分钟到几小时,同步请求-响应模型立刻暴露三个硬伤。
第一个硬伤是超时。Agent 侧和网关侧通常都设了几十秒的请求超时,长任务几乎必然触发超时,调用直接失败,任务却可能还在后台跑,状态彻底失控。第二个硬伤是体验。用户盯着空白界面等几分钟,既看不到进度,也没法中途取消,只能反复刷新。第三个硬伤是资源效率。同步调用会长时间占住连接和线程,并发一上来,连接池就被拖垮,根本没法规模化。
我试过用最朴素的方式硬扛:把超时调到 10 分钟。结果是网关先断,客户端再断,任务状态两头对不上,排查起来非常痛苦。真正可行的方向是把「一次请求等一个结果」拆成「一次请求建立一条事件通道,结果分片推送」。这就是流式 Skill 调用的核心思路,而 MCP 作为工具调用协议,需要在传输层扩展出异步与流式两种能力。
异步调用解决的是「长时间等待」问题:Agent 发起 Action 后立刻拿到一个 task_id,稍后用这个标识符查询状态和结果,适合不需要实时交互的批处理。流式调用解决的是「中间结果可见」问题:Skill 通过一条持久连接持续推送 progress、log、partial_result 等事件,适合用户能从中途结果获益、或需要实时监控进度的任务。两者可以共存,一个长任务既能异步执行,也能流式推送进度。
本文聚焦流式回传的落地设计,以 Server-Sent Events 为通道,拆解从同步阻塞到异步分片推送的改造思路,并给出可复制的 MCP 服务端事件推送配置与客户端消费示例,最后附上超时、断线重连的验证动作。适合正在自建 AI 工具链、被长任务超时折磨过的开发者。
2. TaoToken 前置准备:把 MCP 服务端和模型调用接起来
在写流式推送代码之前,得先把模型调用这条链路打通。MCP 服务端负责执行 Skill 并推送事件,但 Skill 内部往往还要调用大模型做推理或总结,这部分请求需要走一个稳定的 API 入口。我用 TaoToken 作为模型调用的统一入口,它的 API 地址是 https://taotoken.net/api,兼容 OpenAI 风格的请求格式,接入成本很低。
先拿到 API Key。打开 https://taotoken.net/api-keys ,登录后创建一个新的 Key,复制保存。注意 Key 只在创建时完整显示一次,丢了就得重建。拿到 Key 后,把它写进环境变量,避免硬编码进代码仓库:
export TAOTOKEN_API_KEY="sk-你的key" export TAOTOKEN_BASE_URL="https://taotoken.net/api"如果你用的是 Claude Code 这类编码工具,或者想先验证模型是否可用,可以直接在模型对话页面测试:https://taotoken.net/model-chat 。输入一句简单的话,确认返回正常,说明 Key 和网络链路都没问题。这一步别跳过,很多后续报错其实是 Key 或 Base URL 写错导致的。
对于需要长期跑编码任务或 Agent 的场景,可以考虑 Coding Plan,它更适合高频、长时间的调用:https://taotoken.net/coding-plan 。控制台入口在 https://taotoken.net/console ,可以在这里查看用量、管理 Key、排查调用记录。接入文档在 https://taotoken.net/doc ,里面有各语言的完整示例。
这里要强调一个原则:MCP 服务端只负责 Skill 的执行与事件推送,模型调用通过 TaoToken 的 API 完成,两者职责分离。这样流式通道的稳定性不受模型调用耗时影响,即使某次模型请求慢,也不会阻塞事件推送。把 Key 和 Base URL 准备好,接下来进入服务端事件推送的配置。
3. 可复制的 MCP 服务端事件推送配置
流式回传的传输层选型有两个方向:Server-Sent Events 和 WebSocket。SSE 是单向的,服务器到客户端,基于标准 HTTP,实现简单、被广泛支持,适合只读的进度推送;缺点是客户端没法通过同一条通道发消息,提前终止需要另开一个接口。WebSocket 双向,功能更强但实现更复杂。对于大多数长任务的进度推送,SSE 足够用,本文以 SSE 为主。
先看服务端的事件推送配置。下面是一个基于 Node.js 的 MCP 服务端片段,用 SSE 推送 progress、partial_result、heartbeat、complete 四类事件。把它保存为mcp-sse-server.js:
// mcp-sse-server.js import express from "express"; const app = express(); app.use(express.json()); // 任务状态存储,生产环境应换成 Redis 或数据库 const tasks = new Map(); // SSE 流式端点:客户端通过 GET 建立事件通道 app.get("/mcp/stream/:taskId", (req, res) => { const { taskId } = req.params; res.setHeader("Content-Type", "text/event-stream"); res.setHeader("Cache-Control", "no-cache"); res.setHeader("Connection", "keep-alive"); res.setHeader("X-Accel-Buffering", "no"); // 关闭 Nginx 缓冲 res.flushHeaders(); const send = (event, data) => { res.write(`event: ${event}\n`); res.write(`data: ${JSON.stringify(data)}\n\n`); }; // 心跳保活,防止中间设备因空闲关闭连接 const heartbeat = setInterval(() => { send("heartbeat", { ts: Date.now() }); }, 15000); // 模拟长任务:分片推送进度 let percent = 0; const timer = setInterval(() => { percent += 10; send("progress", { taskId, percent, desc: `已处理 ${percent}%` }); if (percent === 50) { send("partial_result", { taskId, chunk: "前一半中间结果" }); } if (percent >= 100) { clearInterval(timer); clearInterval(heartbeat); send("complete", { taskId, result: "最终结果", url: "https://example.com/out" }); res.end(); } }, 1000); // 客户端断开时清理资源 req.on("close", () => { clearInterval(timer); clearInterval(heartbeat); res.end(); }); }); // 异步调用入口:立即返回 task_id app.post("/mcp/invoke", (req, res) => { const taskId = `task_${Date.now()}`; tasks.set(taskId, { status: "running", createdAt: Date.now() }); res.status(202).json({ accepted: true, taskId }); }); app.listen(3000, () => console.log("MCP SSE server on :3000"));这段代码有几个关键点。第一,响应头必须设置Content-Type: text/event-stream,并且调用res.flushHeaders()立即把头发出去,否则客户端会一直等。第二,X-Accel-Buffering: no是给 Nginx 看的,不关掉缓冲,事件会被攒着一起发,流式就变成了批量。第三,心跳事件每 15 秒发一次,防止连接被中间设备判定为空闲而关闭。第四,req.on("close")里清理定时器,避免客户端断开后服务端还在空转。
如果你用 Python,等价的 FastAPI 配置如下,保存为mcp_sse_server.py:
# mcp_sse_server.py import asyncio import json from fastapi import FastAPI, Request from fastapi.responses import StreamingResponse app = FastAPI() async def event_generator(task_id: str, request: Request): percent = 0 while percent < 100: if await request.is_disconnected(): break percent += 10 yield f"event: progress\ndata: {json.dumps({'taskId': task_id, 'percent': percent})}\n\n" if percent == 50: yield f"event: partial_result\ndata: {json.dumps({'taskId': task_id, 'chunk': '中间结果'})}\n\n" await asyncio.sleep(1) yield f"event: complete\ndata: {json.dumps({'taskId': task_id, 'result': 'done'})}\n\n" @app.get("/mcp/stream/{task_id}") async def stream(task_id: str, request: Request): return StreamingResponse( event_generator(task_id, request), media_type="text/event-stream", headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"}, )服务端配置里还有一个容易忽略的点:事件序列号。为了实现断线重连后的断点续传,每个事件应该带一个递增的id字段。SSE 协议原生支持id:行,客户端重连时会通过Last-Event-ID请求头把最后收到的事件 ID 带回来,服务端据此从断点继续推送。改造方式是在send函数里加一行:
let seq = 0; const send = (event, data) => { seq += 1; res.write(`id: ${seq}\n`); res.write(`event: ${event}\n`); res.write(`data: ${JSON.stringify(data)}\n\n`); };服务端收到重连请求时,读取req.headers["last-event-id"],从该 ID 之后继续推送。这一步是长任务可靠性的关键,网络抖动不可避免,没有断点续传,用户就得从头再来。
4. 客户端消费示例与成功结果验证
服务端把事件推出来了,客户端得能正确消费。下面是一个 Node.js 客户端示例,用原生fetch消费 SSE 流,处理 progress、partial_result、complete 三类事件,并实现断线重连。保存为mcp-client.js:
// mcp-client.js const BASE = "http://localhost:3000"; let lastEventId = 0; async function consumeStream(taskId) { const headers = { Accept: "text/event-stream" }; if (lastEventId > 0) { headers["Last-Event-ID"] = String(lastEventId); } const res = await fetch(`${BASE}/mcp/stream/${taskId}`, { headers }); if (!res.ok) throw new Error(`HTTP ${res.status}`); const reader = res.body.getReader(); const decoder = new TextDecoder(); let buffer = ""; while (true) { const { value, done } = await reader.read(); if (done) break; buffer += decoder.decode(value, { stream: true }); // SSE 以空行分隔事件块 const blocks = buffer.split("\n\n"); buffer = blocks.pop(); for (const block of blocks) { let event = "message"; let data = ""; for (const line of block.split("\n")) { if (line.startsWith("id: ")) lastEventId = Number(line.slice(4)); else if (line.startsWith("event: ")) event = line.slice(7); else if (line.startsWith("data: ")) data = line.slice(6); } handleEvent(event, JSON.parse(data || "{}")); } } } function handleEvent(event, data) { switch (event) { case "progress": console.log(`进度 ${data.percent}% - ${data.desc || ""}`); break; case "partial_result": console.log("收到中间结果:", data.chunk); break; case "heartbeat": // 心跳只用于保活,不处理业务 break; case "complete": console.log("任务完成:", data.result); break; case "error": console.error("任务失败:", data.message); break; } } // 带重连的启动逻辑 async function run(taskId) { let retries = 0; while (retries < 5) { try { await consumeStream(taskId); break; // 正常结束 } catch (err) { retries += 1; console.warn(`连接断开,第 ${retries} 次重连,从事件 ${lastEventId} 继续`); await new Promise((r) => setTimeout(r, 1000 * retries)); } } } run("task_demo_001");客户端有几个细节值得说。第一,SSE 的事件块以空行\n\n分隔,必须按块解析,不能按行直接处理,否则会把一个事件的 id、event、data 拆散。第二,buffer = blocks.pop()保留最后一个不完整的块,等下次数据到达再拼,这是处理 TCP 分片的标准做法。第三,重连时把lastEventId通过Last-Event-ID头带回去,服务端就能从断点续传,不会重复推送已经处理过的事件。
验证成功结果的动作很直接。先启动服务端:
node mcp-sse-server.js再开一个终端跑客户端:
node mcp-client.js预期输出如下,进度从 10% 递增到 100%,中途出现一次中间结果,最后是完成事件:
进度 10% - 已处理 10% 进度 20% - 已处理 20% ... 进度 50% - 已处理 50% 收到中间结果: 前一半中间结果 ... 进度 100% - 已处理 100% 任务完成: 最终结果如果想用命令行快速验证 SSE 通道是否正常,可以用 curl:
curl -N http://localhost:3000/mcp/stream/task_demo_001-N参数关闭 curl 的缓冲,事件会实时打印出来。如果看到event: progress一行行滚动,说明服务端推送正常。如果 curl 卡住不动,多半是响应头没设对,或者中间有代理在缓冲。
5. 本篇常见错误排查:401、断流与解析异常
流式调用踩坑的地方比同步调用多,因为多了一条长连接和一套事件协议。下面按真实报错逐条排查。
401 Unauthorized。这个报错通常出现在 Skill 内部调用模型 API 时。检查TAOTOKEN_API_KEY是否设置正确,Base URL 是否为https://taotoken.net/api。常见错误是把 Key 写成了别的服务的,或者环境变量没导出到当前 shell。用echo $TAOTOKEN_API_KEY确认一下。如果 Key 正确仍然 401,去控制台 https://taotoken.net/console 看下 Key 是否被禁用或额度耗尽。
local proxy failed / connection refused。客户端连不上服务端,先确认服务端进程在跑、端口没被占用。lsof -i :3000看下端口状态。如果服务端在容器里,确认端口映射正确。这个报错和模型 API 无关,纯粹是本地网络问题,别往 Key 上找原因。
reading choices 报错。这个通常出现在 Skill 内部解析模型返回时。流式场景下,模型返回可能是分片的,如果代码按完整 JSON 解析,就会在分片边界处报reading 'choices'之类的错。解决办法是先把流式响应拼完整再解析,或者用支持流式增量的解析器。检查你的模型调用是否误用了流式模式却按非流式解析。
OAuth 相关报错。如果你用的是 Claude Code 或类似工具接入,可能遇到 OAuth 认证失败。这类工具通常需要配置 Base URL、Key、Model ID 三件套。以 Claude Code 为例,检查配置文件里的ANTHROPIC_BASE_URL是否指向https://taotoken.net/api,ANTHROPIC_API_KEY是否为你的 Key,Model ID 是否填写正确。三件套缺一不可,任何一个写错都会报认证或模型不存在。
事件收不到 / 流式变成批量。这是最隐蔽的问题。表现是客户端等了很久,然后一次性收到所有事件。原因通常是中间有反向代理在缓冲。检查 Nginx 配置,确保proxy_buffering off;,并且服务端响应头带了X-Accel-Buffering: no。另外确认客户端没有对响应做整体缓冲,比如某些 HTTP 库默认会读完整个 body 才返回。
连接频繁断开。长连接被中间设备判定为空闲而关闭。解决办法是服务端定期发 heartbeat 事件,间隔建议 15 到 30 秒。客户端收到 heartbeat 不处理业务,只用于确认连接还活着。如果断开后没有重连逻辑,任务就丢了,所以客户端必须实现带Last-Event-ID的重连。
背压导致内存暴涨。如果 Skill 产生事件的速度远快于客户端消费速度,事件会在服务端或网关积压,内存持续上涨。解决办法是实现流控:客户端处理慢时,服务端降低发送频率,或者丢弃非关键的 log 事件,只保留 progress 和 complete。生产环境建议给事件缓冲区设上限,超限时主动断开并让客户端重连。
排查顺序建议从外到内:先确认网络连通(curl 能否收到事件),再确认认证(Key 和 Base URL),最后确认解析逻辑(事件块分隔、JSON 解析时机)。大部分问题出在响应头配置和事件解析这两步。
6. 把流式 Skill 调用接进你的工具链
流式回传落地后,长任务的体验会有质的变化。用户能看到进度、拿到中间结果、中途取消,Agent 也不再被超时卡死。把这条链路接进自有工具链时,有几个实践建议。
第一,事件类型保持精简。progress、partial_result、heartbeat、complete、error 这五类基本够用,log 事件按需开启,避免审计日志膨胀。第二,所有事件带递增 ID,客户端重连时用Last-Event-ID续传,这是可靠性的底线。第三,心跳间隔别超过 30 秒,中间设备的空闲超时通常在这个量级。第四,服务端在客户端断开时务必清理定时器和后台任务,否则会积累僵尸任务。
模型调用这条链路,统一走 TaoToken 的 API 入口,Key 在 https://taotoken.net/api-keys 管理,接入细节看 https://taotoken.net/doc 。需要长期跑 Agent 或编码任务的,Coding Plan 更合适:https://taotoken.net/coding-plan 。想先验证模型可用性,直接在 https://taotoken.net/model-chat 试一句就行。
最后留一个实用技巧:在服务端加一个任务状态查询接口,即使流式连接断了,客户端也能通过 task_id 查当前状态和已产出的部分结果。流式和异步不是二选一,组合起来才最稳。