MCP-Airflow-API

MCP-Airflow-API es un servidor MCP que aprovecha el Protocolo de Contexto del Modelo (MCP) para transformar las operaciones de la API REST de Apache Airflow en herramientas de lenguaje natural. Este proyecto oculta la complejidad de las estructuras de la API y permite la gestión intuitiva de clústeres de Airflow mediante comandos en lenguaje natural.

Documentación

🚀 MCP-Airflow-API

Herramienta de Código Abierto Revolucionaria para Gestionar Apache Airflow con Lenguaje Natural

License: MIT Python Docker Pulls BuyMeACoffee

Deploy to PyPI with tag PyPI PyPI - Downloads


Arquitectura e Internals (DeepWiki)

Ask DeepWiki


📋 Resumen

¿Alguna vez te has preguntado lo increíble que sería poder gestionar tus flujos de trabajo de Apache Airflow usando lenguaje natural en lugar de complejas llamadas a la API REST o manipulaciones de la interfaz web? MCP-Airflow-API es el proyecto de código abierto revolucionario que hace realidad este objetivo.

MCP-Airflow-API Screenshot


🎯 ¿Qué es MCP-Airflow-API?

MCP-Airflow-API es un servidor MCP que aprovecha el Protocolo de Contexto de Modelo (MCP) para transformar las operaciones de la API REST de Apache Airflow en herramientas de lenguaje natural. Este proyecto oculta la complejidad de las estructuras de la API y permite una gestión intuitiva de los clústeres de Airflow mediante comandos en lenguaje natural.

🆕 Soporte Multi-Versión de API (¡NUEVO!)

Ahora soporta tanto la API v1 de Airflow (2.x) como la v2 (3.0+) con selección dinámica de versión mediante variable de entorno:

  • API v1: Compatibilidad total con clústeres de Airflow 2.x (43 herramientas) - Documentación
  • API v2: Funciones mejoradas para Airflow 3.0+ incluyendo gestión de activos para programación consciente de datos (45 herramientas) - Documentación

Arquitectura Clave: Un único servidor MCP con herramientas comunes compartidas (43) más herramientas exclusivas de activos para v2 (2) - carga dinámicamente el conjunto de herramientas apropiado según la variable de entorno AIRFLOW_API_VERSION.

Enfoque tradicional (ejemplo):

curl -X GET "http://localhost:8080/api/v1/dags?limit=100&offset=0" \
  -H "Authorization: Basic YWlyZmxvdzphaXJmbG93"

Enfoque MCP-Airflow-API (lenguaje natural):

"Muéstrame los DAGs que se están ejecutando actualmente"


🚀 Inicio Rápido

📝 ¿Necesitas un clúster de Airflow de prueba? Usa nuestro proyecto complementario Airflow-Docker-Compose con soporte para entornos Airflow 2.x y Airflow 3.x.

Diagrama de Flujo del Inicio Rápido/Tutorial

Flow Diagram of Quickstart/Tutorial

🎯 Recomendado: Docker Compose (Entorno de Demostración Completo)

Para evaluación y pruebas rápidas:

git clone https://github.com/call518/MCP-Airflow-API.git
cd MCP-Airflow-API

# Configure your Airflow credentials
cp .env.example .env
# Edit .env with your Airflow API settings

# Start all services
docker-compose up -d

# Access OpenWebUI at http://localhost:3002/
# API documentation at http://localhost:8002/docs

Comenzar con OpenWebUI (Opción Docker)

📌 Nota: Las instrucciones de configuración de la interfaz web se basan en OpenWebUI v0.6.22. Las ubicaciones de los menús y la configuración pueden diferir en versiones más recientes.

  1. Accede a http://localhost:3002/
  2. Inicia sesión con la cuenta de administrador
  3. Ve a "Configuración" → "Herramientas" desde el menú superior
  4. Añade la URL de la herramienta: http://localhost:8002/airflow-api
  5. Configura tu proveedor de LLM (Ollama, OpenAI, etc.)

📦 Métodos de Instalación del Servidor MCP

Método 1: Instalación Directa desde PyPI

uvx --python 3.12 mcp-airflow-api

Método 2: Integración con el Cliente MCP Claude-Desktop

Acceso Local (modo stdio)

