Cloud Composer and Apache Airflow: Building Production Data Pipelines

Cloud Composer is GCP's managed Apache Airflow service for orchestrating complex data pipelines. This guide covers DAG design patterns, Composer environment configuration, GCP operator integration (BigQuery, Dataflow, Dataproc), performance tuning, and monitoring with Cloud Monitoring.

Every organization that builds a data platform eventually needs to orchestrate complex workflows: run transformation A after B completes, trigger a Dataflow job when new data arrives in Cloud Storage, retry failed tasks automatically, and alert when any part of the pipeline fails for more than an hour.

Apache Airflow has become the standard orchestrator for these use cases. It defines workflows as Directed Acyclic Graphs (DAGs) — Python code that describes dependencies between tasks. Cloud Composer is GCP's managed Airflow service: it runs the Airflow infrastructure (scheduler, webserver, workers, metadata database) so you only manage the DAG files.

This guide covers practical Cloud Composer and Airflow usage: setting up a production Composer environment, writing efficient DAGs that use GCP operators, handling failures gracefully, and monitoring pipeline health.

Cloud Composer Environment Setup

# Create a Composer 2 environment (current version)
gcloud composer environments create my-composer-env   --location=europe-west4   --image-version=composer-2.5.0-airflow-2.6.3   --environment-size=ENVIRONMENT_SIZE_SMALL   --service-account=composer-sa@my-project.iam.gserviceaccount.com   --node-count=3   --scheduler-count=2   --web-server-machine-type=composer-n1-webserver-2   --cloud-sql-machine-type=db-n1-standard-2   --network=my-vpc   --subnetwork=composer-subnet

# Grant the Composer service account necessary permissions
gcloud projects add-iam-policy-binding my-project   --member=serviceAccount:composer-sa@my-project.iam.gserviceaccount.com   --role=roles/bigquery.admin

gcloud projects add-iam-policy-binding my-project   --member=serviceAccount:composer-sa@my-project.iam.gserviceaccount.com   --role=roles/dataflow.developer

Writing Production DAGs

A well-structured DAG separates configuration from logic and uses proper dependency management:

# dags/user_analytics_pipeline.py
from datetime import datetime, timedelta
from airflow import DAG
from airflow.providers.google.cloud.operators.bigquery import (
    BigQueryInsertJobOperator,
    BigQueryCheckOperator,
)
from airflow.providers.google.cloud.operators.dataflow import DataflowTemplateOperator
from airflow.providers.google.cloud.operators.gcs import GCSDeleteObjectsOperator
from airflow.providers.google.cloud.sensors.gcs import GCSObjectExistenceSensor
from airflow.utils.trigger_rule import TriggerRule

# Configuration — separate from DAG logic
PROJECT_ID = "my-project"
DATASET = "analytics"
DATAFLOW_TEMPLATE = "gs://my-project-templates/etl-template.json"

default_args = {
    "owner": "data-platform-team",
    "email": ["data-alerts@company.com"],
    "email_on_failure": True,
    "email_on_retry": False,
    "retries": 2,
    "retry_delay": timedelta(minutes=5),
    "retry_exponential_backoff": True,
    "max_retry_delay": timedelta(minutes=30),
}

with DAG(
    dag_id="user_analytics_pipeline",
    default_args=default_args,
    description="Daily user analytics ETL pipeline",
    schedule="0 6 * * *",  # 6am UTC daily
    start_date=datetime(2025, 1, 1),
    catchup=False,  # Don't backfill missed runs
    max_active_runs=1,  # Only one concurrent run
    tags=["analytics", "daily", "bigquery"],
) as dag:

    # Sensor: wait for source data to arrive
    wait_for_source_data = GCSObjectExistenceSensor(
        task_id="wait_for_source_data",
        bucket="my-project-raw-data",
        object="{{ ds_nodash }}/events.parquet",  # Date-stamped path
        timeout=60 * 60 * 4,  # Wait up to 4 hours
        poke_interval=60 * 5,  # Check every 5 minutes
        mode="reschedule",  # Don't hold a worker slot while waiting
    )

    # Validate source data quality
    validate_source = BigQueryCheckOperator(
        task_id="validate_source",
        sql="""
            SELECT COUNT(*) > 0
            FROM `{{ params.project }}.raw.events`
            WHERE DATE(event_timestamp) = DATE('{{ ds }}')
        """,
        use_legacy_sql=False,
        params={"project": PROJECT_ID},
    )

    # Run Dataflow ETL job
    run_etl = DataflowTemplateOperator(
        task_id="run_etl",
        template=DATAFLOW_TEMPLATE,
        parameters={
            "input_date": "{{ ds }}",
            "output_dataset": DATASET,
            "project_id": PROJECT_ID,
        },
        gcp_conn_id="google_cloud_default",
        location="europe-west4",
        job_name="user-events-etl-{{ ds_nodash }}",
        wait_until_finished=True,
    )

    # Transform and aggregate in BigQuery
    transform_events = BigQueryInsertJobOperator(
        task_id="transform_events",
        configuration={
            "query": {
                "query": """
                    CREATE OR REPLACE TABLE `{{ params.project }}.{{ params.dataset }}.user_events_daily`
                    PARTITION BY event_date AS
                    SELECT
                        DATE(event_timestamp) AS event_date,
                        event_type,
                        COUNT(*) AS event_count,
                        COUNT(DISTINCT user_id) AS unique_users,
                        SUM(revenue) AS total_revenue
                    FROM `{{ params.project }}.{{ params.dataset }}.user_events_raw`
                    WHERE DATE(event_timestamp) = DATE('{{ ds }}')
                    GROUP BY 1, 2
                """,
                "useLegacySql": False,
            }
        },
        params={"project": PROJECT_ID, "dataset": DATASET},
    )

    # Data quality check after transformation
    check_output = BigQueryCheckOperator(
        task_id="check_output",
        sql="""
            SELECT
                COUNT(*) > 0
                AND SUM(event_count) > 1000
                AND SUM(unique_users) > 100
            FROM `{{ params.project }}.{{ params.dataset }}.user_events_daily`
            WHERE event_date = DATE('{{ ds }}')
        """,
        use_legacy_sql=False,
        params={"project": PROJECT_ID, "dataset": DATASET},
    )

    # Clean up raw staging data
    cleanup = GCSDeleteObjectsOperator(
        task_id="cleanup_staging",
        bucket="my-project-staging",
        prefix="{{ ds_nodash }}/",
        trigger_rule=TriggerRule.ALL_SUCCESS,
    )

    # Define dependencies
    wait_for_source_data >> validate_source >> run_etl >> transform_events >> check_output >> cleanup

