Repository to version the DAGs used in Airflow2 in pvpdstata02 | Read-only mirror of https://github.com/opendatabs/dags-airflow2 — Kanton Basel-Stadt. Issues & pull requests at the source.
Find a file
Repository files (latest commit first)
Filename Latest commit message Latest commit date
2026-08-18 17:03:08 +02:00
dataspot-connector Add Dataspot connector DAG for StatA DWH 2026-06-05 14:55:35 +02:00
helpers Fix FailureTrackingDockerOperator. This will now allow multiple FailureTrackingDockerOperators per DAG without interference 2025-07-11 00:47:47 +02:00
.gitignore Exclude folders with databases from the dataspot connectors in .gitignore 2025-11-12 17:36:17 +01:00
airflow_db_cleanup.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
airflow_export_env_file.py Update airflow_export_env_file.py 2025-08-08 15:51:25 +02:00
aue_fischereistatistik.py Rename module to aue_fischereistatistik 2025-11-29 12:41:03 +01:00
aue_grundwasser.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
aue_luft_klybeck.py Add API_KEY_DECENTLAB to private environment variables in aue_luft_klybeck.py 2026-08-17 17:41:11 +02:00
aue_rues.py Add cleanup containers task to multiple DAGs using BashOperator 2026-06-08 08:55:28 +02:00
aue_schall.py Update aue_schall to send error on third failure and auto-abort job if it takes too long 2026-04-09 14:31:57 +02:00
aue_umweltlabor.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
awa_bewilligungen.py Add new DAG for AWA Bewilligungen to update datasets with DockerOperator 2026-02-11 15:38:50 +01:00
awa_feiertage.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
bafu_hydro_daten.py Add cleanup containers task to multiple DAGs using BashOperator 2026-06-08 08:55:28 +02:00
bafu_hydro_daten_vorhersagen.py Add cleanup containers task to multiple DAGs using BashOperator 2026-06-08 08:55:28 +02:00
bvb_fahrgastzahlen.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
common_variables.py Add new DAG for FGI STAC dataset processing and include HUWISE_API_KEY in common environment variables 2026-04-30 10:37:59 +02:00
datasette_rsync.py Update datasette_rsync DAG schedule to run daily at midnight 2025-11-19 15:12:09 +01:00
dcc_dataspot_catalog_quality_daily.py Add env variable 2025-12-12 16:45:24 +01:00
dcc_dataspot_connector_aue_brytecube_prod.py Update driver for DCC Dataspot Connectors to 2026.2.2 2026-08-10 14:04:30 +02:00
dcc_dataspot_connector_bvd_stg_p2stg.py Update driver for DCC Dataspot Connectors to 2026.2.2 2026-08-10 14:04:30 +02:00
dcc_dataspot_connector_ed_escada2.py Update driver for DCC Dataspot Connectors to 2026.2.2 2026-08-10 14:04:30 +02:00
dcc_dataspot_connector_fd_itbs_kdm.py Update driver for DCC Dataspot Connectors to 2026.2.2 2026-08-10 14:04:30 +02:00
dcc_dataspot_connector_stata_ad_test.py Update driver for DCC Dataspot Connectors to 2026.2.2 2026-08-10 14:04:30 +02:00
dcc_dataspot_connector_stata_dwh.py Update driver for DCC Dataspot Connectors to 2026.2.2 2026-08-10 14:04:30 +02:00
dcc_dataspot_connector_stata_personen_dwh_test.py Update driver for DCC Dataspot Connectors to 2026.2.2 2026-08-10 14:04:30 +02:00
dcc_dataspot_connector_stata_test.py Update driver for DCC Dataspot Connectors to 2026.2.2 2026-08-10 14:04:30 +02:00
dcc_dataspot_daily_jobs.py Add tasks to sync Swiss law and Basel-Stadt law collections into dataspot 2026-03-26 23:40:15 +01:00
dcc_dataspot_schemes.py Update environment variables from username-password authentication to service user authentication 2025-09-18 11:20:59 +02:00
dcc_datenkatalog_dienststellen_onboarding.py Add new DAG for onboarding dcc_datenkatalog_dienststellen with FailureTrackingDockerOperator 2026-08-10 11:06:38 +02:00
dcc_ki_faq.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
dcc_verzeichnis_personendaten.py Add airflow job for dcc_verzeichnis_personendaten 2026-06-09 14:27:12 +02:00
ed_schulferien.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
ed_swisslos_sportfonds.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
esc_faq.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
euroairport.py Fix: Re-add DockerOperator import 2026-04-10 09:21:42 +02:00
fgi_geodatenshop.py Add Certs 2026-08-11 11:41:51 +02:00
fgi_stac.py Add update certificates here too 2026-08-18 17:03:08 +02:00
gd_abwassermonitoring.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
gd_coronavirus_abwassermonitoring.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
gva_geodatenshop.py Add it everywhere where we have connection to the geo pages 2026-08-11 11:45:26 +02:00
gva_metadata.py Increase failure threshold in DAG configuration from 1 to 4 2025-11-10 09:20:15 +01:00
ibs_parkhaus_bewegungen.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
itbs_klv.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
iwb_gas.py Add new mount for change tracking in iwb_gas.py 2025-09-18 17:17:10 +02:00
iwb_netzlast.py Set user configuration to 'root' in DockerOperator across multiple DAG files 2025-09-18 17:05:32 +02:00
jfs_gartenbaeder.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
kantonslabor_coliminder.py Enhance kantonslabor DAG by adding API credentials to private environment variables 2026-05-08 17:05:50 +02:00
kapo_eventverkehr_stjakob.py Change to Client ID 2026-08-07 10:25:03 +02:00
kapo_geschwindigkeitsmonitoring.py Add rsync DockerOperator to multiple DAGs for file synchronization 2025-11-26 22:33:07 +01:00
kapo_ordnungsbussen.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
kapo_smileys.py Add cleanup containers task to multiple DAGs using BashOperator 2026-06-08 08:55:28 +02:00
LICENSE Initial commit 2024-01-12 09:36:01 +01:00
lufthygiene_rosental.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
luftqualitaet_ch.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
meteoblue_rosental.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
meteoblue_wolf.py Add cleanup containers task to multiple DAGs using BashOperator 2026-06-08 08:55:28 +02:00
meteoschweiz_station_basel_binningen.py Add new DAG for MeteoSwiss Basel-Binningen 2026-02-20 20:13:45 +01:00
mkb_sammlung_europa.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
mobilitaet_dtv.py Do it one hour later since kapo_geschwindigkeitsmonitoring takes mor time now 2025-08-13 12:02:21 +02:00
mobilitaet_mikromobilitaet.py Add it everywhere where we have connection to the geo pages 2026-08-11 11:45:26 +02:00
mobilitaet_mikromobilitaet_stats.py Add rsync as process 2026-06-03 15:29:30 +02:00
mobilitaet_parkflaechen.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
mobilitaet_verkehrszaehldaten.py Add rsync DockerOperator to multiple DAGs for file synchronization 2025-11-26 22:33:07 +01:00
ods_catalog.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
ods_update_temporal_coverage.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
parkendd.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
parkhaeuser.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
parlamentsdienst_gr_abstimmungen.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
parlamentsdienst_grosserrat.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
parlamentsdienst_grosserrat_datasette.py Add it everywhere where we have connection to the geo pages 2026-08-11 11:45:26 +02:00
pyproject.toml Add uv with ruff 2025-04-28 09:18:38 +02:00
README.md Enhance README.md with comprehensive documentation on DAG types, usage guidelines, and related repositories for Airflow2 workflows. 2026-01-22 15:01:03 +01:00
renovate.json Add renovate.json 2025-05-14 13:26:38 +00:00
smarte_strasse_ladestation.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
stadtgaertnerei_spielen.py Add new DAG for stadtgaertnerei_spielen to update datasets using DockerOperator 2025-09-26 17:37:51 +02:00
stadtreinigung_sauberkeitsindex.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
stadtreinigung_wildedeponien.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
staka_abstimmungen.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
staka_briefliche_stimmabgaben.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
staka_gutachten.py Change schedule to less, since it now makes a connection to Sharepoint everytime 2026-03-19 16:46:12 +01:00
staka_kandidaturen.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
staka_kantonsblatt.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
staka_regierungsratsbeschluesse.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
staka_staatskalender.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
staka_vernehmlassungen.py Add staka_vernehmlassungen 2026-01-30 11:55:57 +01:00
staka_verz_verf_persdat.py Deprecate VVP 2026-06-09 11:50:25 +02:00
stata_baselvotes.py Add new DAG for populating baselvotes dataset with FailureTrackingDockerOperator 2026-08-04 15:19:21 +02:00
stata_befragungen.py Add new mounts for Bevoelkerungsbefragung data sources in stata_befragungen.py 2026-03-26 13:34:03 +01:00
stata_bik.py Update stata_bik to send error on third failure and auto-abort job if it takes too long 2026-04-09 14:55:32 +02:00
stata_daily_upload.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
stata_gwr.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
stata_harvester.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
stata_hunde.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
stata_konoer.py Add Requisitionen and Allmend 2026-01-14 16:04:11 +01:00
stata_parzellen.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
stata_pull_changes.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
stata_requisitionen.py Add the argument too 2026-01-14 16:19:15 +01:00
stata_superblock_allmend.py Add the argument too 2026-01-14 16:19:15 +01:00
stata_tourismusdashboard.py Update DockerOperator to use GitHub Container Registry and enable force pull for rsync images 2025-09-22 12:26:08 +02:00
tba_abfuhrtermine.py Update tba_abfuhrtermine.py 2025-12-17 10:06:49 +01:00
tba_baustellen.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
tba_sprayereien.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
tba_wiese.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
unibas_semesterdaten.py Add "unibas_semesterdaten" DAG 2025-09-19 12:01:42 +02:00
uv.lock Add uv with ruff 2025-04-28 09:18:38 +02:00
zefix_handelsregister.py Try adding it again for all of them 2025-08-08 15:00:16 +02:00
zrd_gesetzessammlung.py Update cron schedule in zrd_gesetzessammlung.py 2026-06-09 11:58:02 +02:00

