Script icon
ChatCompletion icon
GoogleGemini icon
If icon
SlackIncomingWebhook icon
Log icon
Schedule icon
Webhook icon

AI-Powered Airflow DAG to Kestra Flow Migration Agent

Automatically convert and validate legacy Apache Airflow Python DAGs into production-grade Kestra 2.0 declarative YAML flows with Google Gemini AI.

Categories
AIDataInfrastructure

Enterprise data teams seeking to modernize their orchestration infrastructure from Apache Airflow to Kestra frequently face friction when migrating hundreds of legacy Python DAG files. Manual rewrites require translating procedural Python operators (BashOperator, PythonOperator, PostgresOperator, SlackWebhookOperator, and >> dependency chains) into declarative Kestra 2.0 YAML primitives, a process prone to syntax errors, macro mismatches, and connection misconfigurations.

This blueprint provides an autonomous migration agent that ingests Airflow Python DAG code, extracts task hierarchies via AST analysis, transpiles operators into modern Kestra 2.0 flow definitions using Google Gemini with strict JSON Schema enforcement, verifies YAML validity and security standards, and posts an itemized migration scorecard to Slack.

How it works

  1. parse_dag_ast (io.kestra.plugin.scripts.python.Script) inspects the source Airflow Python code using AST to identify DAG parameters, schedules, operator types, and task names.
  2. transpile_dag_with_ai (io.kestra.plugin.ai.completion.ChatCompletion) translates Airflow operators to native Kestra 2.0 plugins (mapping PostgresOperator to io.kestra.plugin.jdbc.postgresql.Query, BashOperator to io.kestra.plugin.scripts.shell.Commands, and Airflow Jinja macros to Kestra Pebble expressions) with structured JSON Schema output.
  3. validate_transpiled_flow (io.kestra.plugin.scripts.python.Script) parses the resulting YAML to ensure schema conformance, inlined task credentials, and valid task structure.
  4. evaluate_migration_verdict (io.kestra.plugin.core.flow.If) routes valid flows to an automated Slack notification card with complexity scoring, while routing invalid syntax or unsupported hooks to an advisory log.

What you get

  • Autonomous translation of legacy Python DAGs into clean, declarative Kestra 2.0 YAML.
  • Automated operator mapping covering database queries, shell scripts, Python callables, and alerts.
  • Zero-setup demo execution via the seed_demo_data toggle with pre-populated production DAG samples.
  • Built-in schema validation to prevent malformed YAML or deprecated 1.x task structures from reaching production.

Who it's for

  • Data Platform Engineers migrating legacy infrastructure from Apache Airflow to Kestra.
  • Analytics Engineers modernizing batch data pipelines into declarative event-driven workflows.
  • Platform Engineering teams establishing automated GitOps migration pipelines.

Why orchestrate this with Kestra

Instead of maintaining custom migration scripts that rot as plugin schemas evolve, Kestra orchestrates the entire migration lifecycle declaratively. By combining Python AST extraction, Gemini AI structured schema generation, and flowable conditional evaluation within a single observable workflow, teams gain repeatable, audited pipeline modernization.

Pitfalls

  • Complex dynamic runtime dependencies: Airflow DAGs generating tasks dynamically via database queries during DAG parse time must be decoupled into Kestra dynamic flow tasks.
  • XCom cross-task data transfers: In Airflow, XCom transfers arbitrary pickled Python objects; in Kestra, outputs are structured JSON accessible via {{ outputs.task_id.vars.property }}.
  • Custom proprietary Airflow hooks: In-house company hooks require mapping to standard Kestra HTTP, CLI, or script plugins.
  • Schedule macro translation: Airflow execution date macros ({{ ds }}) map to Kestra trigger variables ({{ trigger.date }}).

Prerequisites

  • Google Gemini API key with access to gemini-2.5-flash.
  • Slack incoming webhook endpoint for migration notifications.

Secrets

  • GEMINI_API_KEY: API token for Google Gemini model access.
  • SLACK_WEBHOOK_URL: Webhook URL for delivering migration scorecards.
  • MIGRATION_WEBHOOK_KEY: Secret authentication key for the event-driven intake webhook.

Quick start

  1. Execute the flow with seed_demo_data: true to test immediate transpilation of the built-in sample DAG.
  2. Configure GEMINI_API_KEY and SLACK_WEBHOOK_URL in your Kestra namespace.
  3. Supply custom Airflow DAG source code via the airflow_dag_code input or trigger via webhook.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.