返回市场
kafka麦普服务器

kafka麦普服务器

作者:tuannvm37 星标更新:2025-11-17

项目介绍

Kafka MCP 服务器

这是一个使用 Go 实现的 Apache Kafka 的 Model Context Protocol (MCP) 服务器,利用了 franz-gomcp-go

该服务器提供了一个通过 MCP 协议与 Kafka 进行交互的实现,使 LLM 模型能够通过标准化接口执行常见的 Kafka 操作。

Go 报告卡 [GitHub 工作流状态](https://github.com/tu uannvm/kafka-mcp-server/actions/workflows/build.yml) Go 版本 Trivy 扫描 SLSA 3 Go 参考 Docker 镜像 GitHub 发布 许可证:MIT

概述

Kafka MCP 服务器弥合了 LLM 模型和 Apache Kafka 之间的差距,允许它们:

  • 从主题中生产和消费消息
  • 列出、描述和管理主题
  • 监控和管理消费者组
  • 评估集群健康状况和配置
  • 执行标准的 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

工作原理:

  1. MCP 客户端(AI 应用程序)通过 stdio 或 HTTP 传输连接到 Kafka MCP 服务器
  2. MCP 服务器暴露三种类型的能力:
    • 工具 - 直接 Kafka 操作(生产/消费消息、描述主题等)
    • 资源 - 集群健康报告和诊断
    • 提示 - 常见操作的预配置工作流程
  3. Kafka 客户端包装器使用 franz-go 库处理所有 Kafka 通信
  4. Apache Kafka 集群处理实际的消息流和存储

传输模式:

  • STDIO:默认模式,适用于本地 MCP 客户端(Claude Desktop、Cursor 等)
  • HTTP:启用远程访问,可选 OAuth 2.1 认证

工具

提示和资源

主要功能

  • Kafka 集成:通过 MCP 实现常见的 Kafka 操作
  • 安全性
    • 支持 SASL(PLAIN、SCRAM-SHA-256、SCRAM-SHA-512)和 TLS 认证
    • HTTP 传输的 OAuth 2.1 认证(原生和代理模式)
    • 支持 Okta、Google、Azure AD 和 HMAC 提供商
  • 灵活的传输:STDIO 用于本地客户端,HTTP 用于远程访问
  • 错误处理:具有有意义反馈的错误处理
  • 配置选项:可定制以适应不同的环境
  • 预配置提示:一组用于常见 Kafka 操作的提示
  • 兼容性:与兼容 MCP 的 LLM 模型一起工作

快速开始

先决条件

  • Go 1.24 或更高版本
  • Docker(用于运行集成测试)
  • 对 Kafka 集群的访问

安装

Homebrew(macOS 和 Linux)

安装 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 客户端集成

此 MCP 服务器可以与多个 AI 应用程序集成。以下是特定平台的说明:

Cursor

编辑 ~/.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 Desktop

编辑您的 Claude 配置文件并添加服务器:

  • macOS~/Library/Application Support/Claude/claude_desktop_config.json
  • Windows:%APPDATA%\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

要与 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

ChatWise

  1. 打开 ChatWise → 设置 → 工具 → "+" → "命令行 MCP"
  2. 配置:
    • IDkafka
    • 命令kafka-mcp-server
    • 参数:(留空)
    • 环境变量:添加环境变量:
      KAFKA_BROKERS=localhost:9092
      KAFKA_CLIENT_ID=kafka-mcp-server
      MCP_TRANSPORT=stdio
      

使用 mcpenetes 简化配置

跨多个客户端管理 MCP 服务器配置可能会变得具有挑战性。mcpenetes 是一个专门的工具,使这个过程显著更容易:

# 安装 mcpenetes
go install github.com/tuannvm/mcpenetes@latest

主要功能

  • 交互式搜索:通过简单的命令查找和选择 Kafka MCP 服务器配置
  • 到处应用:自动同步配置到您所有的 MCP 客户端
  • 配置备份:在进行更改之前安全地备份现有配置
  • 恢复:如果需要,轻松回滚到之前的配置

使用 mcpenetes 快速开始

# 搜索可用的 MCP 服务器,包括 kafka-mcp-server
mcpenetes search 

# 将 kafka-mcp-server 配置一次性应用于所有客户端
mcpenetes apply

# 从剪贴板加载配置
mcpenetes load

使用 mcpenetes,您可以维护多个 Kafka 配置(开发、生产等),并在所有客户端(Cursor、Claude Desktop、Windsurf、ChatWise)之间即时切换,而无需手动编辑每个客户端的配置文件。

MCP 工具

服务器公开以下工具用于 Kafka 交互。有关详细文档,包括示例和样本响应,请参阅 docs/tools.md

  • produce_message:向 Kafka 主题发送消息
  • consume_messages:批量操作从 Kafka 主题消费消息
  • list_brokers:列出所有配置的 Kafka 代理地址
  • describe_topic:提供特定主题的综合元数据
  • list_consumer_groups:枚举集群中的所有消费者组
  • describe_consumer_group:提供详细的消费者组信息,包括滞后指标
  • describe_configs:检索 Kafka 资源的配置设置
  • cluster_overview:提供全面的集群健康概述
  • list_topics:列出所有主题及其元数据,包括分区和复制信息

MCP 资源

服务器提供了可以通过 MCP 协议访问的以下资源。有关详细文档,包括示例响应,请参阅 docs/resources.md

  • kafka-mcp://overview:全面的集群健康概述
  • kafka-mcp://health-check:带有行动建议的详细健康检查
  • kafka-mcp://under-replicated-partitions:分析复制问题的分区
  • kafka-mcp://consumer-lag-report:带有自定义阈值的消费者性能分析

MCP 提示

服务器包含以下预配置提示,用于 Kafka 操作和诊断。有关详细文档,包括参数和示例响应,请参阅 docs/prompts.md

  • kafka_cluster_overview:生成全面的集群健康概述
  • kafka_health_check:执行带有行动建议的详细健康检查
  • kafka_under_replicated_partitions:分析复制问题并提供故障排除指导
  • kafka_consumer_lag_report:提供消费者滞后分析,并提出性能优化建议

端到端示例

这里是一个交互示例,展示了 AI 助手如何使用这些工具帮助用户:

用户查询:"我们需要调查为什么订单处理管道落后了。你能帮我检查一下消费者的滞后吗?"

AI 助手的工作流程:

  1. 首先,检查消费者组

    > 使用 list_consumer_groups 工具
    > 发现 "order-processor" 组
    
  2. 然后,查看消费者组详情

    > 使用 describe_consumer_group 工具,group_id 为 "order-processor",include_offsets=true
    > 发现几个分区有滞后
    
  3. 获取集群概述以检查总体健康状况

    > 使用 cluster_overview 工具
    > 确认没有未充分复制的分区或离线代理
    
  4. 生成消费者滞后报告

    > 使用 consumer_lag_report 提示,threshold=1000
    > 提供消费者组滞后情况的详细分析
    
  5. 提供行动建议:

    根据我的分析,您的 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 客户端 IDkafka-mcp-server
MCP_TRANSPORTMCP 传输方法(stdio/http)stdio
KAFKA_SASL_MECHANISMSASL 机制:plainscram-sha-256scram-sha-512""(禁用)""
KAFKA_SASL_USERSASL 认证的用户名""
KAFKA_SASL_PASSWORDSASL 认证的密码""
KAFKA_TLS_ENABLE启用 Kafka 连接的 TLS(truefalsefalse
KAFKA_TLS_INSECURE_SKIP_VERIFY跳过 TLS 证书验证(truefalsefalse

OAuth 2.1 配置(仅限 HTTP 传输)

当使用 HTTP 传输(MCP_TRANSPORT=http)时,可以启用 OAuth 2.1 认证:

变量描述默认值是否必需
MCP_HTTP_PORTHTTP 服务器端口8080
OAUTH_ENABLED启用 OAuth 2.1 认证false
OAUTH_MODEOAuth 模式:nativeproxynative
OAUTH_PROVIDER提供商:hmacoktagoogleazureokta
OAUTH_SERVER_URL完整的 MCP 服务器 URL(例如 https://localhost:8080-当启用 OAuth 时
OIDC_ISSUEROAuth 发行者 URL-当启用 OAuth 时
OIDC_AUDIENCEOAuth 受众-当启用 OAuth 时
OIDC_CLIENT_IDOAuth 客户端 ID-仅限代理模式
OIDC_CLIENT_SECRETOAuth 客户端密钥-仅限代理模式
OAUTH_REDIRECT_URIS逗号分隔的重定向 URI-仅限代理模式
JWT_SECRETJWT 签名密钥-仅限代理模式

关于详细的 OAuth 设置和示例,请参阅 docs/oauth.md

安全注意事项:

  • 当使用 KAFKA_TLS_INSECURE_SKIP_VERIFY=true 时,服务器将跳过 TLS 证书验证。这仅应在开发或测试环境中使用,或者在使用自签名证书时使用。
  • OAuth 仅在使用 HTTP 传输时可用。STDIO 传输不支持 OAuth。