返回市场
kafka麦普服务器

kafka麦普服务器

作者:aswinayyolath2 星标更新:2025-08-04

项目介绍

Kafka MCP 服务器

一个全面的模型上下文协议(MCP)服务器,用于Apache Kafka操作,支持与Claude Desktop和其他MCP客户端无缝集成。

🚀 功能

  • 集群管理:监控集群健康、代理信息和元数据
  • 主题操作:创建、列出、描述和删除Kafka主题
  • 消息操作:发送和消费具有灵活配置的消息
  • 消费者组管理:列出和描述消费者组
  • 实时监控:获取主题指标、偏移量和性能数据
  • 健康检查:全面的集群健康监控
  • 安全配置:支持SASL、SSL和各种身份验证方法

📋 先决条件

  • Python 3.8+
  • Apache Kafka集群(本地或远程)
  • pip包管理器

🛠️ 安装

  1. 克隆仓库

    git clone https://github.com/aswinayyolath/kafka-mcp-server.git
    cd kafka-mcp-server
    
  2. 创建虚拟环境

    python -m venv venv
    source venv/bin/activate  # 在Windows上:venv\Scripts\activate
    
  3. 安装依赖项

    pip install mcp kafka-python
    
  4. 配置环境变量

    cp .env.template .env
    # 使用您的Kafka配置编辑.env文件
    

⚙️ 配置

环境变量

基于.env.template创建一个.env文件:

# 基本的Kafka配置
KAFKA_BOOTSTRAP_SERVERS=localhost:9092
KAFKA_SECURITY_PROTOCOL=PLAINTEXT

# 如果需要,配置SASL
KAFKA_SASL_MECHANISM=PLAIN
KAFKA_SASL_USERNAME=your_username
KAFKA_SASL_PASSWORD=your_password

# 如果需要,配置SSL
KAFKA_SSL_CAFILE=/path/to/ca.pem
KAFKA_SSL_CERTFILE=/path/to/cert.pem
KAFKA_SSL_KEYFILE=/path/to/key.pem

Docker设置(可选)

使用Docker启动本地Kafka集群:

docker-compose up -d

这将在localhost:9092启动Kafka。

🧪 测试

运行综合测试套件以验证您的设置:

# 设置环境变量
export KAFKA_BOOTSTRAP_SERVERS=localhost:9092
export KAFKA_SECURITY_PROTOCOL=PLAINTEXT

# 运行测试
python test_kafka_mcp.py

测试套件验证:

  • Kafka连接性
  • MCP服务器启动
  • 所有可用工具
  • 主题操作
  • 消息操作
  • 资源和提示

🔧 使用

单独服务器

直接运行MCP服务器:

python kafka_mcp_server.py

Claude Desktop集成

在您的Claude Desktop配置中添加(claude_desktop_config.json):

{
  "mcpServers": {
    "kafka": {
      "command": "python",
      "args": ["/path/to/kafka_mcp_server.py"],
      "env": {
        "KAFKA_BOOTSTRAP_SERVERS": "localhost:9092",
        "KAFKA_SECURITY_PROTOCOL": "PLAINTEXT"
      }
    }
  }
}

🛠️ 可用工具

集群操作

  • get_cluster_info - 获取全面的集群信息
  • health_check - 执行集群健康检查

主题操作

  • list_topics - 列出集群中的所有主题
  • create_topic - 创建新主题
  • describe_topic - 获取详细的主题信息
  • delete_topic - 删除主题
  • get_topic_metrics - 获取全面的主题指标
  • get_topic_offsets - 获取分区偏移量

消息操作

  • send_message - 发送单个消息
  • send_batch_messages - 发送多条消息
  • consume_messages - 从主题中消费消息

消费者组操作

  • list_consumer_groups - 列出所有消费者组
  • describe_consumer_group - 获取详细的消费者组信息

📚 资源

该服务器提供MCP资源,便于访问集群信息:

  • kafka://cluster/info - 实时集群信息
  • kafka://topics/list - 当前主题列表

🎯 提示

内置提示用于常见场景:

  • kafka_monitoring_prompt - 生成监控和故障排除指南
  • kafka_troubleshooting_prompt - 获取特定问题的帮助

💡 示例使用与Claude

一旦与Claude Desktop集成,您可以询问:

  • "你能检查我的Kafka集群的健康状况吗?"
  • "列出所有主题及其分区数量"
  • "创建名为'user-events'的主题,包含3个分区"
  • "向user-events主题发送测试消息"
  • "显示user-events主题最后10条消息"
  • "哪些消费者组是活跃的?"

🔒 安全

认证

服务器支持多种Kafka认证方法:

  • PLAINTEXT:无认证(仅限开发)
  • SASL_PLAINTEXT:通过明文连接进行SASL认证
  • SASL_SSL:通过SSL进行SASL认证
  • SSL:SSL客户端证书认证

最佳实践

  1. 永远不要提交.env文件 - 使用.env.template作为示例
  2. 生产环境中使用SSL - 总是对生产集群的连接进行加密
  3. 限制权限 - 使用具有最小必要权限的专用服务账户
  4. 监控访问 - 记录并监控MCP服务器的使用情况

🐛 故障排除

常见问题

  1. 连接被拒绝

    • 验证Kafka是否正在运行
    • 检查KAFKA_BOOTSTRAP_SERVERS配置
    • 确保网络连接
  2. 认证失败

    • 验证SASL凭据
    • 检查SSL证书路径
    • 验证安全协议设置
  3. 主题未找到

    • 确认主题存在
    • 检查主题名称拼写
    • 验证权限

调试模式

通过设置启用调试日志:

export PYTHONPATH=.
python -c "import logging; logging.basicConfig(level=logging.DEBUG)"
python kafka_mcp_server.py

🤝 贡献

  1. 分叉仓库
  2. 创建功能分支:git checkout -b feature-name
  3. 进行更改
  4. 运行测试:python test_kafka_mcp.py
  5. 提交更改:git commit -am '添加功能'
  6. 推送到分支:git push origin feature-name
  7. 提交拉取请求

📄 许可

该项目根据Apache许可证2.0发布 - 查看LICENSE文件了解详情。

🙏 致谢

📞 支持

对于问题和疑问:

  1. 查看故障排除部分
  2. 运行测试套件以验证您的设置
  3. 查阅Kafka和MCP文档
  4. 提交带有详细错误信息的问题

祝您使用MCP进行愉快的Kafka流处理!🎉