dags-airflow2

Repository to version the DAGs used in Airflow2 in pvpdstata02.

This repository is part of the opendatabs organization. For more information about the organization and its projects, see the opendatabs organization README.

Most of the Docker images referenced in these DAGs are built and maintained in the data-processing repository. For details about how these images are created, their structure, and how to contribute, refer to the data-processing README.

Overview

This repository contains Apache Airflow 2.7.2 DAG definitions that orchestrate data processing workflows. The DAGs primarily use Docker containers to execute ETL (Extract, Transform, Load) processes, with most containers hosted at ghcr.io/opendatabs/data-processing/.

DAG Types

Open Data DAGs

Most DAGs in this repository are for Open Data workflows. These DAGs:

  • Process and publish datasets to the Open Data Portal (data.bs.ch)
  • Must include dataset IDs in the docstring at the top of the file (visible in Airflow UI)
  • Follow the naming convention that matches the folder name in the data-processing repository
  • Use Docker images from ghcr.io/opendatabs/data-processing/{dag_id}:latest

Dataset IDs in Docstrings

Open Data DAGs must document which datasets they affect in the docstring. This information is displayed in the Airflow UI and helps identify affected datasets when errors occur.

Example:

"""
# aue_umweltlabor
This DAG updates the following datasets:

- [100066](https://data.bs.ch/explore/dataset/100066)
- [100067](https://data.bs.ch/explore/dataset/100067)
- [100068](https://data.bs.ch/explore/dataset/100068)
- [100069](https://data.bs.ch/explore/dataset/100069)
"""

