Apache Airflow Data Engineering Pipelines
Production-grade Apache Airflow DAGs for orchestrating healthcare data engineering pipelines. Covers the full spectrum — incremental ingestion, dynamic DAG generation, dbt model triggers, Databricks job orchestration, and cross-system data quality checks — patterns built from real production experience.
┌─────────────────────────────────────────────────────────────┐
│ APACHE AIRFLOW (2.8+) │
│ │
│ ┌─────────────┐ ┌──────────────┐ ┌───────────────────┐ │
│ │ Schedules │ │ DAG Factory │ │ Task Groups │ │
│ │ & Triggers │ │ (Dynamic) │ │ & Dependencies │ │
│ └──────┬──────┘ └──────┬───────┘ └─────────┬─────────┘ │
└─────────┼────────────────┼────────────────────┼────────────┘
│ │ │
┌─────▼──────┐ ┌──────▼───────┐ ┌────────▼────────┐
│ ADLS Gen2 │ │ Databricks │ │ dbt Cloud / │
│ Ingestion │ │ Jobs API │ │ dbt Core │
└─────┬──────┘ └──────┬───────┘ └────────┬────────┘
│ │ │
└────────────────┴────────────────────┘
│
┌──────▼──────┐
│ Azure SQL │
│ Synapse / │
│ Snowflake │
└─────────────┘
apache-airflow-data-pipelines/
│
├── dags/
│ ├── dag_incremental_hl7_ingestion.py # HL7/FHIR healthcare data ingestion
│ ├── dag_patient_etl_orchestrator.py # Master orchestration DAG
│ ├── dag_dbt_daily_run.py # dbt model trigger DAG
│ ├── dag_databricks_job_trigger.py # Databricks job orchestration
│ ├── dag_dynamic_table_loader.py # Dynamic DAG factory pattern
│ └── dag_cross_system_dq_check.py # Cross-system data quality
│
├── plugins/
│ ├── hooks/
│ │ ├── adls_hook.py # Custom ADLS Gen2 hook
│ │ └── databricks_rest_hook.py # Databricks REST API hook
│ ├── operators/
│ │ ├── adls_to_delta_operator.py # ADLS → Delta Lake operator
│ │ └── dq_validation_operator.py # Data quality operator
│ └── sensors/
│ └── adls_file_sensor.py # ADLS file arrival sensor
│
├── docker/
│ └── docker-compose.yaml # Airflow local dev setup
│
├── config/
│ └── airflow_variables.json # Airflow Variables template
│
└── tests/
└── test_dag_integrity.py # DAG integrity tests
Pattern
DAG
Description
Incremental load
dag_incremental_hl7_ingestion.py
Watermark-based ingestion with XCom state passing
Dynamic DAGs
dag_dynamic_table_loader.py
DAG factory generating one DAG per source table
dbt orchestration
dag_dbt_daily_run.py
Run dbt models by tag with retry and alerting
Databricks trigger
dag_databricks_job_trigger.py
Trigger Databricks jobs via REST API, poll for completion
Cross-system DQ
dag_cross_system_dq_check.py
Row count reconciliation across source and target
Custom operator
plugins/operators/
Reusable operators for ADLS and DQ
Custom sensor
plugins/sensors/
File arrival sensor with timeout
git clone https://github.com/donthula9908/apache-airflow-data-pipelines.git
cd apache-airflow-data-pipelines
# Start Airflow with Docker
docker compose -f docker/docker-compose.yaml up -d
# Access UI at http://localhost:8080 (admin/admin)
Component
Technology
Orchestration
Apache Airflow 2.8+
Storage
Azure Data Lake Storage Gen2
Processing
Databricks (PySpark)
Transformation
dbt Core
Containerisation
Docker Compose
Testing
pytest + Airflow TestUtils