Kafka MCP Server

Un servidor MCP para Apache Kafka, que permite a los LLMs realizar operaciones de Kafka como producir y consumir mensajes.

Documentación

Kafka MCP Server

Un servidor de Model Context Protocol (MCP) para Apache Kafka implementado en Go, que utiliza franz-go y mcp-go.

Este servidor proporciona una implementación para interactuar con Kafka a través del protocolo MCP, permitiendo que los modelos LLM realicen operaciones comunes de Kafka mediante una interfaz estandarizada.

Go Report Card GitHub Workflow Status Go Version Trivy Scan SLSA 3 Go Reference Docker Image GitHub Release License: MIT

Descripción general

El Kafka MCP Server cierra la brecha entre los modelos LLM y Apache Kafka, permitiéndoles:

  • Producir y consumir mensajes de topics
  • Listar, describir y gestionar topics
  • Monitorear y gestionar grupos de consumidores
  • Evaluar la salud y configuración del clúster
  • Ejecutar operaciones estándar de Kafka

Todo a través del protocolo estandarizado Model Context Protocol (MCP).

Arquitectura

graph TB
    subgraph "MCP Client (AI Applications)"
        A[Claude Desktop]
        B[Cursor]
        C[Windsurf]
        D[ChatWise]
    end
    
    subgraph "Kafka MCP Server"
        E[MCP Protocol Handler]
        F[Tools Registry]
        G[Resources Registry]
        H[Prompts Registry]
        I[Kafka Client Wrapper]
    end
    
    subgraph "Apache Kafka Cluster"
        J[Broker 1]
        K[Broker 2]
        L[Broker 3]
        M[Topics & Partitions]
        N[Consumer Groups]
    end
    
    A --> E
    B --> E
    C --> E
    D --> E
    
    E --> F
    E --> G
    E --> H
    
    F --> I
    G --> I
    H --> I
    
    I --> J
    I --> K
    I --> L
    
    J --> M
    K --> M
    L --> M
    
    J --> N
    K --> N
    L --> N
    
    classDef client fill:#e1f5fe
    classDef mcp fill:#f3e5f5
    classDef kafka fill:#fff3e0
    
    class A,B,C,D client
    class E,F,G,H,I mcp
    class J,K,L,M,N kafka

Cómo funciona:

  1. Clientes MCP (aplicaciones de IA) se conectan al Kafka MCP Server mediante transporte stdio o HTTP
  2. Servidor MCP expone tres tipos de capacidades:
    • Herramientas - Operaciones directas de Kafka (producir/consumir mensajes, describir topics, etc.)
    • Recursos - Informes de salud del clúster y diagnósticos
    • Prompts - Flujos de trabajo preconfigurados para operaciones comunes
  3. Envoltorio del Cliente Kafka maneja toda la comunicación con Kafka utilizando la librería franz-go
  4. Clúster Apache Kafka procesa el streaming y almacenamiento real de mensajes

Modos de transporte:

  • STDIO: Modo predeterminado, ideal para clientes MCP locales (Claude Desktop, Cursor, etc.)
  • HTTP: Permite acceso remoto con autenticación OAuth 2.1 opcional

Tools

Prompts & Resources

Características principales

  • Integración con Kafka: Implementación de operaciones comunes de Kafka mediante MCP
  • Seguridad:
    • Soporte para autenticación SASL (PLAIN, SCRAM-SHA-256, SCRAM-SHA-512) y TLS
    • Autenticación OAuth 2.1 para transporte HTTP (modos Nativo y Proxy)
    • Soporte para proveedores Okta, Google, Azure AD y HMAC
  • Transporte flexible: STDIO para clientes locales, HTTP para acceso remoto
  • Manejo de errores: Manejo de errores con retroalimentación significativa
  • Opciones de configuración: Personalizable para diferentes entornos
  • Prompts preconfigurados: Conjunto de prompts para operaciones comunes de Kafka
  • Compatibilidad: Funciona con modelos LLM compatibles con MCP

Primeros pasos

Requisitos previos

  • Go 1.24 o posterior
  • Docker (para ejecutar pruebas de integración)
  • Acceso a un clúster de Kafka