The docstring is made visible in the Airflow UI by setting dag.doc_md = __doc__ in the DAG definition.

Naming Convention

Important: The DAG dag_id must match the folder name in the data-processing repository. This ensures:

  • Consistent naming across repositories
  • Correct Docker image references (ghcr.io/opendatabs/data-processing/{dag_id}:latest)
  • Proper mount paths ({PATH_TO_CODE}/data-processing/{dag_id}/...)

For example:

  • DAG file: aue_umweltlabor.py with dag_id="aue_umweltlabor"
  • Corresponding folder in data-processing: data-processing/aue_umweltlabor/
  • Docker image: ghcr.io/opendatabs/data-processing/aue_umweltlabor:latest

Data Catalog DAGs (dcc_)

DAGs starting with dcc_dataspot are mostly for the Data Catalog (Dataspot), not Open Data. These DAGs:

  • Sync metadata and organizational structures to the data catalog
  • May use different Docker image registries (e.g., ghcr.io/dcc-bs/dataspot:latest)
  • Do not publish to the Open Data Portal
  • Examples: dcc_dataspot_daily_jobs, dcc_dataspot_connector_*, dcc_dataspot_catalog_quality_daily

Special Cases: Other Repositories

Some DAGs use Docker images from repositories other than data-processing:

  • rsync: Used for syncing files to remote servers

    • Image: ghcr.io/opendatabs/rsync:latest
    • Used in: stata_konoer, stata_tourismusdashboard, kapo_smileys, kapo_geschwindigkeitsmonitoring, etc.
  • stata_konoer: R-based data processing

    • Image: ghcr.io/opendatabs/stata_konoer:latest
    • Used in: stata_konoer.py
  • tourismusdashboard: Tourism dashboard data processing

    • Image: ghcr.io/opendatabs/tourismusdashboard:latest
    • Used in: stata_tourismusdashboard.py

