Status

StateAccepted
Discussion Thread
Vote Thread[LAZY,CONSENSUS] Example dags
Vote Result ThreadRe: [LAZY,CONSENSUS] Example dags
Progress Tracking (PR/GitHub Project/Issue Label)
Date Created

 

Version Releasedtbd.
Authors

Motivation

Airflow Examples have been grown in number and focus over the past years. They purpose multiple things:

Some example DAGs are in a good quality, some are not following best practices. Current examples do not follow a structure.

There are example DAGs contained in the Airflow core (currently pushed to standard provider/example_dags) as well as there are more examples in other providers. But examples from other providers are lot loaded automatically.

So in the Airflow 3 Dev Calls there was a demand named to clean-up and optimize example DAGs.

Considerations / Targets

Storyline

Current idea:

As for demos and examples a lot of functionality is needed in both decorator as well as classic Dag implementation it would be good to have two similar use cases. Or alternatively describe that the Tailwind south branch prefers to implement all in Pythonic manner whereas the Tailwind North branch data engineers like the classic implementation?

digraph D {
    rankdir = LR;
    
    subgraph cluster_p {
        label = "Dag: Turbine Report";
        
        ev [shape = folder;label = "File\nEvent";];
        report [xlabel = "with failure / retry";];
        ingest;
        
        ev -> ingest -> aggregate -> report;
    }
    
    subgraph cluster_c {
        label = "Dag: City Reporting";
        
        split [label = "split by city";];
        
        subgraph clusterc2 {
            label = "Dynamic Mapped Tasks";
            
            seattle;
            portland;
            vancouver;
        }
        split -> seattle -> summary;
        split -> portland -> summary;
        split -> vancouver -> summary;
    }
    
    ingest -> split [color = "#aaaaaa"; constraint = false; label = "asset event";];
    
    subgraph cluster_m {
        label = "Dag: Measurement correction";
        
        form [label="Trigger form w/\ncity\nmeasure", shape=component]
        notify [shape=diamond]

        form -> file -> notify

        notify -> email -> completion
        notify -> push -> completion
        notify -> none -> completion
    }

    file -> ev [color = "#aaaaaa"; constraint = false; label = "asset event"]

    subgraph cluster_f {
        label = "Dag: Turbine Monitoring";
        
        cron [shape = doublecircle; label = "cron:hourly"]
        inventory [label="get inventory"]

        cron -> inventory

        subgraph cluster_f1 {
            label = "Dynamic Mapped Task Group";

            vpn_on [label="vpn connect\n(setup)"]
            vpn_off [label="vpn disconnect\n(teardown)"]

            health [shape=diamond]

            vpn_on -> get_telemetry -> health -> vpn_off
            vpn_on -> check_alarms -> health

            health -> send_technicial [label="notify if alarm"]
            health -> all_ok
        }

        inventory -> vpn_on
    }

    subgraph cluster_s {
        label = "Dag: Maintenance";
        
        timetable [shape = doublecircle; label = "maintenance\ntimetable"]
        form2 [label="Trigger form w/\nturbine\ntime\nduration", shape=component]

        maintenance_on [label="maintenance start\nbash script"]
        wait [label="wait for maintenance completion\ndeferred"]
        maintenance_off [label="maintenance end\nbash script"]

        timetable -> maintenance_on [label="planned schedule"]
        form2  -> maintenance_on [label="manual trigger"]

        maintenance_on -> wait -> maintenance_off
   }

    subgraph cluster_r2 {
        label = "Dag: Shareholder Reporting\nDynamic Dag Generation";

        cron2 [shape = doublecircle; label = "cron:monthly"]

        cron2 -> load_data

        subgraph cluster_g {
            label = "Generated by loop";

            subgraph cluster_ga {
                label="";
                color=none;

                report_stakeholder_a -> send_report_a
            }

            subgraph cluster_gb {
                label="";
                color=none;

                report_stakeholder_b -> send_report_b
            }

            subgraph cluster_gc {
                label="";
                color=none;

                report_stakeholder_c -> send_report_c
            }
        }

        load_data -> report_stakeholder_a
        load_data -> report_stakeholder_b
        load_data -> report_stakeholder_c
    }
}

Technical Work Packages → Features

Proposed Technical Excellence in Example DAGs

Current Examples → Example DAG Target Panning