DAG Design Patterns

Dynamic DAGs for Multiple Datasets

from airflow import DAG
from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator
from datetime import datetime

DATASETS = ["orders", "customers", "products", "inventory"]

with DAG(
    dag_id="refresh_all_datasets",
    schedule="0 5 * * *",
    start_date=datetime(2025, 1, 1),
    catchup=False,
) as dag:

    previous_task = None

    for dataset_name in DATASETS:
        refresh_task = BigQueryInsertJobOperator(
            task_id=f"refresh_{dataset_name}",
            configuration={
                "query": {
                    "query": f"CALL analytics.refresh_{dataset_name}('{{{{ ds }}}}')",
                    "useLegacySql": False,
                }
            },
        )

        if previous_task:
            previous_task >> refresh_task

        previous_task = refresh_task

TaskGroups for Logical Grouping

from airflow.utils.task_group import TaskGroup

with DAG(dag_id="complex_pipeline", ...) as dag:
    with TaskGroup("extraction") as extraction_group:
        extract_orders = BigQueryInsertJobOperator(task_id="extract_orders", ...)
        extract_customers = BigQueryInsertJobOperator(task_id="extract_customers", ...)

    with TaskGroup("transformation") as transformation_group:
        transform_data = BigQueryInsertJobOperator(task_id="transform", ...)

    with TaskGroup("validation") as validation_group:
        validate_orders = BigQueryCheckOperator(task_id="validate_orders", ...)
        validate_customers = BigQueryCheckOperator(task_id="validate_customers", ...)

    extraction_group >> transformation_group >> validation_group

Performance Tuning

Composer 2 uses Kubernetes for worker autoscaling. Configure it appropriately:

# Update Composer environment configuration
gcloud composer environments update my-composer-env   --location=europe-west4   --scheduler-count=2        # Two schedulers for high availability
  --scheduler-cpu=2   --scheduler-memory=7.5GB   --worker-cpu=2   --worker-memory=7.5GB   --min-workers=2   --max-workers=20 
# Set Airflow configuration for concurrency
gcloud composer environments update my-composer-env   --location=europe-west4   --update-airflow-configs=core-max_active_runs_per_dag=3,core-parallelism=50,core-dag_concurrency=16

Key configuration parameters:

  • parallelism: Maximum number of tasks that can run concurrently across all DAGs
  • dag_concurrency: Maximum concurrent task instances per DAG
  • max_active_runs_per_dag: Maximum concurrent DAG runs

Monitoring Pipeline Health

# View DAG run history
gcloud composer environments run my-composer-env   --location=europe-west4   dags list-runs -- --dag-id user_analytics_pipeline

# Create Cloud Monitoring alert for failed tasks
gcloud alpha monitoring policies create   --display-name="Airflow Task Failure"   --condition-filter='resource.type="cloud_composer_environment" AND metric.type="composer.googleapis.com/environment/dag/task_instance_failure_count"'   --condition-threshold-value=1   --condition-threshold-comparison=COMPARISON_GT   --condition-duration=0s   --notification-channels=$CHANNEL_ID
# Add custom Datadog/PagerDuty/Slack alerting via callbacks
from airflow.operators.python import PythonOperator
from airflow.models import DagRun

def on_failure_callback(context):
    """Send custom alert on task failure."""
    dag_id = context["dag"].dag_id
    task_id = context["task"].task_id
    execution_date = context["execution_date"].isoformat()
    exception = str(context.get("exception", ""))

    message = f"Task failed: {dag_id}.{task_id} at {execution_date}. Error: {exception}"
    send_slack_alert("#data-alerts", message)

default_args = {
    "on_failure_callback": on_failure_callback,
    ...
}

For data pipelines that use Cloud Dataflow for processing, see our Cloud Dataflow guide. For storing pipeline outputs in BigQuery, see our BigQuery architecture guide.