{
  "mcpServers": {
    "mcp-airflow-api": {
      "command": "uvx",
      "args": ["--python", "3.12", "mcp-airflow-api"],
      "env": {
        "AIRFLOW_API_VERSION": "v2",
        "AIRFLOW_API_BASE_URL": "http://localhost:8080/api",
        "AIRFLOW_API_USERNAME": "airflow",
        "AIRFLOW_API_PASSWORD": "airflow"
      }
    }
  }
}git

Acceso Remoto (modo streamable-http sin autenticación)

{
  "mcpServers": {
    "mcp-airflow-api": {
      "type": "streamable-http",
      "url": "http://localhost:8000/mcp"
    }
  }
}

Acceso Remoto (modo streamable-http con autenticación mediante token Bearer - Recomendado)

{
  "mcpServers": {
    "mcp-airflow-api": {
      "type": "streamable-http",
      "url": "http://localhost:8000/mcp",
      "headers": {
        "Authorization": "Bearer your-secure-secret-key-here"
      }
    }
  }
}

Múltiples Clústeres de Airflow con Diferentes Versiones

{
  "mcpServers": {
    "airflow-2x-cluster": {
      "command": "uvx",
      "args": ["--python", "3.12", "mcp-airflow-api"],
      "env": {
        "AIRFLOW_API_VERSION": "v1",
        "AIRFLOW_API_BASE_URL": "http://localhost:38080/api",
        "AIRFLOW_API_USERNAME": "airflow",
        "AIRFLOW_API_PASSWORD": "airflow"
      }
    },
    "airflow-3x-cluster": {
      "command": "uvx",
      "args": ["--python", "3.12", "mcp-airflow-api"],
      "env": {
        "AIRFLOW_API_VERSION": "v2",
        "AIRFLOW_API_BASE_URL": "http://localhost:48080/api",
        "AIRFLOW_API_USERNAME": "airflow",
        "AIRFLOW_API_PASSWORD": "airflow"
      }
    }
  }
}

💡 Consejo Profesional: Usa los clústeres de prueba de Airflow-Docker-Compose para la configuración anterior: se ejecutan en los puertos 38080 (2.x) y 48080 (3.x) respectivamente.

Método 3: Instalación de Desarrollo

git clone https://github.com/call518/MCP-Airflow-API.git
cd MCP-Airflow-API
pip install -e .

# Run in stdio mode
python -m mcp_airflow_api

🌟 Características Principales

  1. Consultas en Lenguaje Natural
    No necesitas aprender sintaxis compleja de API. Simplemente pregunta como lo harías naturalmente:

    • "¿Qué DAGs se están ejecutando actualmente?"
    • "Muéstrame las tareas fallidas"
    • "Encuentra DAGs que contengan ETL"
  2. Capacidades Integrales de Monitoreo
    Monitoreo del estado del clúster en tiempo real:

    • Monitoreo de salud del clúster
    • Análisis de estado y rendimiento de DAGs
    • Seguimiento de registros de ejecución de tareas
    • Gestión de datos XCom
  3. Soporte Dinámico de Versión de API
    Un único servidor MCP se adapta a tu versión de Airflow:

    • API v1: 43 herramientas compartidas para compatibilidad con Airflow 2.x
    • API v2: 43 herramientas compartidas + 2 herramientas de gestión de activos para Airflow 3.0+
    • Control mediante Variable de Entorno: Cambia de versión al instante con AIRFLOW_API_VERSION
    • Cero Cambios de Configuración: Mismos nombres de herramientas, capacidades mejoradas
    • Arquitectura Eficiente: Código base común compartido que elimina la duplicación
  4. Cobertura Integral de Herramientas
    Cubre casi toda la funcionalidad de la API de Airflow:

    • Gestión de DAGs (activar, pausar, reanudar)
    • Monitoreo de instancias de tareas
    • Gestión de pools y variables
    • Configuración de conexiones
    • Consultas de configuración
    • Análisis de registros de eventos
  5. Optimización para Entornos Grandes
    Maneja eficientemente entornos grandes con más de 1000 DAGs:

    • Soporte de paginación inteligente
    • Opciones de filtrado avanzadas
    • Capacidades de procesamiento por lotes

