跳转至正文

RocketMQ 5 智能化控制面 MCP 服务 ​

GitHub RepoApache RocketMQ 5Spring AI 2.0.1Spring Boot 4.1.0JDK 21npm versionLicense

mcp-server-rocketmq 是专为 Apache RocketMQ 5 打造的企业级智能化连接服务,遵循开放标准 Model Context Protocol (MCP) 规范设计,基于 Spring AI 2 + Spring Boot 4 + JDK 21 现代云原生技术栈构建,全面支持 GraalVM 原生二进制与 npx 免安装秒开。

它将大语言模型智能体(Claude Desktop、Cursor、Antigravity 等各类 AI 编程助手与智能运维 Agent)与分布式消息中间件 Apache RocketMQ 5 深度打通。通过暴露标准化的 MCP Tools(工具)、Resources(只读资源) 与 Prompts(专家排障工作流),赋能智能体通过自然语言直接探查集群健康、秒级定位消费堆积、全链路检索消息、排查死信根因,并安全受控地执行消息收发与位点治理。


1. 系统架构与双模拓扑 ​

本项目采用 Remoting 深度运维管理 + gRPC 云原生消息收发 的混合双驱动架构,深度兼顾底层集群运维探测与高吞吐消息收发:

正在渲染架构图表...

2. 四维生产级安全防护门禁 ​

在连接生产集群时,系统内置了严格的安全铁律,防止大模型发生幻觉或越权误操作:

  1. 只读守卫 (Read-Only Guard):当配置 ROCKETMQ_READ_ONLY=true 时,在握手层物理隐藏所有写操作与破坏性工具,仅暴露指标巡检与消息查询能力。
  2. 双层防呆确认 (Dual-Layer Guard):删除主题、重置消费位点、死信重新投递等破坏性高危动作,不仅需要配置开启 ROCKETMQ_ALLOW_DESTRUCTIVE=true,还必须由模型在调用时显式传入 confirm: true,否则拒绝执行。
  3. 4KB 消息截断防护 (Message Body Guard):大消息体内容直读时强制实施 4KB 上限截断并注入警告标识,彻底消除上下文爆炸与 Token 耗尽风险。
  4. 系统内置主题静默过滤 (System Topic Filter):枚举主题时自动静默过滤 RMQ_SYS_* 等系统内部主题,避免污染模型视觉或被误触。

3. MCP Tools 工具契约字典 (共 18 项) ​

服务完整暴露了 18 项标准化原子工具,涵盖集群、主题、消费组、消息与死信全生命周期:

领域模块工具标识 (Tool Name)核心入参 (Parameters)职责说明安全防护级别
集群拓扑域rocketmq_cluster_info(无入参)查询 NameServer / Broker 节点分布、角色与在线状态只读查询
rocketmq_broker_statsbrokerAddr查询指定 Broker 运行时核心指标(吞吐量、写入 TPS、物理磁盘水位)只读查询
主题生命周期rocketmq_list_topicsincludeSystem?列出集群业务 Topic(默认过滤系统内部管理主题)只读查询
rocketmq_topic_routetopic查询指定 Topic 的读写队列分布与 Broker 路由详情只读查询
rocketmq_topic_statustopic查询指定 Topic 各分片队列的最小/最大 Offset 与堆积容量统计只读查询
rocketmq_create_topictopic, readQueueNums?, writeQueueNums?, perm?声明式创建或更新指定 Topic(动态配置队列数与读写权限)受 read-only 约束
rocketmq_delete_topictopic, confirm: true彻底清理下线指定业务 Topic🚨 双层防呆保护
消费组与积压rocketmq_list_consumer_groupsincludeSystem?获取所有已注册的消费组清单只读查询
rocketmq_consumer_statusconsumerGroup查询消费组的在线客户端 ID、IP 端口及订阅详情只读查询
rocketmq_consumer_lagconsumerGroup, topic?精确计算消费组在各分片队列的未消费堆积量 (Lag)只读查询
rocketmq_top_consumer_lagtopN? (默认 10)全集群积压排行榜:极速检出堆积最严重的 TopN 消费组只读查询
rocketmq_reset_consumer_offsetconsumerGroup, topic, resetType, timestamp?, confirm: true按时间戳回溯或按最大位点跳过重置消费点位🚨 双层防呆保护
消息检索排查rocketmq_query_message_by_idtopic, msgId根据 32 位 Message ID 精确检索消息内容与用户属性只读 (4KB 截断)
rocketmq_query_message_by_keytopic, key, beginTimestamp?, endTimestamp?, maxNum?根据业务 Key 在指定时间窗口内扫描匹配的消息列表只读 (4KB 截断)
rocketmq_query_dlq_messagesconsumerGroup, maxNum?检索指定消费组死信队列(DLQ)中的失败堆积消息只读 (4KB 截断)
rocketmq_query_message_tracemsgId, topic?调阅单条消息自 Producer、Broker 至 Consumer 的全链路轨迹耗时只读 (4KB 截断)
消息生产自愈rocketmq_send_messagetopic, body, tag?, keys?, messageGroup?, deliveryTimestamp?发送测试消息(支持普通、分区顺序与定时延时消息)受 read-only 约束
rocketmq_resend_dlq_messageconsumerGroup, msgId, targetTopic, confirm: true将死信队列中的指定消息重新投递回业务 Topic 触发重试🚨 双层防呆保护