These special cases are documented in the individual DAG files.

Helpers

FailureTrackingDockerOperator

Located in helpers/failure_tracking_operator.py, this custom operator extends the standard DockerOperator to include failure tracking capabilities.

When to Use

Use FailureTrackingDockerOperator when you want to tolerate a certain number of consecutive failures before actually failing a task. This is useful for:

  • Intermittent data source issues
  • Network connectivity problems
  • Temporary API outages
  • Frequently scheduled DAGs where occasional failures are expected
  • Any scenario where occasional failures are expected but shouldn't immediately fail the DAG

How to Use

The typical pattern is to define configuration constants at the top of your DAG file:

from helpers.failure_tracking_operator import FailureTrackingDockerOperator
from datetime import datetime, timedelta
from airflow import DAG
from airflow.models import Variable
from docker.types import Mount
from common_variables import COMMON_ENV_VARS, PATH_TO_CODE

# DAG configuration
DAG_ID = "my_dag"
FAILURE_THRESHOLD = 1  # Skip first failure, fail on second
EXECUTION_TIMEOUT = timedelta(minutes=3)
SCHEDULE = "*/5 * * * *"

default_args = {
    "owner": "your.name",
    "depends_on_past": False,
    "start_date": datetime(2024, 1, 1),
    "email": Variable.get("EMAIL_RECEIVERS"),
    "email_on_failure": True,
    "email_on_retry": False,
    "retries": 0,
    "retry_delay": timedelta(minutes=15),
}

with DAG(
    dag_id=DAG_ID,
    default_args=default_args,
    description=f"Run the {DAG_ID} docker container",
    schedule=SCHEDULE,
    catchup=False,
) as dag:
    dag.doc_md = __doc__
    
    upload = FailureTrackingDockerOperator(
        task_id="upload",
        failure_threshold=FAILURE_THRESHOLD,
        execution_timeout=EXECUTION_TIMEOUT,
        image=f"ghcr.io/opendatabs/data-processing/{DAG_ID}:latest",
        force_pull=True,
        api_version="auto",
        auto_remove="force",
        mount_tmp_dir=False,
        command="uv run -m etl",
        private_environment=COMMON_ENV_VARS,
        container_name=DAG_ID,
        docker_url="unix://var/run/docker.sock",
        network_mode="bridge",
        tty=True,
        mounts=[
            Mount(
                source=f"{PATH_TO_CODE}/data-processing/{DAG_ID}/data",
                target="/code/data",
                type="bind",
            ),
            Mount(
                source=f"{PATH_TO_CODE}/data-processing/{DAG_ID}/change_tracking",
                target="/code/change_tracking",
                type="bind",
            ),
        ],
    )

Parameters

  • failure_threshold (int, required): Number of consecutive failures to tolerate before actually failing the task

    • 0 = immediate failure with no skipping (same as regular DockerOperator)
    • 5 = skip first 5 failures, fail on the 6th consecutive failure
    • Common values in the codebase: 1, 2, 4, 5, 16 (depending on DAG frequency and reliability requirements)
  • execution_timeout (timedelta or None, required): Maximum time allowed for task execution. Task fails if exceeded.

    • Use timedelta(minutes=30) for a 30-minute timeout
    • Use timedelta(seconds=50) for a 50-second timeout
    • Use None for no timeout (not recommended)
    • Common values in the codebase: timedelta(minutes=2) to timedelta(minutes=90)

