CDF Kafka MCP服务器
一个支持Apache Knox身份验证的Apache Kafka模型上下文协议(MCP)服务器,灵感来自 SSB MCP服务器 实施。
概述
CDF Kafka MCP服务器通过模型上下文协议在AI模型和Apache Kafka集群之间提供了一个全面的桥梁。它使AI应用程序能够与Kafka主题进行交互,生成和使用消息,并通过一个安全的、企业就绪的Apache Knox身份验证接口管理消费者组。
📸 视觉示例:此README包括显示如何使用SMM(Streams Messaging Manager)web界面创建主题、添加数据和列出主题的屏幕截图。请参阅 视觉示例 有关分步视觉指南的部分。
特性
🔐 企业安全
- Apache Knox网关 管理API集成的身份验证支持
- CDP云 令牌身份验证和基本身份验证支持
- 多种身份验证方法(基于令牌、用户名/密码、CDP令牌)
- 用于安全通信的TLS/SSL配置
- 可配置的SSL证书验证
- 通过Knox网关发现服务
🚀 增强的Kafka操作
- 多方法主题创建:Knox网关、CDP云、Connect API、管理客户端
- 多途径消息生成:Direct、Knox、CDP、Connect API
- 高级主题管理:使用回退方法创建、列出、描述、删除、配置
- 消息操作:使用元数据和多种传输方式进行生产和消费
- 批量消息生成:高效的批量消息处理
- 补偿管理:高级分区和偏移管理
- 消费者群体管理:完整的消费者群体生命周期管理
📊 监测和健康检查
- 综合健康监测:所有服务的实时健康状况
- 性能指标:请求计数、响应时间、成功率
- 服务发现:自动检测可用服务
- 健康史:随时间跟踪健康状况
- 个人健康检查:对每项服务进行细致的健康监测
- 指标收集:详细的性能和使用指标
🔧 诺克斯网关集成
- 拓扑管理:创建和管理Knox拓扑
- 服务配置:通过Knox配置Kafka服务
- 许可证管理:JWT令牌处理和验证
- 服务发现:自动服务端点发现
- 健康监测:诺克斯特定健康检查
- 管理员API:完全Knox Admin API集成
☁️ CDP云支持
- CDP代理API:与CDP代理端点集成
- CDP令牌身份验证:支持CDP特定令牌
- 服务健康:CDP云服务运行状况监控
- API发现:自动发现CDP API端点
- 后备支援:CDP和其他方法之间的巧妙回退
🛠️ 开发者体验
- 完整模型上下文协议 符合40多种工具
- 元数据:全面的错误报告和状态信息
- 灵活配置:YAML和环境变量配置
- 综合录井:用于调试和监控的结构化日志记录
- 多种运输方式:支持多种Kafka访问方式
- 故障弱化:不同方法之间的自动回退
局限性
已知问题
- 管理客户机:
kafka-python管理客户端在某些环境中可能会失败(NodeNotReadyError) - 制片人超时:在某些Cloudera环境中,消息生成可能会超时
- 模拟源连接器:Cloudera MockSourceConnector可能无法可靠地生成消息
- 诺克斯门户:需要正确配置才能实现完整功能
权变措施
为了确保操作的可靠性,请使用以下替代方案:
# List topics
docker exec -it kafka /opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --list
# Create topics
docker exec -it kafka /opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --create --topic my-topic
# Produce messages
docker exec -it kafka /opt/kafka/bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic my-topic
# Consume messages
docker exec -it kafka /opt/kafka/bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic my-topic --from-beginning快速开始
起床跑步
- 启动Docker环境
git clone https://github.com/ibrooks/cdf-kafka-mcp-server.git
cd cdf-kafka-mcp-server
docker-compose up -d- 验证服务
# Check Kafka
docker exec kafka /opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --list
# Check SMM UI
curl -f http://localhost:9991/ || echo "SMM UI not ready yet"- 添加测试数据
echo '{"message": "Hello from Quick Start!", "timestamp": "2024-10-23T20:00:00Z"}' | \
docker exec -i kafka /opt/kafka/bin/kafka-console-producer.sh \
--bootstrap-server localhost:9092 \
--topic cursortest- 在SMM UI中查看
- 打开http://localhost:9991/
- 登录
admin/admin123 - 导航到最热门的主题
- 测试MCP服务器(可选)
export KAFKA_BOOTSTRAP_SERVERS=localhost:9092
uv run python -m cdf_kafka_mcp_server安装
先决条件
- Python 3.10或更高版本
- Apache Kafka集群
- Apache Knox网关(可选,用于企业身份验证)
从源代码安装
git clone https://github.com/ibrooks/cdf-kafka-mcp-server.git
cd cdf-kafka-mcp-server
uv pip install -e .Cloudera代理工作室
与Cloudera Agent Studio一起使用:
{
"mcpServers": {
"cdf-kafka-mcp-server": {
"command": "uvx",
"args": ["--from", "git+https://github.com/ibrooks/cdf-kafka-mcp-server@main", "run-server"],
"env": {
"KAFKA_BOOTSTRAP_SERVERS": "kafka-broker1:9092,kafka-broker2:9092",
"KNOX_TOKEN": "",
"KNOX_GATEWAY": "https://knox-gateway.yourshere.cloudera.site:8443"
}
}
}
}配置
云部署配置
该项目包括主要云Kafka服务的预配置模板:
AWS MSK(Kafka的托管流媒体)
# Use AWS MSK configuration
cp config/kafka_config_aws_msk.yaml config/kafka_config.yaml
# Set AWS credentials
export AWS_ACCESS_KEY_ID="your-access-key"
export AWS_SECRET_ACCESS_KEY="your-secret-key"
export AWS_REGION="us-east-1"汇流云
# Use Confluent Cloud configuration
cp config/kafka_config_confluent_cloud.yaml config/kafka_config.yaml
# Set Confluent credentials
export CONFLUENT_API_KEY="your-api-key"
export CONFLUENT_API_SECRET="your-api-secret"Kafka的Azure事件中心
# Use Azure Event Hubs configuration
cp config/kafka_config_azure_eventhub.yaml config/kafka_config.yaml
# Set Azure connection string
export AZURE_EVENTHUB_CONNECTION_STRING="your-connection-string"通用云配置
# Use generic cloud configuration
cp config/kafka_config_cloud.yaml config/kafka_config.yaml
# Set your cloud provider credentials
export KAFKA_SASL_USERNAME="your-username"
export KAFKA_SASL_PASSWORD="your-password"快速设置脚本
为了简化云部署设置,请使用提供的脚本:
# Make the script executable
chmod +x setup_cloud.sh
# Setup for your cloud provider
./setup_cloud.sh aws-msk
./setup_cloud.sh confluent-cloud
./setup_cloud.sh azure-eventhub
./setup_cloud.sh generic环境变量模板
对于云部署,请使用环境变量模板:
# Copy and customize the template
cp config/env_cloud_template.txt .env
# Edit .env with your actual values云部署注意事项
在部署到云环境时,请考虑以下重要因素:
安全
- 使用IAM角色:对于AWS MSK,更喜欢IAM角色而不是访问密钥
- 旋转凭据:定期轮换API密钥和令牌
- 网络安全:适当使用VPC和安全组
- 加密:确保为所有连接启用TLS/SSL
演出
- 连接池:云服务可能有连接限制
- 重试逻辑:实现重试的指数回退
- 监控:使用云提供商监控工具
- 扩展:考虑根据负载自动缩放
成本优化
- 资源规模:适当调整Kafka集群的大小
- 数据保留:设置适当的保留策略
- 压缩:启用压缩以降低带宽成本
- 监控:监控使用情况以避免意外收费
可靠性
- 多AZ:跨多个可用区部署
- 备份:实施适当的备份策略
- 灾难恢复:灾难恢复方案计划
- 健康检查实施全面的健康监测
配置文件
在以下位置创建配置文件 config/kafka_config.yaml:
kafka:
bootstrap_servers: "localhost:9092"
client_id: "cdf-kafka-mcp-server"
security_protocol: "PLAINTEXT"
timeout: 30
knox:
gateway: "https://knox-gateway.example.com:8443"
token: "your-knox-token-here"
verify_ssl: true
service: "kafka"环境变量
您还可以使用环境变量配置服务器:
export KAFKA_BOOTSTRAP_SERVERS="localhost:9092"
export KNOX_GATEWAY="https://knox-gateway.example.com:8443"
export KNOX_TOKEN="your-knox-token-here"
export KNOX_VERIFY_SSL="true"Claude桌面集成
添加到您的Claude Desktop配置(claude_desktop_config.json):
{
"mcpServers": {
"cdf-kafka-mcp-server": {
"command": "cdf-kafka-mcp-server",
"args": ["--config", "./config/kafka_config.yaml"],
"env": {
"KNOX_GATEWAY": "https://knox-gateway.example.com:8443",
"KNOX_TOKEN": "your-knox-token-here"
}
}
}
}用法
基本设置
直接Kafka连接:
export KAFKA_BOOTSTRAP_SERVERS="localhost:9092"
uv run python -m cdf_kafka_mcp_server诺克斯网关(制作):
export KNOX_GATEWAY="https://your-knox-gateway:8444"
export KNOX_TOKEN="your-bearer-token-here"
export KNOX_SERVICE="kafka"
uv run python -m cdf_kafka_mcp_server配置文件:
uv run python -m cdf_kafka_mcp_server --config config/kafka_config.yaml身份验证方法
诺克斯熊代币(推荐):
export KNOX_GATEWAY="https://knox-gateway.company.com:8444"
export KNOX_TOKEN="your-bearer-token-here"
export KNOX_SERVICE="kafka"
uv run python -m cdf_kafka_mcp_server直接Kafka连接:
export KAFKA_BOOTSTRAP_SERVERS="kafka1:9092,kafka2:9092"
export KAFKA_SECURITY_PROTOCOL="SASL_SSL"
export KAFKA_SASL_MECHANISM="PLAIN"
export KAFKA_SASL_USERNAME="your-username"
export KAFKA_SASL_PASSWORD="your-password"
uv run python -m cdf_kafka_mcp_server配置
环境变量:
KAFKA_BOOTSTRAP_SERVERS-Kafka代理地址(必填)KNOX_GATEWAY-Knox网关URL(用于Knox身份验证)KNOX_TOKEN-承载令牌(用于Knox身份验证)MCP_LOG_LEVEL-日志级别(默认值:INFO)
YAML配置:
kafka:
bootstrap_servers: "localhost:9092"
security_protocol: "PLAINTEXT"
knox:
gateway: "https://knox-gateway.company.com:8444"
token: "your-bearer-token-here"
service: "kafka"可用的MCP工具
服务器提供 40+MCP工具 对于全面的Kafka操作:
🗂️ 主题管理(7个工具):
list_topics-列出所有Kafka主题create_topic- 增强:使用多种方法创建主题(Knox、CDP、Connect、Admin)describe_topic-获取详细的主题信息delete_topic-删除Kafka主题topic_exists-检查主题是否存在get_topic_partitions-获取主题分区计数update_topic_config-更新主题配置
📨 消息操作(6个工具):
produce_message- 增强:使用多种方法(Direct、Knox、CDP、Connect)生成消息consume_messages-消费来自某个主题的消息get_topic_offsets-获取主题分区偏移量get_consumer_groups-列出消费者群体get_consumer_group_details-获取消费者群体详细信息reset_consumer_group_offsets-重置消费者群体偏移
🔌 Kafka连接管理(15个工具):
list_connectors-列出所有连接器create_connector-创建新连接器get_connector_status-获取连接器状态delete_connector-删除连接器pause_connector-暂停连接器resume_connector-恢复连接器list_connector_plugins-列出可用插件validate_connector_config-验证连接器配置
🔧 Knox网关集成(7个工具):
test_knox_connection-测试Knox网关连接get_knox_metadata-获取Knox网关元数据get_knox_gateway_info- 新:获取诺克斯网关信息和状态list_knox_topologies- 新:列出所有Knox拓扑get_knox_topology- 新:获取特定的Knox拓扑配置create_knox_topology- 新:为Kafka服务创建新的Knox拓扑get_knox_service_health- 新:获取诺克斯服务的健康状况get_knox_service_urls- 新:通过Knox网关获取服务URL
☁️ CDP云集成(4个工具):
test_cdp_connection- 新:测试与CDP Cloud的连接get_cdp_apis- 新:获取有关可用CDP API的信息get_cdp_service_health- 新:获取CDP服务的健康状况validate_cdp_token- 新:验证CDP令牌
📊 监测和健康检查(5个工具):
get_health_status- 新:获取所有服务的全面健康状况get_health_summary- 新:获取健康状况摘要get_health_history- 新:获取健康检查历史记录get_service_metrics- 新:获取服务性能指标run_health_check- 新:运行特定的健康检查
🔗 系统信息(3个工具):
get_broker_info-获取Kafka代理信息get_cluster_metadata-获取集群元数据test_connection-测试Kafka连接
使用示例
增强的主题操作:
{"tool": "list_topics", "arguments": {}}
{"tool": "create_topic", "arguments": {"name": "user-events", "partitions": 3, "method": "auto"}}
{"tool": "describe_topic", "arguments": {"name": "user-events"}}增强的消息操作:
{"tool": "produce_message", "arguments": {"topic": "user-events", "value": "Hello Kafka!", "method": "auto"}}
{"tool": "consume_messages", "arguments": {"topic": "user-events", "max_count": 10}}诺克斯门户运营:
{"tool": "get_knox_gateway_info", "arguments": {}}
{"tool": "create_knox_topology", "arguments": {"topology_name": "kafka-topology", "kafka_brokers": ["broker1:9092"]}}
{"tool": "get_knox_service_health", "arguments": {"topology": "default"}}CDP云运营:
{"tool": "test_cdp_connection", "arguments": {}}
{"tool": "get_cdp_apis", "arguments": {}}
{"tool": "validate_cdp_token", "arguments": {"token": "your-cdp-token"}}监测和健康检查:
{"tool": "get_health_status", "arguments": {}}
{"tool": "get_health_summary", "arguments": {}}
{"tool": "run_health_check", "arguments": {"check_name": "kafka"}}
{"tool": "get_service_metrics", "arguments": {}}Kafka连接:
{"tool": "list_connectors", "arguments": {}}
{"tool": "create_connector", "arguments": {"name": "my-connector", "config": {...}}}Claude桌面集成
{
"mcpServers": {
"cdf-kafka-mcp-server": {
"command": "uv",
"args": ["run", "python", "-m", "cdf_kafka_mcp_server"],
"env": {
"KNOX_GATEWAY": "https://your-knox-gateway:8444",
"KNOX_TOKEN": "your-bearer-token-here"
}
}
}
}视觉示例
通过SMM UI创建主题
以下屏幕截图显示了如何使用Streams Messaging Manager(SMM)web界面创建新的Kafka主题:
*此示例演示了SMM web界面中的主题创建过程,显示了主题名称、分区、复制因子和配置选项。*
向主题添加数据
以下屏幕截图显示了如何使用SMM web界面向Kafka主题添加数据:
*此示例演示了SMM web界面中的消息生成过程,展示了如何使用键值对和标头向Kafka主题发送消息。*
通过SMM UI列出主题
以下屏幕截图显示了如何使用Streams消息传递管理器(SMM)web界面查看和管理Kafka主题:
*此示例演示了SMM中的主题列表和管理界面,显示了所有可用的主题及其配置、分区计数和状态信息。*
使用MCP工具与SMM UI
- SMM-UI:通过web界面进行手动交互式操作
- MCP工具:通过人工智能模型实现自动化、程序化操作
- Kafka命令行界面:用于脚本编写和自动化的命令行操作
故障排除
常见问题
MCP服务器无法启动:
# Check environment variables
echo $KAFKA_BOOTSTRAP_SERVERS
echo $KNOX_GATEWAY
echo $KNOX_TOKEN
# Test with debug logging
export MCP_LOG_LEVEL="DEBUG"
uv run python -m cdf_kafka_mcp_serverKnox身份验证问题:
# Verify Knox Gateway is accessible
curl -k https://your-knox-gateway:8444/gateway/admin/v1/version
# Test with different SSL settings
export KNOX_VERIFY_SSL="false" # For testingMCP工具不工作:
# Test connection first
# Use test_connection MCP tool
# Verify Kafka Connect is running
curl http://localhost:28083/connectors调试模式
启用调试日志记录:
export MCP_LOG_LEVEL="DEBUG"
uv run python -m cdf_kafka_mcp_server健康检查命令
# Check all services
docker-compose ps
# Check Kafka health
docker exec kafka /opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --list
# Check SMM health
curl -f http://localhost:9991/
# Check Kafka Connect health
curl -f http://localhost:28083/connectors紧急应对措施
如果MCP服务器不工作:
- 使用Kafka CLI进行所有操作
- 使用SMM UI进行可视化
- 直接使用Kafka Connect REST API
- 检查服务日志中的特定错误
安全注意事项
- 所有敏感数据(密码、令牌、机密)都会在响应中自动编辑
- Knox令牌会自动缓存和刷新
- 对Knox连接实施TLS
- 配置文件应使用适当的权限进行保护
云测试
快速设置
选择云提供商:
# AWS MSK
./Testing/run_cloud_tests.sh aws-msk
# Confluent Cloud
./Testing/run_cloud_tests.sh confluent-cloud
# Azure Event Hubs
./Testing/run_cloud_tests.sh azure-eventhub
# CDP Cloud (Cloudera Data Platform)
./Testing/run_cloud_tests.sh cdp-cloud
# CDP Cloud MCP Tools Test (Comprehensive)
./Testing/run_cdp_cloud_tests.sh
# Generic SASL_SSL
./Testing/run_cloud_tests.sh generic设置环境变量:
# AWS MSK Example
export KAFKA_BOOTSTRAP_SERVERS="your-msk-cluster.kafka.us-east-1.amazonaws.com:9092"
export KAFKA_SECURITY_PROTOCOL="SASL_SSL"
export KAFKA_SASL_MECHANISM="SCRAM-SHA-512"
export KAFKA_SASL_USERNAME="your-iam-username"
export KAFKA_SASL_PASSWORD="your-iam-password"
# Optional Knox Gateway
export KNOX_GATEWAY="https://your-knox-gateway:8444"
export KNOX_TOKEN="your-bearer-token-here"运行测试:
# Run comprehensive cloud tests
./Testing/run_cloud_tests.sh aws-msk --debug
# Or run individual test script
uv run python3 Testing/test_cloud_connection.py云配置文件
config/kafka_config_aws_msk.yaml-AWS MSK配置config/kafka_config_confluent_cloud.yaml-汇流云配置config/kafka_config_azure_eventhub.yaml-Azure事件中心配置config/kafka_config_cdp_cloud.yaml-CDP云(Cloudera数据平台)配置config/kafka_config_cloud.yaml-通用云配置
文档
- CLOUD_TESTING_SETUP.md -全面的云测试指南
- 测试/CDP_CLOUD_Testing.md -CDP云特定测试指南
- 云部署配置 -配置示例
发展
项目结构
├── src/
│ └── cdf_kafka_mcp_server/ # Main package
│ ├── __init__.py # Package initialization
│ ├── main.py # CLI entry point
│ ├── config.py # Configuration management
│ ├── knox_client.py # Knox authentication
│ ├── kafka_client.py # Kafka client implementation
│ └── mcp_server.py # MCP server implementation
├── config/ # Configuration files
├── examples/ # Usage examples
├── Testing/ # Test files
├── pyproject.toml # Project configuration
└── README.md # Documentation建筑
# Install in development mode
pip install -e .
# Build package
python -m build
# Run tests
python -m pytest
# Format code
black src/ Testing/
isort src/ Testing/
# Lint code
flake8 src/ Testing/
mypy src/贡献
- 复刻仓库
- 创建要素分支
- 进行更改
- 如果适用,添加测试
- 提交拉取请求
许可证
Apache许可证2.0
致谢
这个项目的灵感来自 SSB MCP服务器 实现,为MCP服务器开发和企业集成提供了优秀的模式。
