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.

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:
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:
| Setting | Default | Optimized | Impact |
|---|---|---|---|
| Instance tier | db-n1-standard-4 | db-custom-8-32768 | -60% CPU |
| max_connections | 100 | 400 | No connection errors |
| shared_buffers | 1GB | 8GB | -40% disk reads |
| work_mem | 4MB | 32MB | Faster sorts |
| maintenance_work_mem | 64MB | 512MB | Faster vacuum |
| effective_cache_size | 4GB | 24GB | Better query plans |
| Dead row cleanup | Default autovacuum | Aggressive schedule | Prevented 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 < 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:
| Metric | Target | Alert 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
| Metric | At 800 DAGs (before) | At 5,400 DAGs (after) |
|---|---|---|
| Scheduler loop time | 32s | 4.2s |
| Task latency P50 | 45s | 6s |
| Task latency P99 | 180s | 22s |
| DB CPU utilization | 91% | 48% |
| Monthly Composer cost | $3,200 | $8,400 |
| DAGs per dollar | 0.25 | 0.64 |
| Tasks executed/day | 12,000 | 84,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.
Recommended reading

Per-Team Cost Allocation in Shared Kubernetes Clusters: From Chaos to Clarity
Implementing accurate per-namespace cost allocation in multi-tenant Kubernetes clusters, covering request vs. usage attribution, shared resource amortization, and building showback dashboards that drive accountability.

Measuring and Eliminating Toil: From 40% to 12% of Engineering Time
A systematic approach to identifying, measuring, and automating toil—the repetitive operational work that scales linearly with service growth and prevents engineers from doing creative work.

Serverless Postgres in Production: Branching, Scale-to-Zero, and the End of Database Provisioning
Running Neon serverless Postgres in production for 8 months — covering database branching workflows, scale-to-zero economics, connection pooling, and migration from RDS.

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