1. 为什么 AI 调用订单接口会被限流:消息队列 MCP Server 要解决的真实问题
消息队列 MCP Server 是一类把 RabbitMQ、Kafka 这类中间件能力封装成 MCP 工具的服务,让 AI 客户端能通过标准 tool 调用去提交异步任务、发布事件、查询任务状态。它适合谁?适合那些 AI 需要触发耗时操作、又不能让模型调用干等的场景,比如批量数据处理、跨系统同步、订单状态轮询、运维批量操作。
我试过最典型的一个坑:AI 客服每次用户问“我的订单发货了吗”,模型就同步去调订单系统 API。有个用户一天问了 40 多次,订单系统直接把 AI 客服的出口 IP 限流了,后面所有用户都查不了。根因不是订单系统小气,而是架构错了——AI 不应该同步查,应该把查询请求丢进消息队列异步处理,先返回“正在查询中”,结果出来再通知。
另一个场景是 AI 运维助手。运维同学让 AI 批量重启几十台机器,同步等的话 MCP 的 tool 调用直接超时。MCP 协议里 tool 调用一般要求 30 秒内返回,超时客户端就掐断。把耗时操作改成消息队列异步执行后,AI 提交任务拿一个 taskId,然后轮询状态,体验完全不一样。
所以这篇的核心链路是:AI 客户端 → MCP Server(暴露 submit_task / get_task_status 等工具)→ RabbitMQ(任务分发)→ 消费者(真正干活)→ Redis(任务状态存储)。下面从环境准备到一条消息触发任务并验证结果,全部给可复制的配置和代码。
2. TaoToken 统一 Key 前置准备:一个 Key 打通模型调用与 MCP 工具链
在写 MCP Server 之前,先把模型侧的调用凭证统一掉。很多同学做 MCP 项目时,模型调用和工具调用分散在好几套 Key 里,调试时来回切换很痛苦。TaoToken 的做法是给一个统一 Key,模型对话、Coding Plan、API 调用都走同一个入口,MCP Server 里如果需要调用模型做任务编排,直接复用这个 Key 就行。
具体操作路径:
第一步,打开官网 https://taotoken.net/?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= 注册并登录。
第二步,进控制台创建 API Key,地址是 https://taotoken.net/console?utm_source=taotoken_aicg_blog_end&utm_campaign=rewrite&utm_content= 。创建后复制出来,形如sk-xxxxxxxx,这个 Key 后面既用于模型调用,也用于 MCP Server 里需要模型能力的场景。
第三步,如果你要做长期编码或 Agent 类任务,可以看下 Coding Plan:https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_campaign=rewrite&utm_content= 。它适合需要持续调用模型、跑长任务的场景,比按次调用更划算。
第四步,接入文档在 https://taotoken.net/doc?utm_source=taotoken_aicg_blog_end&utm_campaign=rewrite&utm_content= ,里面有 Base URL、Model ID 的完整说明。API 入口是 https://taotoken.net/api ,注意这个地址不带 UTM 参数,配置时直接用。
这里要强调一个三件套概念:不管你用的是 Claude Code、Cline、还是自己写的 MCP Server,接入任何模型服务都需要三个东西——Base URL、API Key、Model ID。TaoToken 的 Base URL 是https://taotoken.net/api,Key 是你在控制台创建的那串,Model ID 在文档里能查到。这三个配齐,模型调用就通了。
如果你用的是 Claude Code 这类工具,配置方式是在 settings 里指定 Base URL 和 Key。如果是自己写的 MCP Server,就在环境变量里注入。下面给一个通用的环境变量配置:
# .env 文件,放在项目根目录 TAOTOKEN_API_KEY=sk-你的Key TAOTOKEN_BASE_URL=https://taotoken.net/api TAOTOKEN_MODEL=claude-sonnet-4-20250514 # RabbitMQ 和 Redis 配置 RABBITMQ_URL=amqp://guest:guest@localhost:5672 REDIS_URL=redis://localhost:6379为什么要先做这一步?因为消息队列 MCP Server 本身不依赖模型,但你的 AI 客户端要能连上模型才能发起 tool 调用。把 Key 统一到 TaoToken 后,模型侧和工具侧的环境变量集中管理,调试时不用改来改去。而且后面如果任务处理器里需要调用模型做内容生成(比如生成报告摘要),直接复用同一个 Key 就行。
3. 可复制的 MCP Server 配置:RabbitMQ + Redis + 工具定义
这一节给完整的可复制配置。项目结构如下:
mq-mcp-server/ ├── package.json ├── tsconfig.json ├── .env └── src/ ├── index.ts # MCP Server 入口,工具定义 ├── rabbitmq.ts # RabbitMQ 连接与交换机/队列声明 ├── task_store.ts # Redis 任务状态存储 ├── publisher.ts # 消息发布器 └── consumer.ts # 消息消费者与任务处理器先看package.json,依赖很干净:
{ "name": "mq-mcp-server", "version": "1.0.0", "type": "module", "scripts": { "build": "tsc", "start": "node dist/index.js" }, "dependencies": { "@modelcontextprotocol/sdk": "^1.0.0", "amqplib": "^0.10.0", "redis": "^4.6.0" }, "devDependencies": { "typescript": "^5.3.0", "@types/amqplib": "^0.10.0", "@types/node": "^20.10.0" } }tsconfig.json关键配置:
{ "compilerOptions": { "target": "ES2022", "module": "ES2022", "moduleResolution": "node", "outDir": "./dist", "rootDir": "./src", "strict": true, "esModuleInterop": true, "skipLibCheck": true }, "include": ["src/**/*"] }RabbitMQ 连接管理模块src/rabbitmq.ts,负责声明交换机、队列、死信队列和绑定关系:
import amqp, { Connection, Channel, ConsumeMessage } from "amqplib"; const RABBITMQ_URL = process.env.RABBITMQ_URL || "amqp://localhost:5672"; let connection: Connection | null = null; let channel: Channel | null = null; const EXCHANGES = { TASKS: "mcp.tasks", EVENTS: "mcp.events", BROADCAST: "mcp.broadcast", DEAD_LETTER: "mcp.dlx", }; const QUEUES = { TASK_DEFAULT: "mcp.task.default", TASK_HIGH: "mcp.task.high", EVENT: "mcp.event", DEAD_LETTER: "mcp.dlq", }; export async function initRabbitMQ(): Promise<void> { connection = await amqp.connect(RABBITMQ_URL); connection.on("close", () => { console.error("RabbitMQ 连接关闭,5 秒后重连..."); setTimeout(() => initRabbitMQ().catch(console.error), 5000); }); connection.on("error", (err) => { console.error("RabbitMQ 连接错误:", err.message); }); channel = await connection.createChannel(); // 死信交换机与队列 await channel.assertExchange(EXCHANGES.DEAD_LETTER, "direct", { durable: true }); await channel.assertQueue(QUEUES.DEAD_LETTER, { durable: true }); await channel.bindQueue(QUEUES.DEAD_LETTER, EXCHANGES.DEAD_LETTER, "task.failed"); // 任务交换机(direct) await channel.assertExchange(EXCHANGES.TASKS, "direct", { durable: true }); await channel.assertQueue(QUEUES.TASK_DEFAULT, { durable: true, deadLetterExchange: EXCHANGES.DEAD_LETTER, deadLetterRoutingKey: "task.failed", messageTtl: 30 * 60 * 1000, }); await channel.bindQueue(QUEUES.TASK_DEFAULT, EXCHANGES.TASKS, "task.default"); await channel.assertQueue(QUEUES.TASK_HIGH, { durable: true, deadLetterExchange: EXCHANGES.DEAD_LETTER, deadLetterRoutingKey: "task.failed", messageTtl: 10 * 60 * 1000, maxPriority: 10, }); await channel.bindQueue(QUEUES.TASK_HIGH, EXCHANGES.TASKS, "task.high"); // 事件交换机(topic) await channel.assertExchange(EXCHANGES.EVENTS, "topic", { durable: true }); await channel.assertQueue(QUEUES.EVENT, { durable: true }); await channel.bindQueue(QUEUES.EVENT, EXCHANGES.EVENTS, "event.#"); // 广播交换机(fanout) await channel.assertExchange(EXCHANGES.BROADCAST, "fanout", { durable: true }); console.error("RabbitMQ 初始化完成"); } export function getChannel(): Channel { if (!channel) throw new Error("RabbitMQ 通道未初始化"); return channel; } export { EXCHANGES, QUEUES }; export type { ConsumeMessage };任务状态存储src/task_store.ts,用 Redis 存任务生命周期:
import { createClient, RedisClientType } from "redis"; const REDIS_URL = process.env.REDIS_URL || "redis://localhost:6379"; const redis: RedisClientType = createClient({ url: REDIS_URL }); export enum TaskStatus { PENDING = "pending", PROCESSING = "processing", COMPLETED = "completed", FAILED = "failed", RETRYING = "retrying", } export interface Task { id: string; type: string; payload: any; status: TaskStatus; result?: any; error?: string; createdAt: number; updatedAt: number; retryCount: number; maxRetries: number; priority: number; } const TASK_KEY_PREFIX = "mcp:task:"; const TASK_STATUS_INDEX_PREFIX = "mcp:task:status:"; const TASK_TTL = 7 * 24 * 60 * 60; export async function initRedis(): Promise<void> { await redis.connect(); console.error("Redis 连接已建立"); } export async function createTask( type: string, payload: any, priority: number = 0, maxRetries: number = 3 ): Promise<Task> { const now = Date.now(); const task: Task = { id: `task_${now}_${Math.random().toString(36).substring(2, 10)}`, type, payload, status: TaskStatus.PENDING, createdAt: now, updatedAt: now, retryCount: 0, maxRetries, priority, }; await redis.set(TASK_KEY_PREFIX + task.id, JSON.stringify(task), { EX: TASK_TTL }); const statusKey = TASK_STATUS_INDEX_PREFIX + TaskStatus.PENDING; await redis.sAdd(statusKey, task.id); await redis.expire(statusKey, TASK_TTL); return task; } export async function updateTaskStatus( taskId: string, status: TaskStatus, extra?: { result?: any; error?: string } ): Promise<Task | null> { const key = TASK_KEY_PREFIX + taskId; const raw = await redis.get(key); if (!raw) return null; const task: Task = JSON.parse(raw); const oldStatus = task.status; task.status = status; task.updatedAt = Date.now(); if (extra?.result !== undefined) task.result = extra.result; if (extra?.error !== undefined) task.error = extra.error; if (status === TaskStatus.RETRYING) task.retryCount += 1; await redis.set(key, JSON.stringify(task), { EX: TASK_TTL }); if (oldStatus !== status) { await redis.sRem(TASK_STATUS_INDEX_PREFIX + oldStatus, taskId); const newStatusKey = TASK_STATUS_INDEX_PREFIX + status; await redis.sAdd(newStatusKey, taskId); await redis.expire(newStatusKey, TASK_TTL); } return task; } export async function getTask(taskId: string): Promise<Task | null> { const raw = await redis.get(TASK_KEY_PREFIX + taskId); return raw ? (JSON.parse(raw) as Task) : null; } export async function getTasksByStatus( status: TaskStatus, limit: number = 50 ): Promise<Task[]> { const taskIds = await redis.sMembers(TASK_STATUS_INDEX_PREFIX + status); const tasks: Task[] = []; for (const id of taskIds.slice(0, limit)) { const task = await getTask(id); if (task) tasks.push(task); } return tasks; } export async function getTaskStats(): Promise<Record<string, number>> { const stats: Record<string, number> = {}; for (const status of Object.values(TaskStatus)) { stats[status] = await redis.sCard(TASK_STATUS_INDEX_PREFIX + status); } return stats; } export { redis };消息发布器src/publisher.ts:
import { getChannel, EXCHANGES, QUEUES } from "./rabbitmq.js"; import { createTask, Task } from "./task_store.js"; export async function submitTask( type: string, payload: any, priority: number = 0 ): Promise<Task> { const task = await createTask(type, payload, priority); const channel = getChannel(); const routingKey = priority >= 5 ? "task.high" : "task.default"; const message = JSON.stringify({ taskId: task.id, type: task.type, priority: task.priority, }); channel.publish(EXCHANGES.TASKS, routingKey, Buffer.from(message), { persistent: true, priority: priority, headers: { "x-task-type": type, "x-created-at": task.createdAt, }, expiration: priority >= 5 ? "600000" : "1800000", }); return task; } export async function publishEvent(eventType: string, data: any): Promise<void> { const channel = getChannel(); const routingKey = `event.${eventType}`; const message = JSON.stringify({ eventType, data, timestamp: Date.now() }); channel.publish(EXCHANGES.EVENTS, routingKey, Buffer.from(message), { persistent: true, }); } export async function broadcastMessage(message: string): Promise<void> { const channel = getChannel(); channel.publish(EXCHANGES.BROADCAST, "", Buffer.from(message), { persistent: true, }); } export async function getDeadLetterCount(): Promise<number> { const channel = getChannel(); const info = await channel.checkQueue(QUEUES.DEAD_LETTER); return info.messageCount; }消费者src/consumer.ts,注意 prefetch 必须在 consume 之前设置:
import { getChannel, QUEUES, ConsumeMessage } from "./rabbitmq.js"; import { getTask, updateTaskStatus, TaskStatus, Task } from "./task_store.js"; type TaskHandler = (task: Task) => Promise<any>; const handlers = new Map<string, TaskHandler>(); export function registerHandler(type: string, handler: TaskHandler): void { handlers.set(type, handler); console.error(`已注册任务处理器: ${type}`); } async function processMessage(msg: ConsumeMessage | null): Promise<void> { if (!msg) return; const channel = getChannel(); const content = JSON.parse(msg.content.toString()); const taskId = content.taskId; try { const task = await getTask(taskId); if (!task) { channel.ack(msg); console.error(`任务 ${taskId} 不存在,消息已丢弃`); return; } // 幂等性检查:已完成的任务直接跳过 if (task.status === TaskStatus.COMPLETED) { channel.ack(msg); console.error(`任务 ${taskId} 已完成,跳过重复消费`); return; } await updateTaskStatus(taskId, TaskStatus.PROCESSING); const handler = handlers.get(task.type); if (!handler) throw new Error(`未找到任务类型 ${task.type} 的处理器`); const result = await handler(task); await updateTaskStatus(taskId, TaskStatus.COMPLETED, { result }); channel.ack(msg); console.error(`任务 ${taskId} 处理完成`); } catch (error) { const errorMsg = error instanceof Error ? error.message : String(error); const task = await getTask(taskId); if (task && task.retryCount < task.maxRetries) { await updateTaskStatus(taskId, TaskStatus.RETRYING); channel.nack(msg, false, true); console.error(`任务 ${taskId} 处理失败,将重试: ${errorMsg}`); } else { await updateTaskStatus(taskId, TaskStatus.FAILED, { error: errorMsg }); channel.nack(msg, false, false); console.error(`任务 ${taskId} 彻底失败,已转入死信队列: ${errorMsg}`); } } } export async function startConsumers(): Promise<void> { const channel = getChannel(); // 关键:prefetch 必须在 consume 之前设置 channel.prefetch(1); channel.consume(QUEUES.TASK_HIGH, processMessage, { noAck: false }); console.error(`已开始消费高优先级队列: ${QUEUES.TASK_HIGH}`); channel.consume(QUEUES.TASK_DEFAULT, processMessage, { noAck: false }); console.error(`已开始消费默认队列: ${QUEUES.TASK_DEFAULT}`); } export function registerBuiltinHandlers(): void { registerHandler("send_email", async (task) => { const { to, subject } = task.payload; await new Promise((resolve) => setTimeout(resolve, 2000)); console.error(`邮件已发送: ${to} - ${subject}`); return { sent: true, to, subject }; }); registerHandler("batch_process", async (task) => { const { dataset, operation } = task.payload; await new Promise((resolve) => setTimeout(resolve, 5000)); return { processed: dataset.length, operation }; }); registerHandler("generate_report", async (task) => { const { reportType } = task.payload; await new Promise((resolve) => setTimeout(resolve, 3000)); return { reportUrl: `https://reports.example.com/${Date.now()}.pdf` }; }); }MCP Server 入口src/index.ts,定义工具列表和调用处理:
#!/usr/bin/env node import { Server } from "@modelcontextprotocol/sdk/server/index.js"; import { StdioServerTransport } from "@modelcontextprotocol/sdk/server/stdio.js"; import { CallToolRequestSchema, ListToolsRequestSchema, } from "@modelcontextprotocol/sdk/types.js"; import { initRabbitMQ } from "./rabbitmq.js"; import { initRedis, getTask, getTasksByStatus, getTaskStats, TaskStatus, } from "./task_store.js"; import { submitTask, publishEvent, broadcastMessage, getDeadLetterCount, } from "./publisher.js"; import { startConsumers, registerBuiltinHandlers } from "./consumer.js"; const server = new Server( { name: "mq-mcp-server", version: "1.0.0" }, { capabilities: { tools: {} } } ); server.setRequestHandler(ListToolsRequestSchema, async () => { return { tools: [ { name: "submit_task", description: "提交一个异步任务到消息队列。任务会在后台异步执行,不会阻塞当前调用。返回任务 ID 用于后续查询状态。", inputSchema: { type: "object", properties: { type: { type: "string", description: "任务类型,如 send_email、batch_process、generate_report", }, payload: { type: "object", description: "任务参数,根据任务类型不同而不同", }, priority: { type: "number", description: "优先级 0-10,0 为普通,5 以上为高优先级,默认 0", default: 0, }, }, required: ["type", "payload"], }, }, { name: "get_task_status", description: "查询指定任务的当前状态和处理结果。", inputSchema: { type: "object", properties: { taskId: { type: "string", description: "任务 ID" }, }, required: ["taskId"], }, }, { name: "list_tasks", description: "按状态列出任务列表。", inputSchema: { type: "object", properties: { status: { type: "string", enum: ["pending", "processing", "completed", "failed", "retrying"], }, limit: { type: "number", default: 20 }, }, required: ["status"], }, }, { name: "publish_event", description: "发布一个事件到事件总线。", inputSchema: { type: "object", properties: { eventType: { type: "string", description: "如 order.created" }, data: { type: "object" }, }, required: ["eventType", "data"], }, }, { name: "broadcast", description: "广播一条消息给所有在线的消费者。", inputSchema: { type: "object", properties: { message: { type: "string" }, }, required: ["message"], }, }, { name: "get_task_stats", description: "获取任务队列的统计信息。", inputSchema: { type: "object", properties: {} }, }, { name: "check_dead_letters", description: "查看死信队列中的消息数量。", inputSchema: { type: "object", properties: {} }, }, ], }; }); server.setRequestHandler(CallToolRequestSchema, async (request) => { const { name, arguments: args } = request.params; try { switch (name) { case "submit_task": { const { type, payload, priority = 0 } = args as any; const task = await submitTask(type, payload, priority); return { content: [ { type: "text", text: JSON.stringify( { success: true, taskId: task.id, status: task.status, message: `任务已提交,ID 为 ${task.id}。使用 get_task_status 查询进度。`, }, null, 2 ), }, ], }; } case "get_task_status": { const { taskId } = args as any; const task = await getTask(taskId); if (!task) { return { content: [ { type: "text", text: JSON.stringify({ success: false, error: `任务 ${taskId} 不存在或已过期` }), }, ], isError: true, }; } return { content: [ { type: "text", text: JSON.stringify( { success: true, task: { id: task.id, type: task.type, status: task.status, retryCount: task.retryCount, result: task.result || null, error: task.error || null, }, }, null, 2 ), }, ], }; } case "list_tasks": { const { status, limit = 20 } = args as any; const tasks = await getTasksByStatus(status as TaskStatus, limit); return { content: [ { type: "text", text: JSON.stringify({ success: true, count: tasks.length, tasks }, null, 2), }, ], }; } case "publish_event": { const { eventType, data } = args as any; await publishEvent(eventType, data); return { content: [ { type: "text", text: JSON.stringify({ success: true, eventType, message: `事件 ${eventType} 已发布` }), }, ], }; } case "broadcast": { const { message } = args as any; await broadcastMessage(message); return { content: [ { type: "text", text: JSON.stringify({ success: true, message: "广播消息已发送" }) }, ], }; } case "get_task_stats": { const stats = await getTaskStats(); const dlqCount = await getDeadLetterCount(); return { content: [ { type: "text", text: JSON.stringify({ success: true, stats: { ...stats, deadLetters: dlqCount } }, null, 2), }, ], }; } case "check_dead_letters": { const dlqCount = await getDeadLetterCount(); return { content: [ { type: "text", text: JSON.stringify({ success: true, totalDeadLetters: dlqCount }, null, 2), }, ], }; } default: return { content: [{ type: "text", text: `未知工具: ${name}` }], isError: true, }; } } catch (error) { return { content: [ { type: "text", text: `执行出错: ${error instanceof Error ? error.message : String(error)}`, }, ], isError: true, }; } }); async function main() { await initRedis(); await initRabbitMQ(); registerBuiltinHandlers(); await startConsumers(); const transport = new StdioServerTransport(); await server.connect(transport); console.error("消息队列 MCP Server 已启动"); } main().catch((error) => { console.error("Server 启动失败:", error); process.exit(1); });MCP 客户端配置片段(以 Claude Desktop 的claude_desktop_config.json为例):
{ "mcpServers": { "mq-mcp-server": { "command": "node", "args": ["/path/to/mq-mcp-server/dist/index.js"], "env": { "RABBITMQ_URL": "amqp://guest:guest@localhost:5672", "REDIS_URL": "redis://localhost:6379", "TAOTOKEN_API_KEY": "sk-你的Key", "TAOTOKEN_BASE_URL": "https://taotoken.net/api" } } } }如果你用的是 Cline 或 Claude Code,配置逻辑一样,把command、args、env三件套填进去即可。注意 Base URL 用https://taotoken.net/api,Key 用控制台创建的那串,Model ID 按文档填。
4. 验证请求:一条消息触发异步任务并确认执行结果
环境准备好后,先启动依赖服务。用 Docker Compose 最省事:
version: "3.8" services: rabbitmq: image: rabbitmq:3.12-management ports: - "5672:5672" - "15672:15672" environment: RABBITMQ_DEFAULT_USER: guest RABBITMQ_DEFAULT_PASS: guest redis: image: redis:7-alpine ports: - "6379:6379"docker compose up -d之后,RabbitMQ 管理界面在http://localhost:15672,默认账号密码都是 guest。
编译并启动 MCP Server:
npm install npm run build node dist/index.js看到RabbitMQ 初始化完成、Redis 连接已建立、已开始消费默认队列就说明服务起来了。
现在验证完整链路。在 MCP 客户端里调用submit_task工具,参数如下:
{ "type": "send_email", "payload": { "to": "user@example.com", "subject": "订单发货通知", "body": "您的订单已发货" }, "priority": 0 }返回结果:
{ "success": true, "taskId": "task_1735000000000_a1b2c3d4", "status": "pending", "message": "任务已提交,ID 为 task_1735000000000_a1b2c3d4。使用 get_task_status 查询进度。" }拿到 taskId 后,立刻调用get_task_status:
{ "taskId": "task_1735000000000_a1b2c3d4" }第一次查询可能返回processing,因为消费者有 2 秒模拟延迟。等 3 秒再查:
{ "success": true, "task": { "id": "task_1735000000000_a1b2c3d4", "type": "send_email", "status": "completed", "retryCount": 0, "result": { "sent": true, "to": "user@example.com", "subject": "订单发货通知" }, "error": null } }状态变成completed,result 里有发送结果,说明整条链路通了。同时 MCP Server 的控制台会打印任务 task_xxx 处理完成。
再验证事件驱动。调用publish_event:
{ "eventType": "order.created", "data": { "orderId": "ORD-20250101-001", "amount": 299.00 } }返回事件 order.created 已发布。因为事件队列绑定了event.#路由键,所有以event.开头的事件都会被投递到事件队列。你可以在消费者里注册事件处理器来响应。
最后看统计。调用get_task_stats:
{ "success": true, "stats": { "pending": 0, "processing": 0, "completed": 1, "failed": 0, "retrying": 0, "deadLetters": 0 } }completed 为 1,deadLetters 为 0,符合预期。
5. 本篇常见错排查:401、local proxy failed、reading choices、OAuth
做 MCP Server 接入时,报错集中在几个地方。下面按真实报错对照排查。
401 Unauthorized。这个最常见,出现在模型调用侧。原因通常是 API Key 没配、配错、或者 Base URL 写成了带路径的地址。检查三件套:Base URL 必须是https://taotoken.net/api,不要多加/v1或结尾斜杠;Key 必须是控制台创建的那串sk-开头;Model ID 必须和文档里一致。如果用的是 Claude Code,检查settings.json里的env字段;如果是 Cline,检查 MCP 配置里的env。改完重启客户端。
local proxy failed。这个报错通常出现在客户端尝试连接 MCP Server 时。原因可能是 MCP Server 进程没起来、路径写错、或者 Node 版本不对。先手动node dist/index.js跑一遍,看有没有报错。如果报Cannot find module,说明npm run build没成功,检查 TypeScript 编译输出。如果报ECONNREFUSED,说明 RabbitMQ 或 Redis 没启动,docker compose ps看一下容器状态。
reading 'choices'。这个报错出现在解析模型响应时,通常是响应体不是预期的 JSON 格式。原因可能是 Base URL 配错导致请求打到了错误端点,或者 Key 无效返回了错误页。用 curl 直接测一下:
curl -X POST https://taotoken.net/api/v1/messages \ -H "Content-Type: application/json" \ -H "x-api-key: sk-你的Key" \ -H "anthropic-version: 2023-06-01" \ -d '{ "model": "claude-sonnet-4-20250514", "max_tokens": 100, "messages": [{"role": "user", "content": "hello"}] }'如果返回正常 JSON,说明 Key 和 Base URL 没问题,问题在客户端配置。如果返回 401 或 404,按上面 401 的排查走。
OAuth 相关报错。如果你用的是 Claude Code 的 OAuth 登录流程,报OAuth token expired或invalid_grant,说明 token 过期了。重新走一遍登录流程,或者在配置里改用 API Key 方式。TaoToken 的 API Key 方式不需要 OAuth,直接填 Key 就行,省去 token 刷新的麻烦。
消息重复消费。这个不是报错但很坑。现象是同一个任务被执行多次。根因是消费者在处理完但还没 ack 时进程挂了,RabbitMQ 会把消息重新投递。解决办法是在processMessage开头做幂等性检查,查 Redis 里任务状态,如果已经是completed就直接 ack 跳过。上面的代码里已经加了这段逻辑。
prefetch 不生效。现象是队列里消息被一次性全推给消费者,unacked 数量暴涨。根因是channel.prefetch(1)写在了channel.consume之后。RabbitMQ 的 QoS 设置必须在消费开始前生效。把 prefetch 挪到 consume 前面就行,代码里已经修正。
死信队列堆积。调用check_dead_letters发现数量一直涨。说明有任务反复失败超过 maxRetries。先看任务失败原因,调用list_tasks传status: "failed"查失败任务,看 error 字段。常见原因是任务处理器抛异常、payload 格式不对、或者依赖的外部服务不可用。修好之后,可以手动重试死信消息。
6. 把消息队列 MCP Server 接进你的 AI 工作流
到这里,一条消息从 AI 客户端发出、经 MCP Server 提交到 RabbitMQ、被消费者处理、状态写回 Redis、再被 AI 查询到结果的完整链路就跑通了。核心就三件事:耗时操作异步化避免 MCP 调用超时,死信队列兜底保证任务不丢,消费者幂等性检查防止重复执行。
如果你要长期跑这类 Agent 任务,模型调用量会比较大,可以看下 Coding Plan:https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_campaign=rewrite&utm_content= 。它适合需要持续调用模型、跑长任务的场景。
需要创建新的 API Key 或者管理现有 Key,去控制台:https://taotoken.net/api-keys?utm_source=taotoken_aicg_blog_end&utm_campaign=rewrite&utm_content= 。接入细节和 Model ID 列表在文档:https://taotoken.net/doc?utm_source=taotoken_aicg_blog_end&utm_campaign=rewrite&utm_content= 。想先试试模型对话效果,可以直接用:https://taotoken.net/chat?utm_source=taotoken_aicg_blog_end&utm_campaign=rewrite&utm_content= 。
实际用下来,消息队列 MCP Server 最大的价值是把 AI 的能力边界从“即时返回”扩展到了“管理后台任务”。AI 提交一个批处理任务,过几分钟回来查结果,整个流程很自然。下一步你可以在这个骨架上加自己的任务处理器,比如对接内部工单系统、触发 CI/CD 流水线、或者做跨系统的数据同步。