这是一个使用 Go 实现的 Apache Kafka 的 Model Context Protocol (MCP) 服务器,利用了 franz-go 和 mcp-go。
该服务器提供了一个通过 MCP 协议与 Kafka 进行交互的实现,使 LLM 模型能够通过标准化接口执行常见的 Kafka 操作。
[
](https://github.com/tu
uannvm/kafka-mcp-server/actions/workflows/build.yml)
Kafka MCP 服务器弥合了 LLM 模型和 Apache Kafka 之间的差距,允许它们:
所有这些都通过标准化的 Model Context Protocol (MCP) 来完成。
graph TB
subgraph "MCP 客户端(AI 应用程序)"
A[Claude Desktop]
B[Cursor]
C[Windsurf]
D[ChatWise]
end
subgraph "Kafka MCP 服务器"
E[MCP 协议处理器]
F[工具注册表]
G[资源注册表]
H[提示注册表]
I[Kafka 客户端包装器]
end
subgraph "Apache Kafka 集群"
J[代理 1]
K[代理 2]
L[代理 3]
M[主题及分区]
N[消费者组]
end
A --> E
B --> E
C --> E
D --> E
E --> F
E --> G
E --> H
F --> I
G --> I
H --> I
I --> J
I --> K
I --> L
J --> M
K --> M
L --> M
J --> N
K --> N
L --> N
classDef 客户端 fill:#e1f5fe
classDef mcp fill:#f3e5f5
classDef kafka fill:#fff3e0
class A,B,C,D 客户端
class E,F,G,H,I mcp
class J,K,L,M,N kafka
工作原理:
传输模式:
安装 kafka-mcp-server 最简单的方法是使用 Homebrew:
# 添加 tap 存储库
brew tap tuannvm/mcp
# 安装 kafka-mcp-server
brew install kafka-mcp-server
更新到最新版本:
brew update && brew upgrade kafka-mcp-server
# 克隆仓库
git clone https://github.com/tuannvm/kafka-mcp-server.git
cd kafka-mcp-server
# 构建服务器
go build -o kafka-mcp-server ./cmd
此 MCP 服务器可以与多个 AI 应用程序集成。以下是特定平台的说明:
编辑 ~/.cursor/mcp.json 并添加 kafka-mcp-server 配置:
{
"mcpServers": {
"kafka": {
"command": "kafka-mcp-server",
"args": [],
"env": {
"KAFKA_BROKERS": "localhost:9092",
"KAFKA_CLIENT_ID": "kafka-mcp-server",
"MCP_TRANSPORT": "stdio"
}
}
}
}
编辑您的 Claude 配置文件并添加服务器:
~/Library/Application Support/Claude/claude_desktop_config.json{
"mcpServers": {
"kafka": {
"command": "kafka-mcp-server",
"args": [],
"env": {
"KAFKA_BROKERS": "localhost:9092",
"KAFKA_CLIENT_ID": "kafka-mcp-server",
"MCP_TRANSPORT": "stdio"
}
}
}
}
重启 Claude Desktop 以应用更改。
要与 Claude Code 一起使用,请使用内置的 MCP 配置命令添加服务器:
# 使用环境变量添加 kafka-mcp-server
claude mcp add kafka \
--env KAFKA_BROKERS=localhost:9092 \
--env KAFKA_CLIENT_ID=kafka-mcp-server \
--env MCP_TRANSPORT=stdio \
--env KAFKA_SASL_MECHANISM= \
--env KAFKA_SASL_USER= \
--env KAFKA_SASL_PASSWORD= \
--env KAFKA_TLS_ENABLE=false \
-- kafka-mcp-server
其他有用命令:
# 列出已配置的 MCP 服务器
claude mcp list
# 删除服务器
claude mcp remove kafka
# 测试服务器连接
claude mcp get kafka
kafkakafka-mcp-serverKAFKA_BROKERS=localhost:9092
KAFKA_CLIENT_ID=kafka-mcp-server
MCP_TRANSPORT=stdio
跨多个客户端管理 MCP 服务器配置可能会变得具有挑战性。mcpenetes 是一个专门的工具,使这个过程显著更容易:
# 安装 mcpenetes
go install github.com/tuannvm/mcpenetes@latest
# 搜索可用的 MCP 服务器,包括 kafka-mcp-server
mcpenetes search
# 将 kafka-mcp-server 配置一次性应用于所有客户端
mcpenetes apply
# 从剪贴板加载配置
mcpenetes load
使用 mcpenetes,您可以维护多个 Kafka 配置(开发、生产等),并在所有客户端(Cursor、Claude Desktop、Windsurf、ChatWise)之间即时切换,而无需手动编辑每个客户端的配置文件。
服务器公开以下工具用于 Kafka 交互。有关详细文档,包括示例和样本响应,请参阅 docs/tools.md。
服务器提供了可以通过 MCP 协议访问的以下资源。有关详细文档,包括示例响应,请参阅 docs/resources.md。
服务器包含以下预配置提示,用于 Kafka 操作和诊断。有关详细文档,包括参数和示例响应,请参阅 docs/prompts.md。
这里是一个交互示例,展示了 AI 助手如何使用这些工具帮助用户:
用户查询:"我们需要调查为什么订单处理管道落后了。你能帮我检查一下消费者的滞后吗?"
AI 助手的工作流程:
首先,检查消费者组
> 使用 list_consumer_groups 工具
> 发现 "order-processor" 组
然后,查看消费者组详情
> 使用 describe_consumer_group 工具,group_id 为 "order-processor",include_offsets=true
> 发现几个分区有滞后
获取集群概述以检查总体健康状况
> 使用 cluster_overview 工具
> 确认没有未充分复制的分区或离线代理
生成消费者滞后报告
> 使用 consumer_lag_report 提示,threshold=1000
> 提供消费者组滞后情况的详细分析
提供行动建议:
根据我的分析,您的 order-processor 消费者组在主题 orders 的分区 2 和 3 上有显著滞后,滞后值分别为 15,420 和 12,305 条消息。
集群本身是健康的,没有未充分复制的分区,这表明这是消费者应用程序的问题,而不是 Kafka 基础设施的问题。
我建议:
1. 检查处理这些分区的消费者实例是否有减速或错误
2. 考虑通过增加更多实例来扩展您的消费者组
3. 审查消费者配置,特别是 max.poll.records 和 max.poll.interval.ms
4. 检查消费者逻辑中的任何处理瓶颈
这种无缝的工作流程展示了 Kafka MCP 工具如何使 LLM 模型执行复杂的诊断并提供行动建议。
服务器可以通过以下环境变量进行配置:
| 变量 | 描述 | 默认值 |
|---|---|---|
KAFKA_BROKERS | 逗号分隔的 Kafka 代理地址列表 | localhost:9092 |
KAFKA_CLIENT_ID | 用于连接的 Kafka 客户端 ID | kafka-mcp-server |
MCP_TRANSPORT | MCP 传输方法(stdio/http) | stdio |
KAFKA_SASL_MECHANISM | SASL 机制:plain、scram-sha-256、scram-sha-512 或 ""(禁用) | "" |
KAFKA_SASL_USER | SASL 认证的用户名 | "" |
KAFKA_SASL_PASSWORD | SASL 认证的密码 | "" |
KAFKA_TLS_ENABLE | 启用 Kafka 连接的 TLS(true 或 false) | false |
KAFKA_TLS_INSECURE_SKIP_VERIFY | 跳过 TLS 证书验证(true 或 false) | false |
当使用 HTTP 传输(MCP_TRANSPORT=http)时,可以启用 OAuth 2.1 认证:
| 变量 | 描述 | 默认值 | 是否必需 |
|---|---|---|---|
MCP_HTTP_PORT | HTTP 服务器端口 | 8080 | 否 |
OAUTH_ENABLED | 启用 OAuth 2.1 认证 | false | 否 |
OAUTH_MODE | OAuth 模式:native 或 proxy | native | 否 |
OAUTH_PROVIDER | 提供商:hmac、okta、google、azure | okta | 否 |
OAUTH_SERVER_URL | 完整的 MCP 服务器 URL(例如 https://localhost:8080) | - | 当启用 OAuth 时 |
OIDC_ISSUER | OAuth 发行者 URL | - | 当启用 OAuth 时 |
OIDC_AUDIENCE | OAuth 受众 | - | 当启用 OAuth 时 |
OIDC_CLIENT_ID | OAuth 客户端 ID | - | 仅限代理模式 |
OIDC_CLIENT_SECRET | OAuth 客户端密钥 | - | 仅限代理模式 |
OAUTH_REDIRECT_URIS | 逗号分隔的重定向 URI | - | 仅限代理模式 |
JWT_SECRET | JWT 签名密钥 | - | 仅限代理模式 |
关于详细的 OAuth 设置和示例,请参阅 docs/oauth.md。
安全注意事项:
- 当使用
KAFKA_TLS_INSECURE_SKIP_VERIFY=true时,服务器将跳过 TLS 证书验证。这仅应在开发或测试环境中使用,或者在使用自签名证书时使用。- OAuth 仅在使用 HTTP 传输时可用。STDIO 传输不支持 OAuth。