📌 注:加粗表示必填项,带 ? 表示可选参数。


4. MCP Resources 与 Prompts 扩展 ​

4.1 MCP Resources (只读上下文) ​

  • rocketmq://cluster/topology:集群物理拓扑快照(Broker 角色、地址与运行时版本)。
  • rocketmq://topics:当前集群所有业务主题清单与队列分布概览。
  • rocketmq://server/status:服务端自身运行时配置(只读状态、高危开关与限制参数)。

4.2 MCP Prompts (专家预置工作流) ​

  • diagnose-consumer-lag:一键触发消费积压全自动根因排查(识别慢消费节点、拉取 TPS 波动、计算堆积预计耗时)。
  • investigate-dlq-root-cause:死信队列消息深度分析(解码异常堆栈、比对投递时间、提供修复建议)。

5. 主流客户端接入配置 ​

5.1 Claude Desktop 配置 (claude_desktop_config.json) ​

通过 npm 官方分发包 @atengk/mcp-server-rocketmq 零安装秒开:

json
{
  "mcpServers": {
    "rocketmq": {
      "command": "npx",
      "args": ["-y", "@atengk/mcp-server-rocketmq"],
      "env": {
        "ROCKETMQ_NAMESRV_ADDR": "127.0.0.1:9876",
        "ROCKETMQ_READ_ONLY": "false",
        "ROCKETMQ_ALLOW_DESTRUCTIVE": "false"
      }
    }
  }
}

5.2 Cursor 与 Antigravity 配置 ​

在 .cursor/mcp.json 或 ~/.gemini/antigravity/mcp/ 中声明:

json
{
  "mcpServers": {
    "rocketmq-prod": {
      "command": "npx",
      "args": ["-y", "@atengk/mcp-server-rocketmq"],
      "env": {
        "ROCKETMQ_NAMESRV_ADDR": "namesrv.internal:9876",
        "ROCKETMQ_PROXY_ADDR": "proxy.internal:8081",
        "ROCKETMQ_ACCESS_KEY": "${ROCKETMQ_AK}",
        "ROCKETMQ_SECRET_KEY": "${ROCKETMQ_SK}",
        "ROCKETMQ_READ_ONLY": "true"
      }
    }
  }
}

5.3 GraalVM 原生二进制或 Java Jar ​

若需要在未安装 Node.js 的独立生产环境运行预编译原生二进制或 Jar:

json
{
  "mcpServers": {
    "rocketmq": {
      "command": "/opt/mcp/rocketmq-mcp-server",
      "env": {
        "ROCKETMQ_NAMESRV_ADDR": "namesrv.internal:9876",
        "ROCKETMQ_PROXY_ADDR": "proxy.internal:8081",
        "ROCKETMQ_ACCESS_KEY": "${ROCKETMQ_AK}",
        "ROCKETMQ_SECRET_KEY": "${ROCKETMQ_SK}"
      }
    }
  }
}

5.4 HTTP SSE 远程长连接模式 ​

适合部署在云原生 Kubernetes 集群或集中运维网关中:

json
{
  "mcpServers": {
    "rocketmq-remote": {
      "url": "http://rocketmq-mcp.infra.internal:8080/sse"
    }
  }
}

6. 环境变量配置全景 ​

环境变量名默认值作用与规范说明
ROCKETMQ_NAMESRV_ADDR127.0.0.1:9876RocketMQ NameServer 集群地址,多个地址以分号分隔
ROCKETMQ_PROXY_ADDR无RocketMQ 5.x gRPC Proxy 地址(发信与死信重投需要)
ROCKETMQ_ACCESS_KEY无ACL 访问控制鉴权 AccessKey(可选)
ROCKETMQ_SECRET_KEY无ACL 访问控制鉴权 SecretKey(可选)
ROCKETMQ_READ_ONLYfalse全局只读安全防线开关,开启后物理隐藏所有写/删工具
ROCKETMQ_ALLOW_DESTRUCTIVEfalse是否允许执行删除主题与重置位点等高危操作
ROCKETMQ_MAX_MESSAGE_BYTES4096消息直读内容截断上限(字节)
SERVER_PORT8080SSE 远程服务模式下的 HTTP 监听端口

7. 本地运行与 Docker 快速启动 ​

bash
# 1. 本地通过 npx 快速启动测试
npx -y @atengk/mcp-server-rocketmq --namesrv 127.0.0.1:9876

# 2. Docker 单命令启动常驻 SSE 服务
docker run -d \
  --name mcp-server-rocketmq \
  -p 8080:8080 \
  -e ROCKETMQ_NAMESRV_ADDR=host.docker.internal:9876 \
  -e ROCKETMQ_READ_ONLY=true \
  ghcr.io/atengk/mcp-server-rocketmq:latest

8. 相关资源与互链 ​

基于 Apache-2.0 协议开源发布