Instalación

Homebrew (macOS y Linux)

La forma más sencilla de instalar kafka-mcp-server es usando Homebrew:

# Add the tap repository
brew tap tuannvm/mcp

# Install kafka-mcp-server
brew install kafka-mcp-server

Para actualizar a la versión más reciente:

brew update && brew upgrade kafka-mcp-server

Desde el código fuente

# Clone the repository
git clone https://github.com/tuannvm/kafka-mcp-server.git
cd kafka-mcp-server

# Build the server
go build -o kafka-mcp-server ./cmd

Integración con clientes MCP

Este servidor MCP puede integrarse con varias aplicaciones de IA. A continuación se presentan instrucciones específicas por plataforma:

Cursor

Edita ~/.cursor/mcp.json y añade la configuración de kafka-mcp-server:

{
  "mcpServers": {
    "kafka": {
      "command": "kafka-mcp-server",
      "args": [],
      "env": {
        "KAFKA_BROKERS": "localhost:9092",
        "KAFKA_CLIENT_ID": "kafka-mcp-server",
        "MCP_TRANSPORT": "stdio"
      }
    }
  }
}

Claude Desktop

Edita tu archivo de configuración de Claude y añade el servidor:

  • macOS: ~/Library/Application Support/Claude/claude_desktop_config.json
  • Windows: %APPDATA%\Claude\claude_desktop_config.json
{
  "mcpServers": {
    "kafka": {
      "command": "kafka-mcp-server",
      "args": [],
      "env": {
        "KAFKA_BROKERS": "localhost:9092",
        "KAFKA_CLIENT_ID": "kafka-mcp-server",
        "MCP_TRANSPORT": "stdio"
      }
    }
  }
}

Reinicia Claude Desktop para aplicar los cambios.

Claude Code

Para usar con Claude Code, añade el servidor usando el comando de configuración MCP integrado:

# Add kafka-mcp-server with environment variables
claude mcp add kafka \
  --env KAFKA_BROKERS=localhost:9092 \
  --env KAFKA_CLIENT_ID=kafka-mcp-server \
  --env MCP_TRANSPORT=stdio \
  --env KAFKA_SASL_MECHANISM= \
  --env KAFKA_SASL_USER= \
  --env KAFKA_SASL_PASSWORD= \
  --env KAFKA_TLS_ENABLE=false \
  -- kafka-mcp-server

Otros comandos útiles:

# List configured MCP servers
claude mcp list

# Remove server
claude mcp remove kafka

# Test server connection
claude mcp get kafka

ChatWise

  1. Abre ChatWise → Configuración → Herramientas → "+" → "Command Line MCP"
  2. Configura:
    • ID: kafka
    • Command: kafka-mcp-server
    • Args: (déjalo vacío)
    • Env: Añade variables de entorno:
      KAFKA_BROKERS=localhost:9092
      KAFKA_CLIENT_ID=kafka-mcp-server
      MCP_TRANSPORT=stdio
      

Simplifica la configuración con mcpenetes

Gestionar las configuraciones del servidor MCP en múltiples clientes puede resultar complicado. mcpenetes es una herramienta dedicada que facilita significativamente este proceso:

# Install mcpenetes
go install github.com/tuannvm/mcpenetes@latest

Características principales

  • Búsqueda interactiva: Encuentra y selecciona configuraciones del servidor Kafka MCP con un simple comando
  • Aplicar en todas partes: Sincroniza automáticamente las configuraciones en todos tus clientes MCP
  • Copia de seguridad de configuración: Realiza copias de seguridad seguras de las configuraciones existentes antes de hacer cambios
  • Restaurar: Revierte fácilmente a configuraciones anteriores si es necesario

Inicio rápido con mcpenetes

# Search for available MCP servers including kafka-mcp-server
mcpenetes search 

# Apply kafka-mcp-server configuration to all your clients at once
mcpenetes apply

# Load a configuration from your clipboard
mcpenetes load