🛠️ Ventajas Técnicas

  • Aprovechamiento del Protocolo de Contexto de Modelo (MCP)
    MCP es un estándar abierto para conexiones seguras entre aplicaciones de IA y fuentes de datos, que proporciona:

    • Interfaz estandarizada
    • Acceso seguro a datos
    • Arquitectura escalable
  • Soporte para Dos Modos de Transporte

    • Modo stdio: Integración directa con el cliente MCP para entornos locales
    • Modo streamable-http: Despliegue basado en HTTP para Docker y acceso remoto

    Control mediante Variable de Entorno:

    FASTMCP_TYPE=stdio          # Default: Direct MCP client mode
    FASTMCP_TYPE=streamable-http # Docker/HTTP mode
    FASTMCP_PORT=8000           # HTTP server port (Docker internal)
    
  • Cobertura Integral de la API de Airflow
    Implementación completa de las API REST oficiales de Airflow:

    • Soporte API v1: Basado en la API REST de Airflow 2.x
    • Soporte API v2: Basado en la API REST de Airflow 3.0+
    • Selección Dinámica de Versión: Cambio en tiempo de ejecución entre versiones de API
    • Paridad de Funciones: Cobertura completa de endpoints para ambas versiones
  • Soporte Completo de Docker
    Configuración completa de Docker Compose con 3 servicios separados:

    • Open WebUI: Interfaz web (puerto 3002)
    • Servidor MCP: Herramientas de API de Airflow (puerto interno 8000, expuesto mediante 18002)
    • Proxy MCPO: Proveedor de endpoints de API REST (puerto 8002)

Casos de Uso en Acción

Capacity Management for Operations Teams

Capacity Management for Operations Teams

Capacity Management for Operations Teams

Capacity Management for Operations Teams

Capacity Management for Operations Teams

Capacity Management for Operations Teams

Capacity Management for Operations Teams

Capacity Management for Operations Teams

Capacity Management for Operations Teams

Capacity Management for Operations Teams

Capacity Management for Operations Teams


⚙️ Configuración Avanzada

Variables de Entorno

# Required - Dynamic API Version Selection (NEW!)
# Single server supports both v1 and v2 - just change this variable!
AIRFLOW_API_VERSION=v1           # v1 for Airflow 2.x, v2 for Airflow 3.0+
AIRFLOW_API_BASE_URL=http://localhost:8080/api

# Test Cluster Connection Examples:
# For Airflow 2.x test cluster (from Airflow-Docker-Compose)
AIRFLOW_API_VERSION=v1
AIRFLOW_API_BASE_URL=http://localhost:38080/api

# For Airflow 3.x test cluster (from Airflow-Docker-Compose)  
AIRFLOW_API_VERSION=v2
AIRFLOW_API_BASE_URL=http://localhost:48080/api

# Authentication
AIRFLOW_API_USERNAME=airflow
AIRFLOW_API_PASSWORD=airflow

# Optional - MCP Server Configuration
MCP_LOG_LEVEL=INFO                   # DEBUG/INFO/WARNING/ERROR/CRITICAL
FASTMCP_TYPE=stdio                   # stdio/streamable-http
FASTMCP_PORT=8000                    # HTTP server port (Docker mode)

# Bearer Token Authentication for streamable-http mode
# Enable authentication (recommended for production)
# Default: false (when undefined, empty, or null)
# Values: true/false, 1/0, yes/no, on/off (case insensitive)
REMOTE_AUTH_ENABLE=false             # true/false
REMOTE_SECRET_KEY=your-secure-secret-key-here

Comparación de Versiones de API

Documentación Oficial:

CaracterísticaAPI v1 (Airflow 2.x)API v2 (Airflow 3.0+)
Total de Herramientas43 herramientas45 herramientas
Herramientas Compartidas43 (100%)43 (96%)
Herramientas Exclusivas02 (Gestión de Activos)
Operaciones Básicas de DAG✅ Mejoradas
Gestión de Tareas✅ Mejorada
Gestión de Conexiones✅ Mejorada
Gestión de Pools✅ Mejorada
Gestión de ActivosNueva
Eventos de ActivosNuevos
Programación Consciente de DatosNueva
Advertencias Mejoradas de DAGNuevas
Filtrado AvanzadoBásicoMejorado

🔐 Seguridad y Autenticación

