Kafka AI消息服务
# 翻译 使AI模型能够通过标准化接口从Apache Kafka主题发布和消费消息,从而轻松地将Kafka消息与大语言模型和代理应用程序集成。
MCP 服务配置
复制以下 JSON 到 OPClaw 或其他 MCP 客户端的配置文件中即可使用
{
"mcpServers": {
"kafka": {
"args": [
"\u003cPATH TO PROJECTS\u003e/main.py"
],
"command": "python"
}
}
}
该服务需要配置环境变量:DEFAULT_GROUP_ID_FOR_CONSUMER、IS_TOPIC_READ_FROM_BEGINNING、KAFKA_BOOTSTRAP_SERVERS、TOOL_CONSUME_DESCRIPTION、TOOL_PUBLISH_DESCRIPTION、TOPIC_NAME
服务介绍
Kafka MCP 服务器
一个与 Apache Kafka 集成的消息上下文协议 (MCP) 服务器,为 LLM 和代理应用程序提供发布和消费功能。
概述
本项目实现了一个服务器,允许 AI 模型通过标准化接口与 Kafka 主题进行交互。它支持:
- 向 Kafka 主题发布消息
- 从 Kafka 主题消费消息
前提条件
- Python 3.8+
- Apache Kafka 实例
- Python 依赖项(见安装部分)
安装
-
克隆仓库:
git clone <repository-url> cd <repository-directory> -
创建并激活虚拟环境:
python -m venv venv source venv/bin/activate # 在 Windows 上使用: venv\Scripts\activate -
安装所需的依赖项:
pip install -r requirements.txt如果没有
requirements.txt文件,则安装以下包:pip install aiokafka python-dotenv pydantic-settings mcp-server
配置
在项目根目录下创建一个 .env 文件,并添加以下变量:
# Kafka Configuration
KAFKA_BOOTSTRAP_SERVERS=localhost:9092
TOPIC_NAME=your-topic-name
IS_TOPIC_READ_FROM_BEGINNING=False
DEFAULT_GROUP_ID_FOR_CONSUMER=kafka-mcp-group
# Optional: Custom Tool Descriptions
# TOOL_PUBLISH_DESCRIPTION="Custom description for the publish tool"
# TOOL_CONSUME_DESCRIPTION="Custom description for the consume tool"
使用
运行服务器
您可以使用提供的 main.py 脚本来运行服务器:
python main.py --transport stdio
可用的传输选项:
stdio: 标准输入/输出(默认)sse: 服务器发送事件
与 Claude Desktop 集成
要将此 Kafka MCP 服务器与 Claude Desktop 一起使用,请将以下配置添加到您的 Claude Desktop 配置文件中:
{
"mcpServers": {
"kafka": {
"command": "python",
"args": [
"<PATH TO PROJECTS>/main.py"
]
}
}
}
将 <PATH TO PROJECTS> 替换为项目目录的绝对路径。
项目结构
main.py: 应用程序入口点kafka.py: Kafka 连接器实现server.py: 带有 Kafka 交互工具的 MCP 服务器实现settings.py: 使用 Pydantic 的配置管理
可用工具
kafka-publish
向配置的 Kafka 主题发布信息。
kafka-consume
从配置的 Kafka 主题消费信息。
- 注意:一旦从主题中读取消息,使用相同的 groupid 就不能再读取该消息