news 2026/10/1 14:31:40

流式 Skill 调用实战:用 MCP 支撑长时间运行任务的异步回传

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
流式 Skill 调用实战:用 MCP 支撑长时间运行任务的异步回传

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 查当前状态和已产出的部分结果。流式和异步不是二选一,组合起来才最稳。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/1 14:31:20

Windows SSH安装全攻略:从客户端到服务端,避开所有坑

装SSH这话题看着基础&#xff0c;实际翻车点全藏在细节里。Windows自带OpenSSH、Git自带SSH、第三方工具Bitvise&#xff0c;再加上VSCode远程插件一搅和&#xff0c;新手很容易装完连不上、连上传不了文件、传了又权限报错。这篇不讲虚的&#xff0c;直接按我自己的实操顺序来…

作者头像 李华
网站建设 2026/10/1 14:30:52

Cursor小团队产品开发实践记录:从模块化设计到TaoToken统一Key接入

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/10/1 14:29:01

在线文本字数统计工具,文案笔记技术文档统计小助手

一、前言 日常工作学习当中&#xff0c;字数统计是十分常见的需求。写 CSDN 博客、撰写技术文档、整理需求说明书、编写投稿内容、整理会议纪要的时候&#xff0c;经常需要了解整篇文档的总字符、汉字数量、英文单词、行数等信息。 如果使用 Word、WPS 可以完成统计&#xff…

作者头像 李华
网站建设 2026/10/1 14:28:07

基于Seed-2.1-pro-0915的电商参考图批量生成工作台实践与验证

做电商视觉的朋友&#xff0c;大概率都被“批量出图”这件事折磨过&#xff1a;产品明明很好看&#xff0c;一进AI生成就变形&#xff1b;张张都要人工盯&#xff0c;出图速度还不如外包。我最近把Seed-2.1-pro-0915接进了一套完整的参考图批量生成工作台&#xff0c;从产品参考…

作者头像 李华
网站建设 2026/10/1 14:28:05

Snappy与Zstandard大对比:大数据压缩格式选型实战指南

这题我太熟了。不管是搞数仓、做实时计算还是维护Hadoop集群&#xff0c;压缩格式选型几乎是每天都要碰的事。早些年大家无脑选Snappy&#xff0c;因为这玩意儿到处都支持&#xff0c;性能也稳&#xff1b;但这两年Zstandard&#xff08;简称zstd&#xff09;势头很猛&#xff…

作者头像 李华