Con mcpenetes, puedes mantener múltiples configuraciones de Kafka (desarrollo, producción, etc.) y cambiar entre ellas al instante en todos tus clientes (Cursor, Claude Desktop, Windsurf, ChatWise) sin editar manualmente los archivos de configuración de cada cliente.

Herramientas MCP

El servidor expone las siguientes herramientas para la interacción con Kafka. Para documentación detallada con ejemplos y respuestas de muestra, consulta docs/tools.md.

  • produce_message: Produce mensajes en topics de Kafka
  • consume_messages: Consume mensajes de topics de Kafka en operaciones por lotes
  • list_brokers: Lista todas las direcciones de brokers Kafka configuradas
  • describe_topic: Proporciona metadatos completos para topics específicos
  • list_consumer_groups: Enumera todos los grupos de consumidores en el clúster
  • describe_consumer_group: Proporciona información detallada del grupo de consumidores, incluyendo métricas de lag
  • describe_configs: Recupera los ajustes de configuración de los recursos de Kafka
  • cluster_overview: Proporciona resúmenes completos de salud del clúster
  • list_topics: Lista todos los topics con metadatos, incluyendo información de particiones y replicación

Recursos MCP

El servidor proporciona los siguientes recursos que se pueden acceder a través del protocolo MCP. Para documentación detallada con respuestas de ejemplo, consulta docs/resources.md.

  • kafka-mcp://overview: Resumen completo de salud del clúster
  • kafka-mcp://health-check: Evaluación detallada de salud con información procesable
  • kafka-mcp://under-replicated-partitions: Análisis de particiones con problemas de replicación
  • kafka-mcp://consumer-lag-report: Análisis de rendimiento del consumidor con umbrales personalizables

Prompts MCP

El servidor incluye los siguientes prompts preconfigurados para operaciones y diagnósticos de Kafka. Para documentación detallada con argumentos y respuestas de ejemplo, consulta docs/prompts.md.

  • kafka_cluster_overview: Genera resúmenes completos de salud del clúster
  • kafka_health_check: Realiza evaluaciones detalladas de salud con recomendaciones procesables
  • kafka_under_replicated_partitions: Analiza problemas de replicación con guía de resolución de problemas
  • kafka_consumer_lag_report: Proporciona análisis de lag del consumidor con recomendaciones de optimización de rendimiento

Ejemplo de extremo a extremo

Aquí tienes un ejemplo de interacción que muestra cómo un asistente de IA podría usar estas herramientas para ayudar a un usuario:

Consulta del usuario: "Necesitamos investigar por qué nuestro pipeline de procesamiento de pedidos se está retrasando. ¿Puedes ayudarme a revisar el lag del consumidor?"

Flujo de trabajo del asistente de IA:

  1. Primero, verifica los grupos de consumidores

    > Using list_consumer_groups tool
    > Discovers "order-processor" group
    
  2. Luego, examina los detalles del grupo de consumidores

    > Using describe_consumer_group tool with group_id "order-processor" and include_offsets=true
    > Discovers lag on several partitions
    
  3. Obtén una visión general del clúster para verificar la salud general

    > Using cluster_overview tool
    > Confirms no under-replicated partitions or offline brokers
    
  4. Genera un informe de lag del consumidor

    > Using consumer_lag_report prompt with threshold=1000
    > Provides detailed analysis of lag across consumer groups
    
  5. Proporciona recomendaciones procesables:

    Based on my analysis, your order-processor consumer group has significant lag on topic orders, partitions 2 and 3, with lag values of 15,420 and 12,305 messages respectively.
    
    The cluster itself is healthy with no under-replicated partitions, suggesting this is an issue with the consumer application rather than Kafka infrastructure.
    
    I recommend:
    1. Check the consumer instances processing those partitions for slowdowns or errors
    2. Consider scaling up your consumer group by adding more instances
    3. Review the consumer configuration, particularly max.poll.records and max.poll.interval.ms
    4. Examine any processing bottlenecks in the consumer logic
    

Este flujo de trabajo fluido demuestra cómo las herramientas Kafka MCP permiten a los modelos LLM realizar diagnósticos sofisticados y proporcionar información procesable.

