mcp-kafka-observer
by Rushi264
README.md
# mcp-kafka-observer
An MCP (Model Context Protocol) server that gives AI agents real-time observability into Apache Kafka clusters. Monitor broker health, track consumer lag, and diagnose issues — all through natural language.
## Why?
Kafka monitoring typically requires juggling multiple dashboards. This MCP server lets any AI assistant (Claude, ChatGPT, Cursor, VS Code Copilot) query your Kafka cluster directly:
- "Is my Kafka cluster healthy?"
- "What's the consumer lag for payment-processor group?"
- "Why is lag spiking on the orders topic?"
## Tools
| Tool | Description |
|------|-------------|
| `get_broker_health` | Cluster state: brokers, controller, under-replicated partitions |
| `list_topics` | All topics with partition counts and replication factors |
| `describe_topic` | Detailed config and partition assignments for a topic |
| `get_consumer_lag` | Per-partition lag for a consumer group |
| `diagnose_lag_spike` | Automated root-cause analysis for lag issues |
| `get_cache_stats` | Cache hit/miss statistics for observability |
## Resources
| Resource URI | Description |
|-------------|-------------|
| `kafka://cluster/overview` | High-level cluster summary |
## Prompts
| Prompt | Description |
|--------|-------------|
| `investigate_lag` | Step-by-step workflow for diagnosing consumer lag |
| `capacity_review` | Template for cluster capacity planning |
## Quick Start
### Prerequisites
- Python 3.12+
- Docker (for local Kafka)
- [uv](https://docs.astral.sh/uv/) package manager
### Setup
```bash
git clone https://github.com/Rushi264/mcp-kafka-observer.git
cd mcp-kafka-observer
# Install dependencies
uv sync
# Start local Kafka
docker compose up -d
# Run tests
uv run pytest -v
```
### Claude Desktop Integration
Add to your `claude_desktop_config.json`:
```json
{
"mcpServers": {
"kafka-observer": {
"command": "uv",
"args": [
"--directory", "/path/to/mcp-kafka-observer",
"run", "python", "-m", "mcp_kafka_observer.server"
],
"env": {
"KAFKA_BOOTSTRAP_SERVERS": "localhost:9092"
}
}
}
}
```
## Architecture
```
MCP Client (Claude / Cursor / VS Code Copilot)
│
│ MCP Protocol (stdio)
▼
mcp-kafka-observer
├── Tools (get_broker_health, get_consumer_lag, ...)
├── Resources (kafka://cluster/overview)
├── Prompts (investigate_lag, capacity_review)
├── TTL Cache (prevents thundering herd on admin API)
└── Analyzer (automated lag diagnosis)
│
│ confluent-kafka AdminClient
▼
Kafka Cluster
```
## Tech Stack
- **Python 3.12** with async/await
- **MCP SDK** (FastMCP) — official Anthropic SDK
- **confluent-kafka** — production-grade Kafka client (librdkafka)
- **Pydantic** — structured output validation
- **Docker Compose** — local Kafka for development
## Testing
```bash
# Unit tests (no Kafka needed)
uv run pytest tests/test_server.py -v
# Integration tests (needs Docker Kafka running)
docker compose up -d
uv run pytest tests/test_kafka_client.py -v
# All tests
uv run pytest -v
# Linter
uv run ruff check src/ tests/
```
## Configuration
Set via environment variables or `.env` file:
| Variable | Default | Description |
|----------|---------|-------------|
| `KAFKA_BOOTSTRAP_SERVERS` | `localhost:9092` | Kafka broker addresses |
| `KAFKA_SASL_MECHANISM` | — | SASL auth mechanism (PLAIN, SCRAM-SHA-256) |
| `KAFKA_SASL_USERNAME` | — | SASL username |
| `KAFKA_SASL_PASSWORD` | — | SASL password |
| `KAFKA_SECURITY_PROTOCOL` | — | Security protocol (SASL_SSL, SASL_PLAINTEXT) |
## License
MITThis server cannot be deployed
Maintenance
ActivityInactive
ResponsivenessNo issues