See https://github.com/apache/airflow/tree/main/airflow-core/src/airflow/example_dags

NameGapsReferenced inFeatures usedProposed Change

Core - airflow-core/src/airflow/example_dags

example_asset_alias.py


-AssetAlias, taskflow

Documentation in https://airflow.apache.org/docs/apache-airflow/stable/authoring-and-scheduling/assets.html#dynamic-data-events-emitting-and-asset-creation-through-assetalias does not use the examples. Either ned to add to real-world example, refernce it in docs or delete them. Unknown User (uranusjr) Do you have an ida how to map this to a real world example (question)

example_asset_alias_with_no_taskflow.py


-AssetAlias

example_asset_decorator.py


-assets, decoratorsBasic examples should be integrated into the storyline. Else too many dags w/o business example. Then drop (error)

example_assets.py


-assets

example_asset_with_watchers.py


-AssetWatcherNeeds to be intgrated into case, not standalone. Then drop (error)

example_branch_labels.py


airflow-core/docs/core-concepts/dags.rst
Needs to be intgrated into case, not standalone. Then drop (error)

example_branch_python_dop_operator_3.py


-
Needs to be intgrated into case, not standalone, then drop (error)

example_complex.py


airflow-core/docs/howto/usage-cli.rst
Lags a real business case. But complexity might be still a good show case. Check for the resulting example, if similar complexity then drop (error)

example_custom_weight.py


airflow-core/docs/administration-and-deployment/priority-weight.rstpriority weights, Just a technical example. Needs to be kept but would be best to integrate in a real use case. (question)

example_dag_decorator.py


airflow-core/docs/core-concepts/dags.rst
Not needed as standalone example if the business example contains similar code. Rework into a real business example then Drop (error)

example_display_name.py


-
Drop standalone example. Ral display names should be added to all examples. → Drop (error)

example_dynamic_task_mapping.py


airflow-core/docs/authoring-and-scheduling/dynamic-task-mapping.rstDynamic task mappingNice example. But use one alternative from real life and drop the individual dag then (error)

example_dynamic_task_mapping_with_no_taskflow_operators.py


airflow-core/docs/authoring-and-scheduling/dynamic-task-mapping.rstDynamic task mapping

example_inlet_event_extra.py


-

Needs a proper example to showcase something useful. Unknown User (uranusjr) do you have a real world example to add to the storyline (question) Then rework on this together with outlet events. Then drop (error) this example as standaone file.

example_kubernetes_executor.py


providers/cncf/kubernetes/docs/kubernetes_executor.rst
Move to K8s Provider

example_latest_only_with_trigger.py


airflow-core/docs/core-concepts/dags.rst

I (= Unknown User (jscheffl)) do not understand the example as well as not the benefit of this Operator. Is there anybody who can provide a real life example of use? Else drop (error)

example_local_kubernetes_executor.py


-
Move to K8s Provider

example_nested_branch_dag.py


-
Not of any use standalone. Drop (error)

example_outlet_event_extra.py


-
Same like example_inlet_event_extra.py (question)

example_params_trigger_ui.py


airflow-core/docs/core-concepts/params.rst
Migrate to storyline, then drop the individual example (error)

example_params_ui_tutorial.py


airflow-core/docs/core-concepts/params.rst
If not all features can be transferred into storyline, keep this as individual example to be able to test all form elements.

example_passing_params_via_test_command.py


-
Migrate to storyline, then drop the individual example (error)

example_setup_teardown.py


-
Migrate to storyline, then drop the individual example (error)

example_setup_teardown_taskflow.py


-

example_simplest_dag.py


-
Merge with tutorial. No value standalone (error)

example_skip_dag.py


-
Should be integrated into storyline and be added as code reference in docs. Then Drop (error) this example

example_task_group_decorator.py


airflow-core/docs/core-concepts/dags.rstTask GroupMigrate to storyline, then drop the individual example (error)

example_task_group.py


-Task Group

example_time_delta_sensor_async.py


-
Merge with tutorial. No value standalone (error)

example_trigger_target_dag.py


-
Ups, this DAG was forgotten to be moved to standard provider (warning) → But anyway: Migrate to storyline, then drop the individual example (error)

example_workday_timetable.py


code copied, not inlined in airflow-core/docs/howto/timetable.rst
Use as a base for one of the examples in the storyline, then Drop (error)

example_xcomargs.py


