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 DAGsdag_concurrency: Maximum concurrent task instances per DAGmax_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.