⚙️ Category 07 • Workflows, CI/CD & Data Governance

Orchestration: Apache Airflow, CI/CD & Governance

Directing complex, mission-critical data pipelines. Master Apache Airflow DAG design, AWS MWAA enterprise scaling, automated CI/CD testing with GitHub Actions, Great Expectations data contracts, and end-to-end data lineage governance.

1. Pipeline Orchestration in Enterprise Data Platforms

The Nervous System

A data pipeline without robust orchestration is a ticking time bomb. Cron jobs fail silently, lack dependency awareness, and offer zero retry policies or execution auditability. Modern data platforms rely on Apache Airflow to schedule, monitor, and manage complex Directed Acyclic Graphs (DAGs) across heterogeneous technologies (Python, Spark, Kafka, Snowflake, dbt).

Directed Acyclic Graphs (DAGs)

Workflows represented as vertices (tasks) and directed edges (dependencies) with zero circular loops, ensuring clean linear execution paths.

Idempotency & Backfills

Every DAG run must produce identical results when re-executed for historical execution dates (logical_date).

Automated CI/CD Gates

Linting DAGs with SQLFluff and Ruff, running PyTest unit tests, and verifying syntax before syncing with production S3 buckets.

Data Quality Assertions

Circuit-breaker validation stopping pipeline execution before downstream executive dashboards are corrupted by anomalous data.

Topic 7.1

Apache Airflow Fundamentals & Core Architecture

Distributed Orchestration

Airflow distributes task execution across scalable worker processes coordinated by a central scheduler and relational metastore.

Airflow Scheduler

Monitors all tasks and DAGs, triggers task instances whose dependencies are met, and passes jobs to the configured executor.

Metadata Database (PostgreSQL/MySQL)

Stores DAG definitions, task instance statuses (queued, running, success, failed), XCom variables, user credentials, and run logs.

Airflow Webserver

Flask-based UI rendering Grid view, Graph view, Gantt execution timings, task log inspector, and manual trigger controls.

Executors (Celery vs Kubernetes)

CeleryExecutor: Fixed pool of long-running worker processes via Redis/RabbitMQ queue. KubernetesExecutor: Spawns an isolated ephemeral pod per task.

Topic 7.2

DAGs, Operators, Sensors & Failure Callbacks

DAG Engineering

Designing resilient, production-grade DAGs using TaskFlow API (@task), external sensors, and retry configurations.

Operators vs Sensors

Operators perform an action (e.g. PythonOperator, SparkSubmitOperator). Sensors wait for an external condition (e.g. S3 file arrival) to become true.

Sensor Mode: Poke vs Reschedule

Poke mode: Holds open an executor worker slot while waiting (wasteful). Reschedule mode: Frees the worker slot and sleeps between checks (mandatory for long waits).

Retries with Exponential Backoff

Configure retries=3, retry_delay=timedelta(minutes=2), and retry_exponential_backoff=True to survive transient network or DB glitches.

Failure Callbacks & SLA Alerts

Triggering on_failure_callback functions that send automated PagerDuty incidents and Slack webhook messages with traceback links.

📄 production_mwaa_dag.py
from datetime import datetime, timedelta
from airflow import DAG
from airflow.sensors.filesystem import FileSensor
from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor
from airflow.operators.python import PythonOperator
from airflow.utils.task_group import TaskGroup

def alert_slack_failure(context):
    """Callback posting task failure alert with execution metadata to Slack."""
    task_id = context['task_instance'].task_id
    exec_date = context['execution_date']
    print(f"CRITICAL: Task {task_id} failed for execution date {exec_date}. Alert sent.")

default_args = {
    'owner': 'pranay_sarode',
    'depends_on_past': False,
    'email_on_failure': True,
    'email': ['padmabushanenterpirses@pranaysarode.com'],
    'retries': 3,
    'retry_delay': timedelta(minutes=3),
    'retry_exponential_backoff': True,
    'on_failure_callback': alert_slack_failure,
}

with DAG(
    dag_id='daily_capital_group_reconciliation',
    default_args=default_args,
    description='Reconcile internal trades vs custodian positions with automated DLQ',
    schedule_interval='0 6 * * 1-5', # 06:00 UTC Monday-Friday
    start_date=datetime(2026, 1, 1),
    catchup=False,
    tags=['production', 'financial', 'capital_group']
) as dag:

    # 1. Wait for Custodian Settlement file on S3 in Reschedule mode
    wait_for_custodian_file = S3KeySensor(
        task_id='wait_for_custodian_settlement',
        bucket_name='enterprise-custodian-landing',
        bucket_key='settlement_.csv',
        mode='reschedule',          # Releases worker slot while sleeping
        poke_interval=120,          # Check every 2 minutes
        timeout=3600                # Timeout after 1 hour
    )

    with TaskGroup("reconciliation_group") as recon_group:
        def run_pyspark_reconciliation():
            print("Executing PySpark 3-Way NAV Reconciliation Job...")

        execute_recon = PythonOperator(
            task_id='execute_3way_recon',
            python_callable=run_pyspark_reconciliation
        )

    wait_for_custodian_file >> recon_group