Autenticación con Token Bearer

Para el modo streamable-http, este servidor MCP soporta autenticación con token Bearer para asegurar el acceso remoto. Esto es especialmente importante cuando se ejecuta el servidor en entornos de producción.

Configuración

Habilitar Autenticación:

# In .env file
REMOTE_AUTH_ENABLE=true
REMOTE_SECRET_KEY=your-secure-secret-key-here

O mediante CLI:

python -m mcp_airflow_api --type streamable-http --auth-enable --secret-key your-secure-secret-key-here

Niveles de Seguridad

  1. Modo stdio (Predeterminado): Acceso solo local, no se necesita autenticación
  2. streamable-http + REMOTE_AUTH_ENABLE=false: Acceso remoto sin autenticación ⚠️ NO RECOMENDADO para producción
  3. streamable-http + REMOTE_AUTH_ENABLE=true: Acceso remoto con autenticación mediante token Bearer ✅ RECOMENDADO para producción

Nota: REMOTE_AUTH_ENABLE tiene como valor predeterminado false cuando no está definido, está vacío o es nulo. Los valores admitidos son true/false, 1/0, yes/no, on/off (sin distinción de mayúsculas y minúsculas).

Configuración del Cliente

Cuando la autenticación está habilitada, los clientes MCP deben incluir el token Bearer en el encabezado de Authorization:

{
  "mcpServers": {
    "mcp-airflow-api": {
      "type": "streamable-http",
      "url": "http://your-server:8000/mcp",
      "headers": {
        "Authorization": "Bearer your-secure-secret-key-here"
      }
    }
  }
}

Buenas Prácticas de Seguridad

  • Habilita siempre la autenticación al usar el modo streamable-http en producción
  • Usa claves secretas fuertes y generadas aleatoriamente (se recomiendan 32+ caracteres)
  • Usa HTTPS cuando sea posible (configura un proxy inverso con SSL/TLS)
  • Restringe el acceso a la red mediante cortafuegos o políticas de red
  • Rota las claves secretas regularmente para mayor seguridad
  • Monitorea los registros de acceso para detectar intentos de acceso no autorizado

Manejo de Errores

Cuando la autenticación falla, el servidor devuelve:

  • 401 No Autorizado para tokens faltantes o inválidos
  • Mensajes de error detallados en formato JSON para depuración

Configuración Personalizada de Docker Compose

version: '3.8'
services:
  mcp-server:
    build: 
      context: .
      dockerfile: Dockerfile.MCP-Server
    environment:
      - FASTMCP_PORT=8000
      - AIRFLOW_API_VERSION=v1
      - AIRFLOW_API_BASE_URL=http://your-airflow:8080/api
      - AIRFLOW_API_USERNAME=airflow
      - AIRFLOW_API_PASSWORD=airflow

Instalación de Desarrollo

git clone https://github.com/call518/MCP-Airflow-API.git
cd MCP-Airflow-API
pip install -e .

# Run in stdio mode
python -m mcp_airflow_api

🧪 Despliegue de Clúster de Airflow de Prueba

Para pruebas y desarrollo, usa nuestro proyecto complementario Airflow-Docker-Compose que soporta entornos Airflow 2.x y 3.x.

Configuración Rápida

  1. Clona el repositorio del entorno de prueba:
    git clone https://github.com/call518/Airflow-Docker-Compose.git
    cd Airflow-Docker-Compose
    

Opción 1: Desplegar Airflow 2.x (LTS)

Para probar la compatibilidad con API v1 con funciones de producción estables:

# Navigate to Airflow 2.x environment
cd airflow-2.x

# (Optional) Customize environment variables
cp .env.template .env
# Edit .env file as needed

# Deploy Airflow 2.x cluster
./run-airflow-cluster.sh

# Access Web UI
# URL: http://localhost:38080
# Username: airflow / Password: airflow