Configuración

El servidor se puede configurar usando las siguientes variables de entorno:

VariableDescripciónPredeterminado
KAFKA_BROKERSLista separada por comas de direcciones de brokers Kafkalocalhost:9092
KAFKA_CLIENT_IDID de cliente Kafka utilizado para las conexioneskafka-mcp-server
MCP_TRANSPORTMétodo de transporte MCP (stdio/http)stdio
KAFKA_SASL_MECHANISMMecanismo SASL: plain, scram-sha-256, scram-sha-512 o "" (deshabilitado)""
KAFKA_SASL_USERNombre de usuario para autenticación SASL""
KAFKA_SASL_PASSWORDContraseña para autenticación SASL""
KAFKA_TLS_ENABLEHabilitar TLS para la conexión Kafka (true o false)false
KAFKA_TLS_INSECURE_SKIP_VERIFYOmitir verificación de certificado TLS (true o false)false

Configuración de OAuth 2.1 (solo transporte HTTP)

Al usar transporte HTTP (MCP_TRANSPORT=http), se puede habilitar la autenticación OAuth 2.1:

VariableDescripciónPredeterminadoRequerido
MCP_HTTP_PORTPuerto del servidor HTTP8080No
OAUTH_ENABLEDHabilitar autenticación OAuth 2.1falseNo
OAUTH_MODEModo OAuth: native o proxynativeNo
OAUTH_PROVIDERProveedor: hmac, okta, google, azureoktaNo
OAUTH_SERVER_URLURL completa del servidor MCP (p. ej., https://localhost:8080)-Cuando OAuth está habilitado
OIDC_ISSUERURL del emisor OAuth-Cuando OAuth está habilitado
OIDC_AUDIENCEAudiencia OAuth-Cuando OAuth está habilitado
OIDC_CLIENT_IDID de cliente OAuth-Solo modo proxy
OIDC_CLIENT_SECRETSecreto de cliente OAuth-Solo modo proxy
OAUTH_REDIRECT_URISURIs de redirección separadas por comas-Solo modo proxy
JWT_SECRETSecreto de firma JWT-Solo modo proxy

Para configuración detallada de OAuth y ejemplos, consulta docs/oauth.md.

Notas de seguridad:

  • Al usar KAFKA_TLS_INSECURE_SKIP_VERIFY=true, el servidor omitirá la verificación del certificado TLS. Esto solo debe usarse en entornos de desarrollo o pruebas, o al usar certificados autofirmados.
  • OAuth solo está disponible al usar transporte HTTP. El transporte STDIO no admite OAuth.
  • Usa siempre HTTPS en producción cuando OAuth esté habilitado.

Consideraciones de seguridad

El servidor está diseñado con seguridad de nivel empresarial:

  • Autenticación:
    • Kafka: Soporte completo para SASL PLAIN, SCRAM-SHA-256 y SCRAM-SHA-512
    • Servidor MCP: Autenticación OAuth 2.1 para transporte HTTP (Okta, Google, Azure AD, HMAC)
  • Cifrado: Soporte TLS para comunicación segura con brokers Kafka
  • Validación de entrada: Validación exhaustiva de todas las entradas del usuario para prevenir ataques de inyección
  • Manejo de errores: Manejo seguro de errores que no expone información sensible
  • Seguridad de tokens: Validación de tokens Bearer con caché de 5 minutos para endpoints protegidos por OAuth

Para mejores prácticas de seguridad OAuth, consulta docs/oauth.md.

Desarrollo

Pruebas

Una cobertura de pruebas integral garantiza la fiabilidad:

# Run all tests (requires Docker for integration tests)
go test ./...

# Run tests excluding integration tests
go test -short ./...

# Run integration tests with specific Kafka brokers
export KAFKA_BROKERS="your-broker:9092"
export SKIP_KAFKA_TESTS="false"
go test ./kafka -v -run Test

Contribuciones

¡Las contribuciones son bienvenidas! No dudes en enviar un Pull Request.

Licencia

Este proyecto está licenciado bajo la Licencia MIT: consulta el archivo LICENSE para más detalles.