Running 5000+ DAGs in Cloud Composer: Scaling Airflow on GCP Without the Pain

How to scale Cloud Composer to handle thousands of DAGs with optimized scheduler configuration, worker autoscaling, and operational patterns that prevent failure.

#gcp#cloud-composer#airflow#data-engineering
Cover image for the article: Running 5000+ DAGs in Cloud Composer: Scaling Airflow on GCP Without the Pain

Cloud Composer is Google's managed Apache Airflow. It handles the infrastructure — GKE cluster, Cloud SQL metadata database, Redis for task queuing — so you can focus on writing DAGs. That promise holds until you hit 500 DAGs. Then the scheduler starts falling behind, the database becomes a bottleneck, and task latency balloons from seconds to minutes.

We scaled our Cloud Composer environment from 200 DAGs to 5,400 DAGs over 14 months. Here's every configuration change, architectural decision, and operational pattern that made it possible.

The Symptoms of Scale Failure

At 800 DAGs with default Composer 2 settings, we hit these symptoms:

  • Scheduler loop time exceeded 30 seconds (target: < 5 seconds)
  • Task latency (time from scheduled to running) averaged 45 seconds, spiking to 3+ minutes
  • Database CPU sustained at 90%+ on Cloud SQL
  • DAG parse time totaled 180 seconds per scheduler loop
  • Worker pods constantly pending due to GKE node pool scaling lag

Cloud Composer 2 Architecture

Understanding the architecture is essential for targeted optimization:

Cloud Composer 2 Architecture

Key components:

  • Scheduler: Parses DAGs, creates task instances, assigns to executors. Runs as 1-N pods.
  • Workers: Execute tasks via KubernetesExecutor or CeleryExecutor pods.
  • Triggerer: Handles deferred operators (async waiting without consuming worker slots).
  • Web Server: Airflow UI and API.
  • Metadata DB: Cloud SQL PostgreSQL storing all DAG state.
  • DAG Storage: GCS bucket synced to all pods.

Scheduler Optimization

The scheduler is the bottleneck in 90% of scale issues. Key configuration:

# airflow_config_overrides in Terraform
resource "google_composer_environment" "production" {
  name   = "data-platform-prod"
  region = "us-central1"

  config {
    software_config {
      image_version = "composer-2.9.0-airflow-2.9.3"

      airflow_config_overrides = {
        # Scheduler performance
        "scheduler-parsing_processes"         = "8"
        "scheduler-max_tis_per_query"         = "512"
        "scheduler-schedule_after_task_execution" = "False"
        "scheduler-min_file_process_interval" = "60"
        "scheduler-dag_dir_list_interval"     = "120"
        "scheduler-scheduler_heartbeat_sec"   = "5"

        # Reduce DB load
        "core-dagbag_import_timeout"          = "120"
        "core-min_serialized_dag_update_interval" = "60"
        "core-min_serialized_dag_fetch_interval"  = "30"

        # Task execution
        "core-parallelism"                    = "256"
        "core-max_active_tasks_per_dag"       = "32"
        "core-max_active_runs_per_dag"        = "4"

        # Performance tuning
        "scheduler-use_job_schedule"          = "True"
        "core-dag_file_processor_timeout"     = "180"
      }
    }

    workloads_config {
      scheduler {
        cpu    = 4
        memory_gb = 16
        storage_gb = 10
        count  = 3  # Multiple schedulers for HA and throughput
      }
      triggerer {
        cpu    = 2
        memory_gb = 4
        count  = 2
      }
      web_server {
        cpu    = 2
        memory_gb = 8
      }
      worker {
        cpu    = 4
        memory_gb = 16
        storage_gb = 20
        min_count = 4
        max_count = 40
      }
    }

    environment_size = "ENVIRONMENT_SIZE_LARGE"
  }
}

Critical Parameters Explained

parsing_processes = 8: Number of parallel processes parsing DAG files. Match this to scheduler CPU count. Each process handles a subset of DAG files.

min_file_process_interval = 60: Don't re-parse a DAG file more frequently than every 60 seconds. Default is 30 — increasing this cuts scheduler CPU by 40% with negligible freshness impact.

max_tis_per_query = 512: How many task instances the scheduler processes per loop iteration. Higher values mean fewer database round-trips but longer individual queries.

DAG Design Patterns for Scale

Pattern 1: Dynamic DAG Generation with Lazy Loading

Don't generate 5,000 separate DAG files. Use dynamic generation with parameter files:

# dags/generator/dynamic_etl_dags.py
"""Generate ETL DAGs from configuration without individual files."""

import json
import os
from airflow import DAG
from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator
from airflow.providers.google.cloud.transfers.gcs_to_bigquery import GCSToBigQueryOperator
from datetime import datetime, timedelta

# Load configurations from a single JSON file (not separate DAG files)
CONFIG_PATH = os.path.join(os.path.dirname(__file__), "etl_configs.json")

with open(CONFIG_PATH) as f:
    etl_configs = json.load(f)

default_args = {
    "owner": "data-platform",
    "depends_on_past": False,
    "email_on_failure": True,
    "email_on_retry": False,
    "retries": 2,
    "retry_delay": timedelta(minutes=5),
}


