kafka连接mcp

一个MCP服务器,它公开 Kafka Connect REST API 操作作为工具,让LLM通过自然语言管理连接器、任务和插件。
工具
| 工具 | Kafka连接端点 | 说明 |
|---|---|---|
get_cluster_info | GET / | 集群版本和元数据 |
list_connectors | GET /connectors | 列出所有连接器名称 |
get_connector | GET /connectors/{name} | 连接器信息、配置和任务 |
get_connector_status | GET /connectors/{name}/status | 连接器和任务状态 |
get_connector_config | GET /connectors/{name}/config | 连接器配置 |
create_connector | POST /connectors | 创建新连接器 |
update_connector_config | PUT /connectors/{name}/config | 更换连接器配置 |
delete_connector | DELETE /connectors/{name} | 删除连接器 |
pause_connector | PUT /connectors/{name}/pause | 暂停连接器 |
resume_connector | PUT /connectors/{name}/resume | 恢复暂停的连接器 |
restart_connector | POST /connectors/{name}/restart | 重新启动连接器(可选任务) |
get_task_status | GET /connectors/{name}/tasks/{id}/status | 特定任务的状态 |
restart_task | POST /connectors/{name}/tasks/{id}/restart | 重新启动特定任务 |
list_connector_plugins | GET /connector-plugins | 集群上可用的插件 |
validate_connector_config | PUT /connector-plugins/{name}/config/validate | 根据插件架构验证配置 |
设置
先决条件
- Python 3.12+
- 紫外线
- 正在运行的Kafka Connect集群(或使用附带的Docker Compose)
安装
uv sync添加到克劳德代码
claude mcp add kafka-connect \
-e KAFKA_CONNECT_URL=http://localhost:8083 \
-- uv --directory /path/to/kafka-connect-mcp run kafka-connect-mcp配置
| 环境变量 | 默认值 | 描述 |
|---|---|---|
KAFKA_CONNECT_URL | http://localhost:8083 | Kafka Connect REST API基础URL |
KAFKA_CONNECT_ENABLE_CREATE | false | 允许 create_connector |
KAFKA_CONNECT_ENABLE_UPDATE | false | 允许 update_connector_config |
KAFKA_CONNECT_ENABLE_DELETE | false | 允许 delete_connector |
KAFKA_CONNECT_ENABLE_PAUSE_RESUME | false | 允许 pause_connector 和 resume_connector |
KAFKA_CONNECT_ENABLE_RESTART | false | 允许 restart_connector 和 restart_task |
KAFKA_CONNECT_MUTATION_ALLOWLIST | _(空)_ | 用于变异操作的可选逗号分隔连接器列表 |
安全模式(能力门控)
此服务器是 只读 默认情况下,因为所有变异功能默认为 false. 除非您明确启用特定功能,否则将阻止修改工具。
示例:
所选连接器的仅重新启动模式:
KAFKA_CONNECT_ENABLE_RESTART=true
KAFKA_CONNECT_MUTATION_ALLOWLIST=payments-sink,inventory-source启用所有变异操作(仅用于开发):
KAFKA_CONNECT_ENABLE_CREATE=true
KAFKA_CONNECT_ENABLE_UPDATE=true
KAFKA_CONNECT_ENABLE_DELETE=true
KAFKA_CONNECT_ENABLE_PAUSE_RESUME=true
KAFKA_CONNECT_ENABLE_RESTART=true跑步
stdio(默认值,用于克劳德代码)
KAFKA_CONNECT_URL=http://localhost:8083 uv run kafka-connect-mcpSSE(用于Docker/远程)
uv run kafka-connect-mcp --transport sse --host 0.0.0.0 --port 8000Docker Compose(全栈)
启动Zookeeper、Kafka、Kafka Connect(使用Datagen插件)和MCP服务器:
docker compose up --build -d| 服务 | 端口 | 描述 |
|---|---|---|
| 动物园管理员 | 2181 | zookeeper |
| 卡夫卡 | 9092 | 卡夫卡经纪人 |
| kafka-connect | 8083 | kafka-connect REST API |
| mcp服务器 | 8000 | mcp服务器(SSE传输) |
示例:创建Datagen连接器
堆栈完成后,让Claude创建一个数据生成器连接器,或直接创建:
curl -X POST http://localhost:8083/connectors \
-H "Content-Type: application/json" \
-d '{
"name": "datagen-users",
"config": {
"connector.class": "io.confluent.kafka.connect.datagen.DatagenConnector",
"kafka.topic": "users",
"quickstart": "users",
"key.converter": "org.apache.kafka.connect.storage.StringConverter",
"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter.schemas.enable": "false",
"max.interval": "1000",
"tasks.max": "1"
}
}'测试
uv run pytest tests/ -v测试使用 respx 模拟对Kafka Connect API的HTTP调用——不需要运行集群。
发布
从以下位置自动发布 main:
- 碰撞
project.version在pyproject.toml(例如0.1.0->0.1.1). - 合并到
main. - GitHub Actions创建标签/发布
vX.Y.Z. - 已发布的版本会自动构建并发布到PyPI。
一次性PyPI设置
为项目设置PyPI可信发布 kafka-connect-mcp:
- 业主:
lawrencemq - 存储库:
kafka-connect-mcp - 工作流程:
.github/workflows/release.yml - 环境: _(此工作流不需要)_
项目结构
kafka-connect-mcp/
├── pyproject.toml
├── Dockerfile
├── docker-compose.yml
├── src/kafka_connect_mcp/
│ ├── __init__.py
│ ├── safety.py # Read-only and mutation policy gates
│ └── server.py # MCP tools and entry point
└── tests/
├── conftest.py # Shared fixtures (mock URL + respx router)
├── test_cluster.py
├── test_connectors.py
├── test_safety.py
├── test_tasks.py
└── test_plugins.py许可证
根据Apache许可证2.0版授权。看 许可证.