Note: Both failure_threshold and execution_timeout are required parameters and must be provided when using FailureTrackingDockerOperator. All other parameters are the same as DockerOperator.

How It Works

The operator tracks consecutive failures using Airflow Variables (stored as {dag_id}_{task_id}_consecutive_failures).

  • On each failure, the count increments
  • If the count exceeds the threshold, the task fails and raises an exception
  • If the count is under the threshold, the task is skipped (raises AirflowSkipException)
  • On success, the counter resets to 0

This allows DAGs to continue running even when individual task executions fail, as long as the failure count stays below the threshold. This is particularly useful for frequently scheduled DAGs where occasional failures are expected.

Docker Container Cleanup

Some DAGs include a cleanup task at the beginning to remove any old containers that may have been left behind from previous runs. This prevents container name conflicts and ensures clean execution.

When to Use

Add a cleanup task when:

  • Your DAG uses a fixed container_name (not auto-generated)
  • Previous failed runs might have left containers running
  • You're experiencing container name conflicts

How to Implement

from airflow.operators.bash import BashOperator

# Cleanup task to remove any old containers at the beginning
cleanup_containers = BashOperator(
    task_id="cleanup_old_containers",
    bash_command=f'''
        docker rm -f {DAG_ID} 2>/dev/null || true
        ''',
)

# Set the task dependency
cleanup_containers >> your_main_task

The 2>/dev/null || true ensures the command doesn't fail if the container doesn't exist.

Common Variables

The common_variables.py module provides shared configuration:

  • COMMON_ENV_VARS: Dictionary of common environment variables (proxy settings, email configuration, FTP credentials, ODS API keys)
  • PATH_TO_CODE: Path to the code directory (set via Airflow Variable)

These are typically used via:

from common_variables import COMMON_ENV_VARS, PATH_TO_CODE

DAG Configuration Fields

DAG Object

The DAG object is the main container for your workflow. Reference: Airflow 2.7.2 DAG Documentation

Common Fields

  • dag_id (str, required): Unique identifier for the DAG. Must be unique across all DAGs.

  • description (str, optional): Description of what the DAG does. Displayed in the Airflow UI.

  • default_args (dict, optional): Default arguments applied to all tasks in the DAG. Can be overridden at the task level.

  • schedule (str, timedelta, or None, optional): Schedule for the DAG execution.

    • Cron expression: "0 6 * * *" (daily at 6 AM)
    • None for manually triggered DAGs
    • Reference: Scheduling
  • catchup (bool, optional): Whether to backfill missed DAG runs. Set to False to prevent catchup.

default_args Dictionary

Common fields used in default_args:

  • owner (str): Owner of the DAG (typically an email or username)

  • depends_on_past (bool): Whether a task depends on the success of the previous task instance

  • start_date (datetime): The first execution date for the DAG

  • email (str or list): Email address(es) to send notifications to

  • email_on_failure (bool): Send email on task failure

  • email_on_retry (bool): Send email on task retry

  • retries (int): Number of retries for failed tasks

  • retry_delay (timedelta): Delay between retries

DockerOperator Fields

The DockerOperator runs Docker containers as Airflow tasks. Reference: Docker Operator Documentation

Common DockerOperator Fields

  • task_id (str, required): Unique identifier for the task within the DAG

  • image (str, required): Docker image to use (e.g., "ghcr.io/opendatabs/data-processing/my_dag:latest")

  • command (str or list, optional): Command to run in the container

    • Can be a string: "uv run -m etl"
    • Or a list: ["java", "-jar", "/app/app.jar"]
    • Reference: Command
  • container_name (str, optional): Name for the Docker container. If not provided, a name is auto-generated.

  • force_pull (bool, optional): Whether to pull the image even if it already exists locally. Set to True to always get the latest version.

  • api_version (str, optional): Docker API version. Use "auto" to auto-detect.

  • auto_remove (str, optional): Whether to remove the container after execution.

    • "force" = always remove, even on failure
    • True = remove on success
    • False = never remove
    • Reference: Auto Remove
  • mount_tmp_dir (bool, optional): Whether to mount a temporary directory. Set to False to disable.

  • docker_url (str, optional): URL to the Docker daemon. Use "unix://var/run/docker.sock" for local Docker.

  • network_mode (str, optional): Docker network mode (e.g., "bridge", "host").

  • tty (bool, optional): Whether to allocate a pseudo-TTY. Set to True for better log output.

    • Reference: TTY
  • private_environment (dict, optional): Environment variables to pass to the container (sensitive values are hidden in logs).

  • environment (dict, optional): Environment variables to pass to the container (visible in logs).

  • mounts (list, optional): List of docker.types.Mount objects to mount volumes into the container.

    • Example:

      from docker.types import Mount
      
      mounts=[
          Mount(
              source="/host/path",
              target="/container/path",
              type="bind",
          ),
      ]
      
    • Reference: Mounts

  • working_dir (str, optional): Working directory inside the container.

  • execution_timeout (timedelta, optional): Maximum time allowed for task execution. Task fails if exceeded.