def create_etl_dag(config: dict) -> DAG:
    """Factory function to create a DAG from config."""
    dag_id = f"etl_{config['source']}_{config['destination']}"

    dag = DAG(
        dag_id=dag_id,
        default_args=default_args,
        description=config.get("description", f"ETL: {config['source']} -> {config['destination']}"),
        schedule_interval=config.get("schedule", "@hourly"),
        start_date=datetime(2026, 1, 1),
        catchup=False,
        max_active_runs=config.get("max_active_runs", 2),
        tags=["etl", config["source"], config.get("team", "data")],
    )

    with dag:
        extract = GCSToBigQueryOperator(
            task_id="extract_load",
            bucket=config["gcs_bucket"],
            source_objects=[config["source_path"]],
            destination_project_dataset_table=config["staging_table"],
            source_format=config.get("format", "NEWLINE_DELIMITED_JSON"),
            write_disposition="WRITE_TRUNCATE",
        )

        transform = BigQueryInsertJobOperator(
            task_id="transform",
            configuration={
                "query": {
                    "query": config["transform_sql"],
                    "useLegacySql": False,
                    "destinationTable": {
                        "projectId": config["project"],
                        "datasetId": config["dataset"],
                        "tableId": config["destination"],
                    },
                    "writeDisposition": "WRITE_TRUNCATE",
                }
            },
        )

        extract >> transform

    return dag


# Register all DAGs in the global namespace
for config in etl_configs:
    dag_id = f"etl_{config['source']}_{config['destination']}"
    globals()[dag_id] = create_etl_dag(config)

This approach reduces parse time from O(n) file reads to O(1) file read + O(n) lightweight object creation.

Pattern 2: Deferred Operators for Long-Running Tasks

Deferred operators free worker slots during waiting periods:

# dags/operators/deferred_bq_check.py
from airflow.triggers.base import BaseTrigger, TriggerEvent
from airflow.models import BaseOperator
from typing import Any, AsyncIterator
import asyncio
from google.cloud import bigquery


class BigQueryJobTrigger(BaseTrigger):
    """Async trigger that polls BigQuery job status without consuming a worker."""

    def __init__(self, job_id: str, project_id: str, poll_interval: int = 30):
        super().__init__()
        self.job_id = job_id
        self.project_id = project_id
        self.poll_interval = poll_interval

    def serialize(self) -> tuple[str, dict[str, Any]]:
        return (
            "dags.operators.deferred_bq_check.BigQueryJobTrigger",
            {"job_id": self.job_id, "project_id": self.project_id, "poll_interval": self.poll_interval},
        )

    async def run(self) -> AsyncIterator[TriggerEvent]:
        client = bigquery.Client(project=self.project_id)
        while True:
            job = client.get_job(self.job_id)
            if job.state == "DONE":
                if job.error_result:
                    yield TriggerEvent({"status": "error", "message": str(job.error_result)})
                else:
                    yield TriggerEvent({"status": "success", "job_id": self.job_id})
                return
            await asyncio.sleep(self.poll_interval)

Database Optimization

Cloud SQL backing the metadata database is the second bottleneck. Our optimizations:

SettingDefaultOptimizedImpact
Instance tierdb-n1-standard-4db-custom-8-32768-60% CPU
max_connections100400No connection errors
shared_buffers1GB8GB-40% disk reads
work_mem4MB32MBFaster sorts
maintenance_work_mem64MB512MBFaster vacuum
effective_cache_size4GB24GBBetter query plans
Dead row cleanupDefault autovacuumAggressive schedulePrevented bloat

Additionally, we run a daily maintenance DAG:

# dags/maintenance/db_cleanup.py
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.utils.db import provide_session
from datetime import datetime, timedelta

@provide_session
def cleanup_old_task_instances(session=None, **kwargs):
    """Remove task instances older than 90 days to prevent table bloat."""
    from airflow.models import TaskInstance, DagRun
    cutoff = datetime.utcnow() - timedelta(days=90)

    deleted = session.query(TaskInstance).filter(
        TaskInstance.execution_date &#x3C; cutoff,
        TaskInstance.state.in_(["success", "skipped"])
    ).delete(synchronize_session=False)

    session.commit()
    print(f"Cleaned up {deleted} old task instances")

Monitoring Dashboard

Key metrics to track at scale:

MetricTargetAlert Threshold
Scheduler loop duration< 5s> 15s
DAG parse time (total)< 30s> 60s
Task latency (scheduled → running)< 10s> 60s
Worker pod pending time< 20s> 60s
Metadata DB CPU< 60%> 80%
Metadata DB connections< 200> 350
GCS sync lag< 30s> 120s
Zombie tasks detected/hour< 5> 20

Results at 5,400 DAGs

MetricAt 800 DAGs (before)At 5,400 DAGs (after)
Scheduler loop time32s4.2s
Task latency P5045s6s
Task latency P99180s22s
DB CPU utilization91%48%
Monthly Composer cost$3,200$8,400
DAGs per dollar0.250.64
Tasks executed/day12,00084,000

The cost increased 2.6x while throughput increased 7x — a 2.7x improvement in cost efficiency.

Conclusion

Scaling Cloud Composer to 5,000+ DAGs requires coordinated optimization across scheduler configuration, DAG design patterns, database tuning, and worker autoscaling. The single most impactful change was switching from individual DAG files to dynamic generation — it cut parse time by 80%. The second was properly sizing and tuning the metadata database.

Don't wait until the scheduler is falling behind to address these issues. Proactively configure for your next order of magnitude, not your current one.

Comments

    No comments yet. Be the first to share your thoughts.