---
name: create-dlt-docker-operator-dag
description: Create an Airflow DAG that runs containerized dlt pipelines using DockerOperator. Use this when asked to create a Docker container task, a DockerOperator DAG, or a containerized dlt Airflow DAG.
disable-model-invocation: true
---

# Create a DockerOperator DAG

Create an Airflow DAG that runs containerized workloads (such as dlt connectors or dbt) using `DockerOperator`. This pattern decouples pipeline logic from Airflow, running each task in its own isolated container.

## Context

Instead of installing pipeline dependencies directly in the Airflow environment, `DockerOperator` runs each task as a standalone Docker container. This requires:

- A **docker-proxy** service in `docker-compose.yml` (see the `setup-dlt-docker-operator-airflow` skill)
- A **shared utilities module** (`dags/utils/common.py`) for reading secrets from Airflow Variables
- Pre-built Docker images for each workload (e.g., dlt connectors, dbt)

The DAG reads secrets from Airflow Variables (not from the filesystem) and passes them to containers via environment variables.

## Prerequisites

- Airflow setup with DockerOperator support (see `setup-dlt-docker-operator-airflow` skill)
- Docker images built and available locally or in a registry
- Airflow Variables configured:
  - `dlt_destination` - Target destination (e.g., `postgres`, `snowflake`)
  - `dlt_secrets_toml` - Full TOML content of `secrets.toml` (stored as a Variable, base64-encoded at runtime)

## Steps

1. Ensure the shared utilities module exists at `dags/utils/common.py`. Use the template from `assets/utils_common.py` if it doesn't.

   > **Compatibility warning:** `assets/utils_common.py` uses `tomllib` (Python 3.11+). If the Airflow environment runs Python 3.8–3.10 (e.g., EWAH), replace `tomllib` with a `tomli` backport or remove the TOML validation step. For Python 3.8-compatible setups, prefer the uv connector approach instead (see `create-dlt-uv-connector-dag` skill).

2. Create a new DAG file at `dags/dag_extract_load__<connector_name>.py`.

3. Import the required modules and shared utilities.

   ```python
   import pendulum

   from airflow.sdk import dag, task
   from airflow.providers.docker.operators.docker import DockerOperator

   from utils.common import get_dlt_destination, get_dlt_secrets_toml_base64
   ```

4. Define the DAG using the `@dag` decorator.

   ```python
   @dag(
       schedule=None,  # Set schedule as needed (e.g., "@daily")
       start_date=pendulum.datetime(2025, 1, 1, tz="UTC"),
       catchup=False,
       tags=["extract_load", "dlt"],
   )
   def extract_load_<connector_name>():
       ...
   ```

5. Create a task that runs the dlt connector via DockerOperator.

   ```python
   @task()
   def run_dlt_pipeline(destination: str, secrets_toml_base64: str):
       return DockerOperator(
           task_id="run_dlt_pipeline",
           image="dlt-connector-<connector_name>:latest",
           container_name="dlt-<connector_name>",
           api_version="auto",
           auto_remove="force",
           docker_url="tcp://docker-proxy:2375",
           environment={
               "DLT_DESTINATION": destination,
               "DLT_SOURCE_NAME": "<source_name>",
               "DLT_SECRETS_TOML_BASE64": secrets_toml_base64,
           },
       ).execute(context={})
   ```

6. Wire up the DAG flow by reading shared secrets and passing them to tasks.

   ```python
       secrets_toml_base64 = get_dlt_secrets_toml_base64()
       destination = get_dlt_destination()
       run_dlt_pipeline(destination, secrets_toml_base64)

   dag = extract_load_<connector_name>()
   ```

7. For connectors with multiple instances (e.g., multiple ad accounts), create separate tasks for each.

   ```python
   @task()
   def load_account_germany(destination: str, secrets_toml_base64: str):
       return DockerOperator(
           task_id="account_germany",
           image="dlt-connector-google_ads:latest",
           container_name="dlt-google-ads-germany",
           api_version="auto",
           auto_remove="force",
           docker_url="tcp://docker-proxy:2375",
           environment={
               "DLT_DESTINATION": destination,
               "DLT_SOURCE_NAME": "google_ads_germany",
               "DLT_SECRETS_TOML_BASE64": secrets_toml_base64,
           },
       ).execute(context={})
   ```

8. For passing additional configuration (e.g., resource configs), serialize to JSON and pass via environment variable.

   ```python
   import json

   RESOURCES_CONFIG = [
       {"table": "products", "schema": "public", "incremental": "modified", "primary_key": "uuid"},
   ]

   # In the DockerOperator environment:
   environment={
       ...,
       "DLT_RESOURCES_CONFIG": json.dumps(RESOURCES_CONFIG),
   }
   ```

## Validation

- [ ] DAG appears in Airflow UI without import errors
- [ ] Tasks run successfully and containers complete
- [ ] Data loads to the expected destination
- [ ] Secrets are not exposed in Airflow logs (use `show_return_value_in_logs=False` for sensitive tasks)

## Examples

**Key DockerOperator parameters:**

| Parameter | Value | Purpose |
|-----------|-------|---------|
| `docker_url` | `tcp://docker-proxy:2375` | Connect via the socat proxy service |
| `api_version` | `"auto"` | Auto-detect Docker API version |
| `auto_remove` | `"force"` | Always clean up containers |
| `container_name` | Unique per task | Avoid container name conflicts |

**DAG naming convention:** `dag_extract_load__<connector_name>.py` for EL pipelines, `dag_transform_dbt.py` for transformations.