Example DAG Structure

Open Data DAG Example

"""
# my_dag
This DAG updates the following datasets:

- [100001](https://data.bs.ch/explore/dataset/100001)
- [100002](https://data.bs.ch/explore/dataset/100002)
"""

from datetime import datetime, timedelta
from airflow import DAG
from airflow.models import Variable
from airflow.providers.docker.operators.docker import DockerOperator
from docker.types import Mount
from common_variables import COMMON_ENV_VARS, PATH_TO_CODE

DAG_ID = "my_dag"  # Must match folder name in data-processing repository

default_args = {
    "owner": "your.name",
    "depends_on_past": False,
    "start_date": datetime(2024, 1, 1),
    "email": Variable.get("EMAIL_RECEIVERS"),
    "email_on_failure": True,
    "email_on_retry": False,
    "retries": 0,
    "retry_delay": timedelta(minutes=15),
}

with DAG(
    dag_id=DAG_ID,
    default_args=default_args,
    description=f"Run the {DAG_ID} docker container",
    schedule="0 6 * * *",  # Daily at 6 AM
    catchup=False,
) as dag:
    dag.doc_md = __doc__  # Makes the docstring visible in Airflow UI
    
    process_task = DockerOperator(
        task_id="process",
        image=f"ghcr.io/opendatabs/data-processing/{DAG_ID}:latest",
        force_pull=True,
        api_version="auto",
        auto_remove="force",
        mount_tmp_dir=False,
        command="uv run -m etl",
        private_environment=COMMON_ENV_VARS,
        container_name=DAG_ID,
        docker_url="unix://var/run/docker.sock",
        network_mode="bridge",
        tty=True,
        mounts=[
            Mount(
                source=f"{PATH_TO_CODE}/data-processing/{DAG_ID}/data",
                target="/code/data",
                type="bind",
            ),
            Mount(
                source=f"{PATH_TO_CODE}/data-processing/{DAG_ID}/change_tracking",
                target="/code/change_tracking",
                type="bind",
            ),
        ],
    )

Data Catalog DAG Example (dcc_)

"""
# dcc_dataspot_sync_example
This DAG syncs data catalog metadata.
"""

from datetime import datetime, timedelta
from airflow import DAG
from airflow.models import Variable
from airflow.providers.docker.operators.docker import DockerOperator
from common_variables import COMMON_ENV_VARS

default_args = {
    "owner": "your.name",
    "depends_on_past": False,
    "start_date": datetime(2024, 1, 1),
    "email": Variable.get("EMAIL_RECEIVERS"),
    "email_on_failure": True,
    "email_on_retry": False,
    "retries": 0,
    "retry_delay": timedelta(minutes=15),
}

with DAG(
    "dcc_dataspot_sync_example",
    default_args=default_args,
    description="Sync data catalog metadata",
    schedule="0 3 * * *",
    catchup=False,
) as dag:
    dag.doc_md = __doc__
    
    sync_task = DockerOperator(
        task_id="sync",
        image="ghcr.io/dcc-bs/dataspot:latest",
        force_pull=True,
        api_version="auto",
        auto_remove="force",
        mount_tmp_dir=False,
        command="python -m scripts.sync_example",
        private_environment=COMMON_ENV_VARS,
        container_name="dcc_dataspot_sync_example",
        docker_url="unix://var/run/docker.sock",
        network_mode="bridge",
        tty=True,
    )

Additional Resources