Detalles del entorno:

  • Imagen: apache/airflow:2.10.2
  • Puerto: 38080 (configurable mediante AIRFLOW_WEBSERVER_PORT)
  • API: endpoints /api/v1/*
  • Autenticación: Autenticación Básica
  • Caso de uso: Listo para producción, funciones estables

Opción 2: Desplegar Airflow 3.x (Última Versión)

Para probar API v2 con las últimas funciones, incluida la gestión de Activos:

# Navigate to Airflow 3.x environment  
cd airflow-3.x

# (Optional) Customize environment variables
cp .env.template .env
# Edit .env file as needed

# Deploy Airflow 3.x cluster
./run-airflow-cluster.sh

# Access API Server
# URL: http://localhost:48080
# Username: airflow / Password: airflow

Detalles del entorno:

  • Imagen: apache/airflow:3.0.6
  • Puerto: 48080 (configurable mediante AIRFLOW_APISERVER_PORT)
  • API: endpoints /api/v2/* + gestión de Activos
  • Autenticación: Token JWT (FabAuthManager)
  • Caso de uso: Desarrollo, prueba de nuevas funciones

Opción 3: Desplegar Ambas Versiones Simultáneamente

Para pruebas integrales en diferentes versiones de Airflow:

# Start Airflow 2.x (port 38080)
cd airflow-2.x && ./run-airflow-cluster.sh

# Start Airflow 3.x (port 48080) 
cd ../airflow-3.x && ./run-airflow-cluster.sh

Diferencias Clave

CaracterísticaAirflow 2.xAirflow 3.x
AutenticaciónAutenticación BásicaTokens JWT (FabAuthManager)
Puerto Predeterminado3808048080
Endpoints de API/api/v1/*/api/v2/*
Soporte de Activos❌ Limitado/Experimental✅ Soporte Completo
Paquetes de Proveedoresprovidersdistributions
Estabilidad✅ Listo para Producción🧪 Beta/Desarrollo

Limpieza

Para detener y limpiar los entornos de prueba:

# For Airflow 2.x
cd airflow-2.x && ./cleanup-airflow-cluster.sh

# For Airflow 3.x
cd airflow-3.x && ./cleanup-airflow-cluster.sh

🌈 Arquitectura Preparada para el Futuro

  • Diseño escalable y estructura modular para añadir fácilmente nuevas funciones
  • Protocolo conforme a estándares para integración con otras herramientas
  • Operaciones nativas en la nube e interfaz preparada para LLM
  • Procesamiento de consultas consciente del contexto y capacidades de gestión automatizada de flujos de trabajo

🎯 ¿Para Quién es Esta Herramienta?

  • Ingenieros de Datos — Reducen el tiempo de depuración, mejoran la productividad, minimizan la curva de aprendizaje
  • Ingenieros DevOps — Automatizan el monitoreo de infraestructura, reducen el tiempo de respuesta ante incidentes
  • Administradores de Sistemas — Gestión fácil de usar sin APIs complejas, monitoreo del estado del clúster en tiempo real

🚀 Contribución de Código Abierto y Comunidad

Repositorio: https://github.com/call518/MCP-Airflow-API

Cómo Contribuir

  • Informes de errores y sugerencias de funciones
  • Mejoras de documentación
  • Contribuciones de código

Considera dar una estrella al proyecto si te resulta útil.


🔮 Conclusión

MCP-Airflow-API cambia el paradigma de la ingeniería de datos y la gestión de flujos de trabajo:
No necesitas memorizar llamadas a la API REST — simplemente pregunta en lenguaje natural:

"Muéstrame el estado de los trabajos ETL que se están ejecutando actualmente."


🏷️ Etiquetas

#Apache-Airflow #MCP #ModelContextProtocol #DataEngineering #DevOps #WorkflowAutomation #NaturalLanguage #OpenSource #Python #Docker #AI-Integration


📚 Ejemplos de Consultas y Casos de Uso

Esta sección proporciona ejemplos completos de cómo usar las herramientas de MCP-Airflow-API con consultas en lenguaje natural.

Operaciones Básicas de DAG

  • list_dags: "Listar todos los DAGs con límite de 10 en formato de tabla." → Devuelve hasta 10 DAGs
  • list_dags: "Listar todos los DAGs en formato de tabla." → Devuelve todos los DAGs (ADVERTENCIA: Se necesitan muchos tokens)
  • list_dags: "Mostrar la siguiente página de DAGs." → Usar offset para paginación
  • list_dags: "Listar DAGs 21-40." → list_dags(limit=20, offset=20)
  • list_dags: "Filtrar DAGs cuyo ID contenga 'tutorial'." → list_dags(id_contains="etl")
  • list_dags: "Filtrar DAGs cuyo nombre visible contenga 'tutorial'." → list_dags(name_contains="daily")
  • get_dags_detailed_batch: "Obtener información detallada de todos los DAGs con estado de ejecución." → get_dags_detailed_batch(fetch_all=True)
  • get_dags_detailed_batch: "Obtener detalles de DAGs activos y no pausados con ejecuciones recientes." → get_dags_detailed_batch(is_active=True, is_paused=False)
  • get_dags_detailed_batch: "Obtener información detallada de DAGs que contengan 'example' con historial de ejecuciones." → get_dags_detailed_batch(id_contains="example", limit=50)
  • running_dags: "Mostrar DAGs en ejecución."
  • failed_dags: "Mostrar DAGs fallidos."
  • trigger_dag: "Disparar el DAG 'example_complex'."
  • pause_dag: "Pausar el DAG 'example_complex' en formato de tabla."
  • unpause_dag: "Reanudar el DAG 'example_complex' en formato de tabla."

Gestión de Clúster y Salud

  • get_health: "Verificar la salud del clúster de Airflow."
  • get_version: "Obtener información de la versión de Airflow."

Gestión de Pools

  • list_pools: "Listar todos los pools."
  • list_pools: "Mostrar estadísticas de uso de pools."
  • get_pool: "Obtener detalles del pool 'default_pool'."
  • get_pool: "Verificar la utilización del pool."

Gestión de Variables

  • list_variables: "Listar todas las variables."
  • list_variables: "Mostrar todas las variables de Airflow con sus valores."
  • get_variable: "Obtener la variable 'database_url'."
  • get_variable: "Mostrar el valor de la variable 'api_key'."

Gestión de Instancias de Tareas

  • list_task_instances_all: "Listar todas las instancias de tareas del DAG 'example_complex'."
  • list_task_instances_all: "Mostrar instancias de tareas en ejecución."
  • list_task_instances_all: "Mostrar instancias de tareas filtradas por el pool 'default_pool'."
  • list_task_instances_all: "Listar instancias de tareas con duración mayor a 300 segundos."
  • list_task_instances_all: "Mostrar instancias de tareas fallidas de la semana pasada."
  • list_task_instances_all: "Listar instancias de tareas fallidas de ayer."
  • list_task_instances_all: "Mostrar instancias de tareas que comenzaron después de las 9 AM de hoy."
  • list_task_instances_all: "Listar instancias de tareas de los últimos 3 días con estado 'failed'."
  • get_task_instance_details: "Obtener detalles de la tarea 'data_processing' en el DAG 'example_complex' ejecución 'scheduled__xxxxx'."
  • list_task_instances_batch: "Listar instancias de tareas fallidas del mes pasado."
  • list_task_instances_batch: "Mostrar instancias de tareas en lote para múltiples DAGs de esta semana."
  • get_task_instance_extra_links: "Obtener enlaces adicionales para la tarea 'data_processing' en la última ejecución."
  • get_task_instance_logs: "Recuperar registros de la tarea 'create_entry_gcs' intento número 2 del DAG 'example_complex'."

Gestión de XCom

  • list_xcom_entries: "Listar entradas XCom para la tarea 'data_processing' en el DAG 'example_complex' ejecución 'scheduled__xxxxx'."
  • list_xcom_entries: "Mostrar todas las entradas XCom para la tarea 'data_processing' en la última ejecución."
  • get_xcom_entry: "Obtener entrada XCom con clave 'result' para la tarea 'data_processing' en una ejecución específica."
  • get_xcom_entry: "Recuperar el valor XCom para la clave 'processed_count' de la tarea 'data_processing'."

Gestión de Configuración

  • get_config: "Mostrar todas las secciones y opciones de configuración de Airflow." → Devuelve la configuración completa o 403 si expose_config=False
  • list_config_sections: "Listar todas las secciones de configuración con información resumida."
  • get_config_section: "Obtener todos los ajustes de la sección 'core'." → get_config_section("core")
  • get_config_section: "Mostrar opciones de configuración del servidor web." → get_config_section("webserver")
  • search_config_options: "Encontrar todas las opciones de configuración relacionadas con la base de datos." → search_config_options("database")
  • search_config_options: "Buscar ajustes de tiempo de espera en la configuración." → search_config_options("timeout")

Importante: Las herramientas de configuración requieren expose_config = True en airflow.cfg sección [webserver]. Incluso los usuarios administradores reciben errores 403 si esto está deshabilitado.

Análisis y Monitoreo de DAGs

  • get_dag: "Obtener detalles del DAG 'example_complex'."
  • get_dags_detailed_batch: "Obtener detalles completos de todos los DAGs con historial de ejecución." → get_dags_detailed_batch(fetch_all=True)
  • get_dags_detailed_batch: "Obtener detalles de DAGs activos con información de la última ejecución." → get_dags_detailed_batch(is_active=True)
  • get_dags_detailed_batch: "Obtener información detallada de DAGs ETL con datos de ejecución recientes." → get_dags_detailed_batch(id_contains="etl")

Nota: get_dags_detailed_batch devuelve cada DAG con detalles de configuración (de get_dag()) y un campo latest_dag_run que contiene la información de ejecución más reciente (run_id, estado, fecha de ejecución, fecha de inicio, fecha de fin, etc.).

  • dag_graph: "Mostrar el grafo de tareas del DAG 'example_complex'."
  • list_tasks: "Listar todas las tareas del DAG 'example_complex'."
  • dag_code: "Obtener el código fuente del DAG 'example_complex'."
  • list_event_logs: "Listar registros de eventos del DAG 'example_complex'."
  • list_event_logs: "Mostrar registros de eventos con ID de ayer para todos los DAGs."
  • get_event_log: "Obtener la entrada de registro de eventos con ID 12345."
  • all_dag_event_summary: "Mostrar resumen de conteo de eventos para todos los DAGs."
  • list_import_errors: "Listar errores de importación con ID."
  • get_import_error: "Obtener error de importación con ID 67890."
  • all_dag_import_summary: "Mostrar resumen de errores de importación para todos los DAGs."
  • dag_run_duration: "Obtener estadísticas de duración de ejecución del DAG 'example_complex'."
  • dag_task_duration: "Mostrar la última ejecución del DAG 'example_complex'."
  • dag_task_duration: "Mostrar duraciones de tareas de la última ejecución de 'manual__xxxxx'."
  • dag_calendar: "Obtener información del calendario del DAG 'example_complex' del mes pasado."
  • dag_calendar: "Mostrar el horario del DAG 'example_complex' de esta semana."

Ejemplos de Cálculo de Fechas

Las herramientas calculan automáticamente las fechas relativas basándose en la fecha/hora actual del servidor:

Entrada del UsuarioMétodo de CálculoFormato de Ejemplo
"ayer"fecha_actual - 1 díaAAAA-MM-DD (1 día antes de la fecha actual)
"semana pasada"fecha_actual - 7 días hasta fecha_actual - 1 díaAAAA-MM-DD a AAAA-MM-DD (rango de 7 días)
"últimos 3 días"fecha_actual - 3 días hasta fecha_actualAAAA-MM-DD a AAAA-MM-DD (rango de 3 días)
"esta mañana"fecha_actual 00:00 a 12:00Formato AAAA-MM-DDTHH:mm:ssZ

El servidor siempre usa su fecha/hora actual para estos cálculos.

Gestión de Activos (Solo API v2)

Disponible solo cuando AIRFLOW_API_VERSION=v2 (Airflow 3.0+):

  • list_assets: "Mostrar todos los activos registrados en el sistema." → Lista todos los activos de datos para programación basada en datos
  • list_assets: "Encontrar activos con URI que contenga 's3://data-lake'." → list_assets(uri_pattern="s3://data-lake")
  • list_asset_events: "Mostrar eventos de activos recientes." → Lista cuándo se crearon o actualizaron los activos
  • list_asset_events: "Mostrar eventos de activos para un URI específico." → list_asset_events(asset_uri="s3://bucket/file.csv")
  • list_asset_events: "Encontrar eventos producidos por DAGs ETL." → list_asset_events(source_dag_id="etl_pipeline")

Ejemplos de Programación Basada en Datos:

  • "Muéstrame qué activos disparan el DAG customer_analysis."
  • "Lista todos los activos creados por el DAG data_ingestion esta semana."
  • "Encuentra activos que no se hayan actualizado recientemente."
  • "Muestra el linaje de datos de nuestro pipeline de entrenamiento de ML."

Contribuciones

🤝 ¿Tienes ideas? ¿Encontraste errores? ¿Quieres agregar funciones interesantes?

¡Siempre estamos emocionados de dar la bienvenida a nuevos contribuyentes! Ya sea que corrijas un error tipográfico, agregues una nueva herramienta de monitoreo o mejores la documentación, cada contribución hace que este proyecto sea mejor.

Formas de contribuir:

  • 🐛 Reportar problemas o errores
  • 💡 Sugerir nuevas funciones de monitoreo de Airflow
  • 📝 Mejorar la documentación
  • 🚀 Enviar solicitudes de extracción
  • ⭐ ¡Marca el repositorio con una estrella si te resulta útil!

Consejo profesional: El código base está diseñado para ser muy amigable al agregar nuevas herramientas. Revisa las funciones @mcp.tool() existentes en airflow_api.py.


🛠️ Agregar Herramientas Personalizadas (Avanzado)

Este servidor MCP está diseñado para una fácil extensibilidad. Después de explorar las funciones principales y el inicio rápido, puedes agregar tus propias herramientas personalizadas de la siguiente manera:

Guía Paso a Paso

1. Agregar Funciones Auxiliares (Opcional)

Agrega funciones de datos reutilizables a src/mcp_airflow_api/functions.py:

async def get_your_custom_data(target_resource: str = None) -> List[Dict[str, Any]]:
  """Your custom data retrieval function."""
  # Example implementation - adapt to your service
  data_source = await get_data_connection(target_resource)
  results = await fetch_data_from_source(
    source=data_source,
    filters=your_conditions,
    aggregations=["count", "sum", "avg"],
    sorting=["count DESC", "timestamp ASC"]
  )
  return results

2. Crear Tu Herramienta MCP

Agrega tu función de herramienta a src/mcp_airflow_api/airflow_api.py:

@mcp.tool()
async def get_your_custom_analysis(limit: int = 50, target_name: Optional[str] = None) -> str:
  """
  [Tool Purpose]: Brief description of what your tool does
    
  [Exact Functionality]:
  - Feature 1: Data aggregation and analysis
  - Feature 2: Resource monitoring and insights
  - Feature 3: Performance metrics and reporting
    
  [Required Use Cases]:
  - When user asks "your specific analysis request"
  - Your business-specific monitoring needs
    
  Args:
    limit: Maximum results (1-100)
    target_name: Target resource/service name
    
  Returns:
    Formatted analysis results
  """
  try:
    limit = max(1, min(limit, 100))  # Always validate input
    results = await get_your_custom_data(target_resource=target_name)
    if results:
      results = results[:limit]
    return format_table_data(results, f"Custom Analysis (Top {len(results)})")
  except Exception as e:
    logger.error(f"Failed to get custom analysis: {e}")
    return f"Error: {str(e)}"

3. Actualizar Importaciones (Si es Necesario)

Agrega tu función auxiliar a las importaciones en src/mcp_airflow_api/airflow_api.py:

from .functions import (
  # ...existing imports...
  get_your_custom_data,  # Add your new function
)

4. Actualizar la Plantilla de Prompts (Recomendado)

Agrega la descripción de tu herramienta a src/mcp_airflow_api/prompt_template.md para un mejor reconocimiento del lenguaje natural:

### **Your Custom Analysis Tool**

### X. **get_your_custom_analysis**
**Purpose**: Brief description of what your tool does
**Usage**: "Show me your custom analysis" or "Get custom analysis for database_name"
**Features**: Data aggregation, resource monitoring, performance metrics
**Required**: `target_name` parameter for specific resource analysis

5. Probar Tu Herramienta

# Local testing
./scripts/run-mcp-inspector-local.sh

# Or with Docker
docker-compose up -d
docker-compose logs -f mcp-server

# Test with natural language:
# "Show me your custom analysis"
# "Get custom analysis for target_name"

¡Eso es todo! Tu herramienta personalizada está lista para usarse con consultas en lenguaje natural.

Licencia

Úsalo, modifícalo y distribúyelo libremente bajo la Licencia MIT.