-
Use as a base for one of the examples in the storyline, then Drop (error)

example_xcom.py


-
Use as a base for one of the examples in the storyline, then Drop (error)

tutorial_dag.py


-

Consolidate with tutorial.py

(warning) Note is referenced in some tests and needs to be replaces with the other tutorial

tutorial_objectstorage.py


airflow-core/docs/tutorial/objectstorage.rst
(question)

tutorial.py


airflow-core/docs/core-concepts/dag-run.rst

airflow-core/docs/tutorial/fundamentals.rst


Keep as starter tutorial, beautify using technical rules (tick)

tutorial_taskflow_api.py


airflow-core/docs/tutorial/taskflow.rst

tutorial_taskflow_api_virtualenv.py


-

Check if features can be added to storyline for technical completeness, then drop (error)

tutorial_taskflow_templates.py


-
Check if features can be added to storyline for technical completeness, then drop (error)

Standard - providers/standard/tests/system/standard

example_bash_decorator.py


providers/standard/docs/operators/bash.rst
Keep

example_bash_operator.py


providers/standard/docs/operators/bash.rst

airflow-core/docs/core-concepts/debug.rst

contributing-docs/quick-start-ide/contributors_quick_start_vscode.rst


Keep

example_branch_datetime_operator.py


providers/standard/docs/operators/datetime.rst
Keep, but need a better idea

example_branch_day_of_week_operator.py


providers/standard/docs/operators/datetime.rst
needs a better idea (warning)

example_branch_operator_decorator.py


providers/standard/docs/operators/python.rst
Used in the general example storyline, then Drop (error)

example_branch_operator.py


providers/standard/docs/operators/python.rst
Used in the general example storyline, then Drop (error)

example_external_task_child_deferrable.py


-
Drop (error)

example_external_task_marker_dag.py


providers/standard/docs/sensors/external_task_sensor.rst
Drop (error)

example_external_task_parent_deferrable.py


providers/standard/docs/sensors/external_task_sensor.rst
Drop (error)

example_latest_only.py


providers/standard/docs/operators/latest_only.rst
needs a better idea (warning)

example_python_decorator.py


airflow-core/docs/tutorial/taskflow.rst

providers/standard/docs/operators/python.rst


Used in the general example storyline, then Drop (error)

example_python_operator.py


airflow-core/docs/best-practices.rst

providers/standard/docs/operators/python.rst


Used in the general example storyline, then Drop (error)

example_sensor_decorator.py


airflow-core/docs/tutorial/taskflow.rst

providers/standard/docs/sensors/python.rst


Used in the general example storyline, then Drop (error)

example_sensors.py


providers/standard/docs/sensors/bash.rst

providers/standard/docs/sensors/datetime.rst

providers/standard/docs/sensors/file.rst

providers/standard/docs/sensors/python.rst


Used in the general example storyline, then Drop (error)

example_short_circuit_decorator.py


providers/standard/docs/operators/python.rst
needs a better idea (warning)

example_short_circuit_operator.py


providers/standard/docs/operators/python.rst
needs a better idea (warning)

example_trigger_controller_dag.py


providers/standard/docs/operators/trigger_dag_run.rst
See example_trigger_target_dag.py

ArangoDB - providers/arangodb/src/airflow/providers/arangodb/example_dags

example_arangodb.py


providers/arangodb/docs/operators/index.rst
Keep

Oracle - providers/oracle/src/airflow/providers/oracle/example_dags

example_oracle.py


providers/oracle/docs/operators.rst
Keep

Edge - providers/edge3/src/airflow/providers/edge3/example_dags

integration_test.py


-
Keep for the moment

win_notepad.py


-
Keep for the moment

win_test.py


-
Keep for the moment


List of features to cover in new examples

FeatureDescription
Dag authoringSimplest code that allows you to create a dag.
TaskflowUsing decorators to define tasks and dags.
Trigger rulesDecide condition for a task to run.
Setup/TeardownAllow setting setup/teardowns for tasks.
OperatorsShowing a way to use pre-existing pieces to achieve tasks.
Custom weightsUsed to prioritize tasks.
Dynamic task mappingAllow dynamically at runtime generate number of tasks based on inputs.
DAG paramsLet user provide dags for a dagrun.
AssetsLogical grouping of data, can be updated by dags and used to trigger downstream dags
AssetWatcherEvent driven scheduling