This is an automated email from the ASF dual-hosted git repository.
shahar1 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new d7b672f44e0 Update docs and diagrams for the Airflow 3 Dag processor
(#74073)
d7b672f44e0 is described below
commit d7b672f44e075fba9e5bf04a4eb7edd6286e293d
Author: 백형준 <[email protected]>
AuthorDate: Sat Oct 3 01:53:52 2026 +0900
Update docs and diagrams for the Airflow 3 Dag processor (#74073)
The architecture diagrams and several docs pages still described the
Airflow 2 layout, where the scheduler parses Dag files and workers talk to the
metadata database directly. In Airflow 3 the Dag processor is required, the
scheduler never accesses Dag bundles, and workers reach the database only
through the API server's Execution API, as the overview text above the diagrams
already says.
---
.../modules_management.rst | 2 +-
airflow-core/docs/best-practices.rst | 10 ++++-----
airflow-core/docs/core-concepts/dags.rst | 6 +++---
airflow-core/docs/core-concepts/overview.rst | 4 +---
.../img/diagram_basic_airflow_architecture.md5sum | 2 +-
.../img/diagram_basic_airflow_architecture.png | Bin 116194 -> 127853 bytes
.../docs/img/diagram_basic_airflow_architecture.py | 9 +++++---
...agram_dag_processor_airflow_architecture.md5sum | 2 +-
.../diagram_dag_processor_airflow_architecture.png | Bin 237657 -> 237572 bytes
.../diagram_dag_processor_airflow_architecture.py | 23 ++++++++++++---------
...diagram_distributed_airflow_architecture.md5sum | 2 +-
.../diagram_distributed_airflow_architecture.png | Bin 192497 -> 215541 bytes
.../diagram_distributed_airflow_architecture.py | 20 +++++++++++-------
13 files changed, 45 insertions(+), 35 deletions(-)
diff --git
a/airflow-core/docs/administration-and-deployment/modules_management.rst
b/airflow-core/docs/administration-and-deployment/modules_management.rst
index f5125d05ffe..74b76db92b8 100644
--- a/airflow-core/docs/administration-and-deployment/modules_management.rst
+++ b/airflow-core/docs/administration-and-deployment/modules_management.rst
@@ -229,7 +229,7 @@ You should import such shared Dags using full path
(starting from the directory
from my_company.my_custom_dags.base_dag import BaseDag # This is cool
The relative imports are counter-intuitive, and depending on how you start
your python code, they can behave
-differently. In Airflow the same Dag file might be parsed in different
contexts (by schedulers, by workers
+differently. In Airflow the same Dag file might be parsed in different
contexts (by the Dag processor, by workers
or during tests) and in those cases, relative imports might behave
differently. Always use full
python package paths when you import anything in Airflow Dags, this will save
you a lot of troubles.
You can read more about relative import caveats in
diff --git a/airflow-core/docs/best-practices.rst
b/airflow-core/docs/best-practices.rst
index 8959125f29c..3dc9adbc1da 100644
--- a/airflow-core/docs/best-practices.rst
+++ b/airflow-core/docs/best-practices.rst
@@ -103,7 +103,7 @@ You should avoid writing the top level code which is not
necessary to create Ope
and build Dag relations between them. This is because of the design decision
for the scheduler of Airflow
and the impact the top-level code parsing speed on both performance and
scalability of Airflow.
-Airflow scheduler executes the code outside the Operator's ``execute`` methods
with the minimum interval of
+Airflow Dag processor executes the code outside the Operator's ``execute``
methods with the minimum interval of
:ref:`min_file_process_interval<config:dag_processor__min_file_process_interval>`
seconds. This is done in order
to allow dynamic scheduling of the Dags - where scheduling and dependencies
might change over time and
impact the next schedule of the Dag. Airflow scheduler tries to continuously
make sure that what you have
@@ -411,7 +411,7 @@ or if you need to deserialize a json object from the
variable :
{{ var.json.<variable_name> }}
In top-level code, variables using jinja templates do not produce a request
until a task is running, whereas,
-``Variable.get()`` produces a request every time the Dag file is parsed by the
scheduler if caching is not enabled.
+``Variable.get()`` produces a request every time the Dag file is parsed by the
Dag processor if caching is not enabled.
Using ``Variable.get()`` without :ref:`enabling
caching<config:secrets__use_cache>` will lead to suboptimal
performance in the Dag file processing.
In some cases this can cause the Dag file to timeout before it is fully parsed.
@@ -501,10 +501,10 @@ Avoid triggering Dags immediately after changing them or
any other accompanying
Dag folder.
You should give the system sufficient time to process the changed files. This
takes several steps.
-First the files have to be distributed to scheduler - usually via distributed
filesystem or Git-Sync, then
-scheduler has to parse the Python files and store them in the database.
Depending on your configuration,
+First the files have to be distributed to the Dag processor - usually via
distributed filesystem or Git-Sync, then
+the Dag processor has to parse the Python files and store them in the
database. Depending on your configuration,
speed of your distributed filesystem, number of files, number of Dags, number
of changes in the files,
-sizes of the files, number of schedulers, speed of CPUS, this can take from
seconds to minutes, in extreme
+sizes of the files, number of Dag processors, speed of CPUS, this can take
from seconds to minutes, in extreme
cases many minutes. You should wait for your Dag to appear in the UI to be
able to trigger it.
In case you see long delays between updating it and the time it is ready to be
triggered, you can look
diff --git a/airflow-core/docs/core-concepts/dags.rst
b/airflow-core/docs/core-concepts/dags.rst
index 18acf9bbdec..2909eb5430e 100644
--- a/airflow-core/docs/core-concepts/dags.rst
+++ b/airflow-core/docs/core-concepts/dags.rst
@@ -814,8 +814,8 @@ the drain until its Dag runs have been created and
finished. While a Dag is drai
to make the Dag active again.
Dags can be deactivated (do not confuse it with ``Active`` tag in the UI) by
removing them from the
-``DAGS_FOLDER``. When scheduler parses the ``DAGS_FOLDER`` and misses the Dag
that it had seen
-before and stored in the database it will set is as deactivated. The metadata
and history of the
+``DAGS_FOLDER``. When the Dag processor parses the ``DAGS_FOLDER`` and misses
the Dag that it had seen
+before and stored in the database it will set it as deactivated. The metadata
and history of the
Dag is kept for deactivated Dags and when the Dag is re-added to the
``DAGS_FOLDER`` it will be again
activated and history will be visible. You cannot activate/deactivate Dag via
UI or API, this
can only be done by removing files from the ``DAGS_FOLDER``. Once again - no
data for historical runs of the
@@ -829,7 +829,7 @@ see the information about those you will see the error that
the Dag is missing.
You can also delete the Dag metadata from the metadata database using UI or
API, but it does not
always result in disappearing of the Dag from the UI - which might be also
initially a bit confusing.
If the Dag is still in ``DAGS_FOLDER`` when you delete the metadata, the Dag
will re-appear as
-Scheduler will parse the folder, only historical runs information for the Dag
will be removed.
+the Dag processor will parse the folder, only historical runs information for
the Dag will be removed.
This all means that if you want to actually delete a Dag and its all
historical metadata, you need to do
it in three steps:
diff --git a/airflow-core/docs/core-concepts/overview.rst
b/airflow-core/docs/core-concepts/overview.rst
index 6c7d374a2d5..d793527e731 100644
--- a/airflow-core/docs/core-concepts/overview.rst
+++ b/airflow-core/docs/core-concepts/overview.rst
@@ -126,12 +126,10 @@ The meaning of the different connection types in the
diagrams below is as follow
* **black dashed lines** represent control flow of workers by the *scheduler*
(via executor)
* **black solid lines** represent accessing the UI to manage execution of the
workflows
* **red dashed lines** represent accessing the *metadata database*
+* **green solid lines** represent *workers* communicating with the *API
server* through the *Execution API*
.. _overview-basic-airflow-architecture:
-..
- TODO AIP-72: These diagrams need to be updated to reflect AF3 changes like
bundles, required Dag processor, execution api, etc.
-
Basic Airflow deployment
........................
diff --git a/airflow-core/docs/img/diagram_basic_airflow_architecture.md5sum
b/airflow-core/docs/img/diagram_basic_airflow_architecture.md5sum
index d81957047b0..f4df323d82f 100644
--- a/airflow-core/docs/img/diagram_basic_airflow_architecture.md5sum
+++ b/airflow-core/docs/img/diagram_basic_airflow_architecture.md5sum
@@ -1 +1 @@
-2c62bb0602aab89a86375daeab7543ba
+6f76e94c45a7b8c8c8505c202d557d11
diff --git a/airflow-core/docs/img/diagram_basic_airflow_architecture.png
b/airflow-core/docs/img/diagram_basic_airflow_architecture.png
index ffef80a834c..11d9b30b7e2 100644
Binary files a/airflow-core/docs/img/diagram_basic_airflow_architecture.png and
b/airflow-core/docs/img/diagram_basic_airflow_architecture.png differ
diff --git a/airflow-core/docs/img/diagram_basic_airflow_architecture.py
b/airflow-core/docs/img/diagram_basic_airflow_architecture.py
index a523188bb37..77a5d16d380 100644
--- a/airflow-core/docs/img/diagram_basic_airflow_architecture.py
+++ b/airflow-core/docs/img/diagram_basic_airflow_architecture.py
@@ -68,14 +68,16 @@ def generate_basic_airflow_diagram():
):
user = User("Airflow User")
- dag_files = Custom("DAG files", MULTIPLE_FILES_IMAGE.as_posix())
- user >> Edge(color="brown", style="solid", reverse=False,
label="author\n\n") >> dag_files
+ dag_bundle = Custom("Dag bundle", MULTIPLE_FILES_IMAGE.as_posix())
+ user >> Edge(color="brown", style="solid", reverse=False,
label="author\n\n") >> dag_bundle
with Cluster("Parsing, Scheduling & Executing"):
+ dag_processor = Python("Dag processor")
scheduler = Python("Scheduler")
metadata_db = Custom("Metadata DB", DATABASE_IMAGE.as_posix())
scheduler >> Edge(color="red", style="dotted", reverse=True) >>
metadata_db
+ dag_processor >> Edge(color="red", style="dotted", reverse=True) >>
metadata_db
plugins_and_packages = Custom(
"Plugin folder\n& installed packages", PACKAGES_IMAGE.as_posix(),
color="transparent"
@@ -90,9 +92,10 @@ def generate_basic_airflow_diagram():
metadata_db >> Edge(color="red", style="dotted", reverse=True) >>
webserver
- dag_files >> Edge(color="brown", style="solid", label="read\n\n") >>
scheduler
+ dag_bundle >> Edge(color="brown", style="solid", label="read\n\n") >>
dag_processor
plugins_and_packages >> Edge(color="blue", style="solid",
label="install\n\n") >> scheduler
+ plugins_and_packages >> Edge(color="blue", style="solid",
label="install\n\n") >> dag_processor
plugins_and_packages >> Edge(color="blue", style="solid",
label="install\n\n") >> webserver
console.print(f"[green]Generating architecture image {image_file}")
diff --git
a/airflow-core/docs/img/diagram_dag_processor_airflow_architecture.md5sum
b/airflow-core/docs/img/diagram_dag_processor_airflow_architecture.md5sum
index ffd3b442a02..726d91140e3 100644
--- a/airflow-core/docs/img/diagram_dag_processor_airflow_architecture.md5sum
+++ b/airflow-core/docs/img/diagram_dag_processor_airflow_architecture.md5sum
@@ -1 +1 @@
-56006541287fe451c20e5fdd373d5456
+0a11d5b907ff20ed12776c2a66e389c4
diff --git
a/airflow-core/docs/img/diagram_dag_processor_airflow_architecture.png
b/airflow-core/docs/img/diagram_dag_processor_airflow_architecture.png
index c793e4f801b..086a7259660 100644
Binary files
a/airflow-core/docs/img/diagram_dag_processor_airflow_architecture.png and
b/airflow-core/docs/img/diagram_dag_processor_airflow_architecture.png differ
diff --git
a/airflow-core/docs/img/diagram_dag_processor_airflow_architecture.py
b/airflow-core/docs/img/diagram_dag_processor_airflow_architecture.py
index 8a1ba58e2f1..cc5a12268d9 100644
--- a/airflow-core/docs/img/diagram_dag_processor_airflow_architecture.py
+++ b/airflow-core/docs/img/diagram_dag_processor_airflow_architecture.py
@@ -67,7 +67,7 @@ def generate_dag_processor_airflow_diagram():
operations_user = User("Operations User")
deployment_manager = User("Deployment Manager")
- with Cluster("Security perimeter with no DAG code execution",
graph_attr={"bgcolor": "lightgrey"}):
+ with Cluster("Security perimeter with no Dag code execution",
graph_attr={"bgcolor": "lightgrey"}):
with Cluster("Scheduling\n\n"):
schedulers = Custom("Scheduler(s)",
PYTHON_MULTIPROCESS_LOGO.as_posix())
@@ -78,19 +78,19 @@ def generate_dag_processor_airflow_diagram():
metadata_db = Custom("Metadata DB", DATABASE_IMAGE.as_posix())
- dag_author = User("DAG Author")
+ dag_author = User("Dag Author")
- with Cluster("Security perimeter with DAG code execution"):
+ with Cluster("Security perimeter with Dag code execution"):
with Cluster("Execution"):
workers = Custom("Worker(s)",
PYTHON_MULTIPROCESS_LOGO.as_posix())
triggerer = Custom("Triggerer(s)",
PYTHON_MULTIPROCESS_LOGO.as_posix())
with Cluster("Parsing"):
- dag_processors = Custom("DAG\nProcessor(s)",
PYTHON_MULTIPROCESS_LOGO.as_posix())
- dag_files = Custom("DAG files", MULTIPLE_FILES_IMAGE.as_posix())
+ dag_processors = Custom("Dag\nprocessor(s)",
PYTHON_MULTIPROCESS_LOGO.as_posix())
+ dag_bundle = Custom("Dag bundle", MULTIPLE_FILES_IMAGE.as_posix())
plugins_and_packages = Custom("Plugin folder\n& installed packages",
PACKAGES_IMAGE.as_posix())
- dag_author >> Edge(color="brown", style="dashed", reverse=False,
label="author\n\n") >> dag_files
+ dag_author >> Edge(color="brown", style="dashed", reverse=False,
label="author\n\n") >> dag_bundle
(
deployment_manager
>> Edge(color="blue", style="solid", reverse=False,
label="install\n\n")
@@ -109,12 +109,15 @@ def generate_dag_processor_airflow_diagram():
metadata_db >> Edge(color="red", style="dotted", reverse=True) >>
webservers
metadata_db >> Edge(color="red", style="dotted", reverse=True) >>
schedulers
dag_processors >> Edge(color="red", style="dotted", reverse=True) >>
metadata_db
- workers >> Edge(color="red", style="dotted", reverse=True) >>
metadata_db
+ (
+ workers
+ >> Edge(color="darkgreen", style="solid", reverse=True,
label="Execution API\n\n")
+ >> webservers
+ )
triggerer >> Edge(color="red", style="dotted", reverse=True) >>
metadata_db
- dag_files >> Edge(color="brown", style="solid", label="sync\n\n") >>
workers
- dag_files >> Edge(color="brown", style="solid", label="sync\n\n") >>
dag_processors
- dag_files >> Edge(color="brown", style="solid", label="sync\n\n") >>
triggerer
+ dag_bundle >> Edge(color="brown", style="solid", label="sync\n\n") >>
workers
+ dag_bundle >> Edge(color="brown", style="solid", label="sync\n\n") >>
dag_processors
console.print(f"[green]Generating architecture image
{dag_processor_architecture_image_file}")
diff --git
a/airflow-core/docs/img/diagram_distributed_airflow_architecture.md5sum
b/airflow-core/docs/img/diagram_distributed_airflow_architecture.md5sum
index bcc881d52cc..4d6451684f7 100644
--- a/airflow-core/docs/img/diagram_distributed_airflow_architecture.md5sum
+++ b/airflow-core/docs/img/diagram_distributed_airflow_architecture.md5sum
@@ -1 +1 @@
-9715885d7403de716a7f5114c47f8027
+c029d2722da36b45d17fb62a31839759
diff --git a/airflow-core/docs/img/diagram_distributed_airflow_architecture.png
b/airflow-core/docs/img/diagram_distributed_airflow_architecture.png
index 66117d74e0b..a60bae9d426 100644
Binary files
a/airflow-core/docs/img/diagram_distributed_airflow_architecture.png and
b/airflow-core/docs/img/diagram_distributed_airflow_architecture.png differ
diff --git a/airflow-core/docs/img/diagram_distributed_airflow_architecture.py
b/airflow-core/docs/img/diagram_distributed_airflow_architecture.py
index 7c3ae3c3cd8..7791c705e41 100644
--- a/airflow-core/docs/img/diagram_distributed_airflow_architecture.py
+++ b/airflow-core/docs/img/diagram_distributed_airflow_architecture.py
@@ -65,13 +65,14 @@ def generate_distributed_airflow_diagram():
graph_attr=graph_attr,
edge_attr=edge_attr,
):
- dag_author = User("DAG Author")
+ dag_author = User("Dag Author")
deployment_manager = User("Deployment Manager")
- dag_files = Custom("DAG files", MULTIPLE_FILES_IMAGE.as_posix(),
height="1.8")
- dag_author >> Edge(color="brown", style="solid", reverse=False,
label="author\n\n") >> dag_files
+ dag_bundle = Custom("Dag bundle", MULTIPLE_FILES_IMAGE.as_posix(),
height="1.8")
+ dag_author >> Edge(color="brown", style="solid", reverse=False,
label="author\n\n") >> dag_bundle
with Cluster("Parsing, Scheduling & Executing"):
+ dag_processors = Custom("Dag processor(s)",
PYTHON_MULTIPROCESS_LOGO.as_posix())
schedulers = Custom("Scheduler(s)",
PYTHON_MULTIPROCESS_LOGO.as_posix())
workers = Custom("Worker(s)", PYTHON_MULTIPROCESS_LOGO.as_posix())
triggerer = Custom("Triggerer(s)",
PYTHON_MULTIPROCESS_LOGO.as_posix())
@@ -79,6 +80,7 @@ def generate_distributed_airflow_diagram():
metadata_db = Custom("Metadata DB", DATABASE_IMAGE.as_posix())
schedulers - Edge(color="black", style="dashed",
taillabel="[Executor]") - workers
+ dag_processors >> Edge(color="red", style="dotted", reverse=True) >>
metadata_db
schedulers >> Edge(color="red", style="dotted", reverse=True) >>
metadata_db
plugins_and_packages = Custom(
@@ -91,7 +93,6 @@ def generate_distributed_airflow_diagram():
>> plugins_and_packages
)
- workers >> Edge(color="red", style="dotted", reverse=True) >>
metadata_db
triggerer >> Edge(color="red", style="dotted", reverse=True) >>
metadata_db
operations_user = User("Operations User")
@@ -101,10 +102,14 @@ def generate_distributed_airflow_diagram():
webservers >> Edge(color="black", style="solid", reverse=True,
label="operate\n\n") >> operations_user
metadata_db >> Edge(color="red", style="dotted", reverse=True) >>
webservers
+ (
+ workers
+ >> Edge(color="darkgreen", style="solid", reverse=True,
label="Execution API\n\n")
+ >> webservers
+ )
- dag_files >> Edge(color="brown", style="solid", label="sync\n") >>
workers
- dag_files >> Edge(color="brown", style="solid", label="sync\n") >>
schedulers
- dag_files >> Edge(color="brown", style="solid", label="sync\n") >>
triggerer
+ dag_bundle >> Edge(color="brown", style="solid", label="sync\n") >>
workers
+ dag_bundle >> Edge(color="brown", style="solid", label="sync\n") >>
dag_processors
plugins_and_packages >> Edge(color="blue", style="solid",
label="install\n\n") >> workers
plugins_and_packages >> Edge(color="blue", style="solid",
label="install\n\n") >> schedulers
@@ -114,6 +119,7 @@ def generate_distributed_airflow_diagram():
>> triggerer
)
plugins_and_packages >> Edge(color="blue", style="solid",
label="install\n\n") >> webservers
+ plugins_and_packages >> Edge(color="blue", style="solid",
label="install\n\n") >> dag_processors
console.print(f"[green]Generating architecture image {image_file}")