Topic 7.3

CI/CD, Automated Testing & GitOps for Pipelines

DevOps & Testing

Deploying pipeline code directly to production without testing is unacceptable. A modern GitOps workflow incorporates automated DAG syntax checking, SQL linting, and PySpark unit tests in GitHub Actions or Jenkins.

DAG Validation Testing

PyTest scripts that import every DAG file in the repo, ensuring zero import errors, missing task dependencies, or circular loops.

SQLFluff Linting

Automated linter enforcing consistent SQL formatting, column alias naming conventions, and preventing anti-patterns (e.g. SELECT *).

Blue/Green & Canary Rollouts

Deploying DAGs to a staging MWAA environment before triggering automated synchronization with the production S3 DAGs bucket.

Secret Management via KMS

Never hardcode credentials in Airflow variables. Integrate Airflow with AWS Secrets Manager or HashiCorp Vault backends.

Topic 7.4

Data Quality Frameworks & Column-Level Lineage

Data Governance

Data governance provides visibility into data lineage, data contracts, and quality assertions before downstream consumption.

Great Expectations

Declarative assertions (e.g. expect_column_values_to_not_be_null, expect_column_values_to_be_between) validated inside Airflow tasks.

OpenLineage & Marquez

Open standard for metadata collection. Automatically emits run, job, and dataset lineage events to trace data origins from source to dashboard.

Data Contracts

Agreements between data producers and data consumers defining explicit schemas, SLA delivery times, and null percentage guarantees.

Data Catalogs (DataHub / Amundsen)

Searchable enterprise catalog providing data discovery, ownership metadata, schema change histories, and compliance classification.

2. Production MWAA Orchestration & CI/CD Blueprint

Enterprise Architecture
[GitHub / GitLab Repo] → [GitHub Actions: PyTest • SQLFluff • MyPy] ↓ (Artifact Sync) [AWS S3 DAGs & Requirements Bucket] ↓ [Amazon MWAA (Managed Airflow)] • Auto-scaling Workers (Celery) • Secrets Manager • Task 1: S3KeySensor (Reschedule) • Task 2: Great Expectations Validation Gate • Task 3: PySpark EMR / Databricks Execution • Task 4: OpenLineage Metadata Emission → DataHub Catalog • On Failure: PagerDuty & Slack Webhook Callback
Hands-On Lab

Lab 07: Deploy an Airflow Pipeline with Great Expectations & Slack Alerts

Lab Guide
1

Create Great Expectations Suite

Define an expectation suite validating trade price bounds, non-null ISIN codes, and valid exchange identifier codes.

2

Build Airflow DAG with Circuit Breaker

Create an Airflow task running the validation suite. If validation fails, halt execution and route anomalous records to S3 DLQ.

3

Configure Slack Webhook Notifications

Implement an on_failure_callback formatting an automated Slack alert with the failing task ID and error message.

4

Automate CI/CD Validation

Configure a GitHub Actions workflow that runs pytest tests/test_dags.py on every pull request before merging.

3. Orchestration & Airflow Interview Questions

Orchestration Scenarios
Q1: Why should you avoid top-level code and database queries in Airflow DAG files? ▲

Answer: The Airflow Scheduler continuously parses every Python file in the DAG directory every few seconds (governed by min_file_process_interval). Any code written outside of task callables (e.g. database connections, API requests, heavy computation) executes on every single scheduler heartbeat.
Consequences: Massive CPU saturation on the scheduler, high database connection pool exhaustion, and DAG parsing timeouts. Always place dynamic logic, DB queries, and heavy imports inside operators or @task callables.

Q2: What is the difference between execution_date (logical_date) and start_date in Airflow? ▼

Answer:
• logical_date (formerly execution_date): Represents the beginning of the data interval the DAG run is processing, NOT when the run physically started. For a daily DAG covering October 5, the logical date is 2026-10-05 00:00, but it physically runs on 2026-10-06 00:00 after the entire day's data has closed.
• start_date: The initial timestamp from which the scheduler begins creating DAG runs.
Writing idempotent pipelines requires querying partitions based on logical_date rather than wall-clock time datetime.now().

4. Related Tracks

Explore Next

Explore the Master Education & Interview Academy

Access structured learning paths, scenario-based interview guides, and interactive mock interview rubrics.

🎓 Explore Education Track →
Open Orchestration interview questions →