AsyncAPI MCP Kafka生成器
自定义AsyncAPI生成器模板,自动生成 模型上下文协议(MCP) Python中的服务器,来自AsyncAPI v3规范。
目标是让LLM(Claude、Cursor等)使用AsyncAPI规范作为唯一的真实来源,与基于Kafka的事件架构进行交互。
______________________________________________________________________
运作原理
AsyncAPI spec (.yaml)
│
▼
AsyncAPI Generator + template/
│
▼
generated-code/
├── mcp_server.py ← FastMCP server with one @mcp.tool per send operation
└── kafka_producer.py ← Confluent Kafka producer with Schema Registry对于每一个 send 在规范中定义的操作中,生成器创建了一个类型化的Python函数,用 @mcp.tool该函数根据从规范中提取的JSON模式验证和序列化有效载荷,然后将消息生成到相应的Kafka主题。
从读取安全配置(代理主机、身份验证方案) servers[] 规范的一部分,并直接嵌入到生成的代码中。
______________________________________________________________________
项目结构
.
├── template/
│ ├── mcp_server.js # Dynamic template (AsyncAPI React SDK) → generates mcp_server.py
│ └── kafka_producer.py # Static file — copied as-is to generated-code/
├── yaml/ # AsyncAPI v3 spec examples
│ ├── streets-lights.yaml # PLAINTEXT + SCRAM-SHA-256 servers
│ ├── streets-lights-oauth.yaml # OAuth2 client_credentials server
│ └── temperature.yaml # IoT temperature sensors, PLAINTEXT
├── generated-code/ # Output of the generator (gitignored — regenerate as needed)
│ ├── mcp_server.py
│ └── kafka_producer.py
├── consumer/ # Standalone Kafka consumer for testing/observability
│ ├── kafka_consumer.py
│ └── run_consumer.py
├── docker/ # Local dev environment
│ ├── docker-compose.yml # Kafka SASL_SSL (SCRAM) + Schema Registry (Basic Auth) + nginx
│ ├── docker-compose-plain.yml # Kafka PLAINTEXT only — for quick local testing
│ ├── docker-compose-oauth.yml # Kafka SASL_SSL (OAUTHBEARER) + Schema Registry (Basic Auth) + Keycloak
│ ├── keycloak-realm.example.json # Keycloak realm config template (copy to keycloak-realm.json)
│ ├── gen-certs.sh # Generates SSL certs for SASL_SSL listeners (run once)
│ ├── kafka-init.sh # Creates SCRAM users at broker startup
│ └── launch.sh # Interactive launcher — choose which stack to bring up
├── api/ # FastAPI service that wraps the generator (REST interface)
│ └── app.py
├── scripts/
│ ├── run_asyncapi_generator.sh # Wrapper around asyncapi generate
│ └── run_mcp_inspector.sh # Launches the MCP Inspector UI
├── package.json # AsyncAPI generator template config
└── pyproject.toml # Python dependencies (uv)______________________________________________________________________
先决条件
- Node.js (用于AsyncAPI生成器)
- Python≥3.13 + 紫外线
- Docker+Docker组合 (适用于本地Kafka堆栈)
安装AsyncAPI生成器和Python依赖项:
npm install
uv sync______________________________________________________________________
生成MCP服务器
使用提供的脚本,传递规范文件和服务器名称(可选):
./scripts/run_asyncapi_generator.sh yaml/streets-lights.yaml plain-connections这 server 参数选择在生成的代码中嵌入哪个代理主机和安全方案。如果省略,则使用规范中的第一个服务器。
符合规范的可用服务器:
| 规格 | 服务器名称 | 安全性 |
|---|---|---|
streets-lights.yaml | plain-connections | 无(PLAINTEXT) |
streets-lights.yaml | scram-connections | SCRAM-SHA-256 |
streets-lights-oauth.yaml | oauth-connections | OAuth2客户端_凭据 |
______________________________________________________________________
当地发展环境
以交互方式启动堆栈:
cd docker/
bash launch.sh有三种堆栈可供选择:
| 堆栈 | 编写文件 | 端口 | 服务 |
|---|---|---|---|
| 明文 | docker-compose-plain.yml | 9092 | Kafka+模式注册表(打开)+Kafka-ui |
| 紧急停堆 | docker-compose.yml | 9095 | Kafka(SASL_SSL+SCRAM-SHA-256)+模式注册表(基本认证)+Kafka-ui+nginx |
| OAuth2 | docker-compose-oauth.yml | 9095 | Kafka(SASL_SSL+OAUTHBEARER)+模式注册表(基本认证)+Kafka-ui+Keycloak |
对于SCRAM和OAuth2堆栈,首先生成SSL证书(只需要一次):
cd docker/
bash gen-certs.shOAuth2/钥匙斗篷设置
OAuth2堆栈需要 docker/keycloak-realm.json 文件(gitignored)。复制示例并设置您的客户端密码:
cp docker/keycloak-realm.example.json docker/keycloak-realm.json
# Edit keycloak-realm.json and replace "CHANGE_ME" with a real secret王国(masorange)和客户(mcp-app)在首次启动时自动导入。不需要手动Keycloak配置。
______________________________________________________________________
运行生成的MCP服务器
配置 generated-code/.env 为您的安全方案设置适当的变量:
SCHEMA_REGISTRY_URL=http://localhost:8081
# PLAINTEXT — no extra config needed
# SCRAM-SHA-256
KAFKA_USERNAME=testuser
KAFKA_PASSWORD=testpassword
KAFKA_SSL_CA_LOCATION=../docker/certs/ca.crt
# OAuth2 client_credentials
OAUTH_CLIENT_ID=mcp-app
OAUTH_CLIENT_SECRET=your-secret
OAUTH_TOKEN_URL=http://localhost:9090/realms/masorange/protocol/openid-connect/token
# Schema Registry Basic Auth (SCRAM and OAuth2 stacks only — not needed for PLAINTEXT)
SCHEMA_REGISTRY_USERNAME=admin
SCHEMA_REGISTRY_PASSWORD=testpassword运行服务器:
cd generated-code/
uv run fastmcp dev mcp_server.py # development mode (MCP Inspector)
uv run mcp_server.py # production mode______________________________________________________________________
通过安全方案快速启动
明文
不需要证书或凭据。
# 1. Start the stack
cd docker && bash launch.sh # option 1
# 2. Generate the MCP server
./scripts/run_asyncapi_generator.sh yaml/streets-lights.yaml plain-connections
# 3. Run
cd generated-code && uv run mcp_server.py______________________________________________________________________
SCRAM-SHA-256
# 1. Generate SSL certificates (broker only — run once)
cd docker && bash gen-certs.sh
# 2. Start the stack
bash launch.sh # option 2
# 3. Generate the MCP server
./scripts/run_asyncapi_generator.sh yaml/streets-lights.yaml scram-connections
# 4. Configure generated-code/.env
KAFKA_USERNAME=testuser
KAFKA_PASSWORD=testpassword
KAFKA_SSL_CA_LOCATION=../docker/certs/ca.crt
SCHEMA_REGISTRY_URL=http://localhost:8081
SCHEMA_REGISTRY_USERNAME=admin
SCHEMA_REGISTRY_PASSWORD=testpassword
# 5. Run
cd generated-code && uv run mcp_server.py______________________________________________________________________
OAuth2(客户端证书)
# 1. Set up the Keycloak realm
cp docker/keycloak-realm.example.json docker/keycloak-realm.json
# Edit keycloak-realm.json and replace "CHANGE_ME" with a real client secret
# 2. Generate SSL certificates (broker only — run once)
cd docker && bash gen-certs.sh
# 3. Start the stack
bash launch.sh # option 3
# 4. Generate the MCP server
./scripts/run_asyncapi_generator.sh yaml/streets-lights-oauth.yaml oauth-connections
# 5. Configure generated-code/.env
KAFKA_SSL_CA_LOCATION=../docker/certs/ca.crt
SCHEMA_REGISTRY_URL=http://localhost:8081
SCHEMA_REGISTRY_USERNAME=admin
SCHEMA_REGISTRY_PASSWORD=testpassword
OAUTH_CLIENT_ID=mcp-app
OAUTH_CLIENT_SECRET=
OAUTH_TOKEN_URL=http://localhost:9090/realms/masorange/protocol/openid-connect/token
# 6. Run
cd generated-code && uv run mcp_server.py______________________________________________________________________
安全支持
从规范中读取安全配置 servers[].security 字段并嵌入生成的 mcp_server.py.
| AsyncAPI方案 | Kafka协议 | 身份验证 |
|---|---|---|
| _(无)_ | PLAINTEXT | 无 |
scramSha256 | SASL_SSL | SCRAM-SHA-256(用户名+密码) |
oauth2 / clientCredentials | SASL_SSL | OAUTHBEARER——通过客户端密钥获取M2M令牌 |
OAuth2客户端_凭证流
MCP服务器使用 client_id + client_secret无需人工干预——只要Kafka客户端需要,令牌就会在后台自动获取。Kafka代理根据Keycloak的JWKS端点验证JWT签名,并检查 iss 和 aud 声称。
这非常适合MCP服务器:生产者是一个具有自己的应用程序身份的自动化服务,而不是人类用户。
OAuth2 client_credentials flow
______________________________________________________________________
使用的AsyncAPI规范功能
路径参数
通道地址参数(例如。 {streetlightId})被提取并作为函数参数公开。默认情况下,第一个路径参数用作Kafka分区键。
自定义分区密钥(x-kafka-key)
用以下命令覆盖分区键字段 x-kafka-key 操作扩展:
operations:
sendReading:
action: send
x-kafka-key: sensorId
channel:
$ref: '#/channels/readingsChannel'可选字段
有效载荷属性未在以下列出 required 生成为 Optional[T] = None 并在发送前过滤掉。
文档字符串
生成的工具函数包括从以下内容构建的文档字符串 summary, description和财产 description 规范中的字段。
______________________________________________________________________
消耗事件(可观察性)
提供了一个简单的消费者进行测试。通过配置 consumer/.env:
KAFKA_BOOTSTRAP_SERVERS=localhost:9095
TOPICS=smartylighting.streetlights.1.0.action.home.turn.on
CONSUMER_GROUP_ID=streetlights-logger
KAFKA_USERNAME=testuser
KAFKA_PASSWORD=testpassword
KAFKA_SSL_CA_LOCATION=../docker/certs/ca.crtcd consumer/
uv run run_consumer.py______________________________________________________________________
示例规格
| 文件 | 服务器 | 描述 |
|---|---|---|
yaml/streets-lights.yaml | plain-connections (9092), scram-connections (9095) | 路灯API-路径参数、PLAINTEXT和SCRAM |
yaml/streets-lights-oauth.yaml | oauth-connections (9095) | 路灯API-OAuth2客户端_凭证 |
yaml/temperature.yaml | _(普通)_ | 物联网温度传感器——必填字段,无身份验证 |
