kaxil commented on code in PR #74333: URL: https://github.com/apache/airflow/pull/74333#discussion_r4221356359
########## airflow-core/docs/administration-and-deployment/dag-importers.rst: ########## @@ -0,0 +1,154 @@ + .. Licensed to the Apache Software Foundation (ASF) under one + or more contributor license agreements. See the NOTICE file + distributed with this work for additional information + regarding copyright ownership. The ASF licenses this file + to you under the Apache License, Version 2.0 (the + "License"); you may not use this file except in compliance + with the License. You may obtain a copy of the License at + + .. http://www.apache.org/licenses/LICENSE-2.0 + + .. Unless required by applicable law or agreed to in writing, + software distributed under the License is distributed on an + "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + KIND, either express or implied. See the License for the + specific language governing permissions and limitations + under the License. + +.. _dag-importers: + +Dag Importers +============= + +|experimental| + +.. versionadded:: 3.4.0 + +.. warning:: + + The Dag importer interface may still change in a minor release until Airflow 3.5.0, while it is + being stabilized. If you maintain a custom importer, check the release notes when you upgrade. + +A Dag importer turns the contents of a :doc:`Dag bundle <dag-bundles>` into Dags. Importers let Airflow +load Dags from formats other than Python files, such as YAML or JSON pipeline definitions, without changes +to the Dag processor, the scheduler, or the workers. + +Airflow ships with two importers: + +**airflow.sdk.importers.PythonDagImporter** (``.py`` and ``.pyc``) + Loads a Python file as a module and collects the Dags it defines, as described in + :ref:`concepts-dag-loading`. + +**airflow.sdk.importers.ZipImporter** (``.zip``) + Imports each Python file in a zip archive, as described in :ref:`concepts-packaging-dags`. + +Each :doc:`language SDK <../authoring-and-scheduling/language-sdks/index>` with a configured coordinator +also registers an importer for its bundle format. + +How Airflow uses importers +-------------------------- + +An importer handles a Dag definition: a unit of Dag source, such as a Python file or a member of a zip +archive. + +* The Dag processor asks each importer to list the Dag definitions it can handle in a bundle, and queues the + files that hold them. A file holding several definitions, such as a zip archive, is parsed as one. +* A parsing process imports each definition with the importer that listed it. The resulting Dags go through + the same validation and :doc:`cluster policies <cluster-policies>` as Dags defined in Python, and the + importer's errors are shown as import errors. +* The source code the importer returns for a definition is shown in the **Code** tab of the UI. +* A worker uses the same importer to load the Dag before it runs a task. + +A Dag definition must be a file in the bundle, or nested in one, as a zip member is. The Dag processor +ignores any other definition with a warning. + +Configuring importers +--------------------- + +Register importers in :ref:`config:dag_processor__dag_importer_configs`. Each entry supplies the +``classpath`` of the importer and, optionally, the ``extensions`` it handles and the ``kwargs`` to create it +with: + +.. code-block:: ini + + [dag_processor] + dag_importer_configs = [ + { + "classpath": "my_company.importers.YamlDagImporter", + "extensions": [".yaml", ".yml"] + } + ] + +Without ``extensions``, the importer handles the extensions its class declares. An entry with only a +classpath can be given as a string, for example ``["my_company.importers.YamlDagImporter"]``. + +To register an importer for one bundle only, add an ``importers`` key, in the same format, to that bundle's +entry in :ref:`config:dag_processor__dag_bundle_config_list`: + +.. code-block:: ini + + [dag_processor] + dag_bundle_config_list = [ + { + "name": "pipelines", + "classpath": "airflow.dag_processing.bundles.local.LocalDagBundle", + "kwargs": {"path": "/opt/airflow/pipelines"}, + "importers": ["my_company.importers.YamlDagImporter"] + } + ] + +Each extension is handled by one importer. Importers are registered in this order, and a later importer +takes over an extension from an earlier one, with a warning in the logs: + +#. the built-in importers; +#. the importers of the language SDK coordinators; +#. ``dag_importer_configs``; +#. the bundle's ``importers``. + +This lets you replace a built-in importer, for example to handle ``.py`` files differently. + +.. important:: + + Workers import Dags with the same importers as the Dag processor. Install the package that provides your + importer, and set the same importer configuration, on the Dag processor, on the workers, and anywhere + else Dags are parsed, such as where you run ``airflow dags`` commands. + +Writing a custom importer +------------------------- + +An importer subclasses :class:`~airflow.sdk.importers.AbstractDagImporter`, typed with the class of Dag +definition it handles, and implements: + +``can_handle(definition)`` + Whether the importer can import a definition, a path, or a file name. +``list_dag_definitions(bundle, *, safe_mode)`` + Yield the Dag definitions in the bundle that the importer handles. +``import_definition(definition, bundle)`` + Import the Dags of one definition and return them in a :class:`~airflow.sdk.importers.DagImportResult`. +``get_source_code(definition)`` Review Comment: On current main this is `get_source_code(definition, dag_id=None)`, and `DagBag` calls it as `importer.get_source_code(definition, dag.dag_id)` ([dagbag.py](https://github.com/apache/airflow/blob/ef9feefd83fcd1cf7b18c245d424fd2eb5634211/airflow-core/src/airflow/dag_processing/dagbag.py#L344)). An importer written to the one-argument signature here raises `TypeError` on that call, which is caught and only logged, so its Dags show an empty Code tab. Can you rebase and document `dag_id` (the Dag whose own source is wanted; an importer that can't tell the Dags of one definition apart ignores it)? ########## airflow-core/docs/administration-and-deployment/dag-importers.rst: ########## @@ -0,0 +1,154 @@ + .. Licensed to the Apache Software Foundation (ASF) under one + or more contributor license agreements. See the NOTICE file + distributed with this work for additional information + regarding copyright ownership. The ASF licenses this file + to you under the Apache License, Version 2.0 (the + "License"); you may not use this file except in compliance + with the License. You may obtain a copy of the License at + + .. http://www.apache.org/licenses/LICENSE-2.0 + + .. Unless required by applicable law or agreed to in writing, + software distributed under the License is distributed on an + "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + KIND, either express or implied. See the License for the + specific language governing permissions and limitations + under the License. + +.. _dag-importers: + +Dag Importers +============= + +|experimental| + +.. versionadded:: 3.4.0 + +.. warning:: + + The Dag importer interface may still change in a minor release until Airflow 3.5.0, while it is + being stabilized. If you maintain a custom importer, check the release notes when you upgrade. + +A Dag importer turns the contents of a :doc:`Dag bundle <dag-bundles>` into Dags. Importers let Airflow +load Dags from formats other than Python files, such as YAML or JSON pipeline definitions, without changes +to the Dag processor, the scheduler, or the workers. + +Airflow ships with two importers: + +**airflow.sdk.importers.PythonDagImporter** (``.py`` and ``.pyc``) + Loads a Python file as a module and collects the Dags it defines, as described in + :ref:`concepts-dag-loading`. + +**airflow.sdk.importers.ZipImporter** (``.zip``) + Imports each Python file in a zip archive, as described in :ref:`concepts-packaging-dags`. + +Each :doc:`language SDK <../authoring-and-scheduling/language-sdks/index>` with a configured coordinator +also registers an importer for its bundle format. + +How Airflow uses importers +-------------------------- + +An importer handles a Dag definition: a unit of Dag source, such as a Python file or a member of a zip +archive. + +* The Dag processor asks each importer to list the Dag definitions it can handle in a bundle, and queues the + files that hold them. A file holding several definitions, such as a zip archive, is parsed as one. +* A parsing process imports each definition with the importer that listed it. The resulting Dags go through + the same validation and :doc:`cluster policies <cluster-policies>` as Dags defined in Python, and the + importer's errors are shown as import errors. +* The source code the importer returns for a definition is shown in the **Code** tab of the UI. +* A worker uses the same importer to load the Dag before it runs a task. Review Comment: These bullets hold for importers you write yourself, but not for the language SDK importers mentioned at line 45. For those, the Dag processor runs the SDK runtime and never calls `import_definition`, and a worker fails the task unless its queue is routed through `[sdk] queue_to_coordinator` ([`_fail_lang_sdk_task`](https://github.com/apache/airflow/blob/ef9feefd83fcd1cf7b18c245d424fd2eb5634211/task-sdk/src/airflow/sdk/execution_time/task_runner.py#L1033-L1050)). Maybe scope this list to Python-side importers and link to the language SDK page for the native case? ########## airflow-core/docs/administration-and-deployment/dag-importers.rst: ########## @@ -0,0 +1,154 @@ + .. Licensed to the Apache Software Foundation (ASF) under one + or more contributor license agreements. See the NOTICE file + distributed with this work for additional information + regarding copyright ownership. The ASF licenses this file + to you under the Apache License, Version 2.0 (the + "License"); you may not use this file except in compliance + with the License. You may obtain a copy of the License at + + .. http://www.apache.org/licenses/LICENSE-2.0 + + .. Unless required by applicable law or agreed to in writing, + software distributed under the License is distributed on an + "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + KIND, either express or implied. See the License for the + specific language governing permissions and limitations + under the License. + +.. _dag-importers: + +Dag Importers +============= + +|experimental| + +.. versionadded:: 3.4.0 + +.. warning:: + + The Dag importer interface may still change in a minor release until Airflow 3.5.0, while it is + being stabilized. If you maintain a custom importer, check the release notes when you upgrade. + +A Dag importer turns the contents of a :doc:`Dag bundle <dag-bundles>` into Dags. Importers let Airflow +load Dags from formats other than Python files, such as YAML or JSON pipeline definitions, without changes +to the Dag processor, the scheduler, or the workers. + +Airflow ships with two importers: + +**airflow.sdk.importers.PythonDagImporter** (``.py`` and ``.pyc``) + Loads a Python file as a module and collects the Dags it defines, as described in + :ref:`concepts-dag-loading`. + +**airflow.sdk.importers.ZipImporter** (``.zip``) + Imports each Python file in a zip archive, as described in :ref:`concepts-packaging-dags`. + +Each :doc:`language SDK <../authoring-and-scheduling/language-sdks/index>` with a configured coordinator +also registers an importer for its bundle format. + +How Airflow uses importers +-------------------------- + +An importer handles a Dag definition: a unit of Dag source, such as a Python file or a member of a zip +archive. + +* The Dag processor asks each importer to list the Dag definitions it can handle in a bundle, and queues the + files that hold them. A file holding several definitions, such as a zip archive, is parsed as one. +* A parsing process imports each definition with the importer that listed it. The resulting Dags go through + the same validation and :doc:`cluster policies <cluster-policies>` as Dags defined in Python, and the + importer's errors are shown as import errors. +* The source code the importer returns for a definition is shown in the **Code** tab of the UI. +* A worker uses the same importer to load the Dag before it runs a task. + +A Dag definition must be a file in the bundle, or nested in one, as a zip member is. The Dag processor +ignores any other definition with a warning. + +Configuring importers +--------------------- + +Register importers in :ref:`config:dag_processor__dag_importer_configs`. Each entry supplies the +``classpath`` of the importer and, optionally, the ``extensions`` it handles and the ``kwargs`` to create it +with: + +.. code-block:: ini + + [dag_processor] + dag_importer_configs = [ + { + "classpath": "my_company.importers.YamlDagImporter", + "extensions": [".yaml", ".yml"] + } + ] + +Without ``extensions``, the importer handles the extensions its class declares. An entry with only a +classpath can be given as a string, for example ``["my_company.importers.YamlDagImporter"]``. + +To register an importer for one bundle only, add an ``importers`` key, in the same format, to that bundle's +entry in :ref:`config:dag_processor__dag_bundle_config_list`: + +.. code-block:: ini + + [dag_processor] + dag_bundle_config_list = [ + { + "name": "pipelines", + "classpath": "airflow.dag_processing.bundles.local.LocalDagBundle", + "kwargs": {"path": "/opt/airflow/pipelines"}, + "importers": ["my_company.importers.YamlDagImporter"] + } + ] + +Each extension is handled by one importer. Importers are registered in this order, and a later importer +takes over an extension from an earlier one, with a warning in the logs: + +#. the built-in importers; +#. the importers of the language SDK coordinators; +#. ``dag_importer_configs``; +#. the bundle's ``importers``. + +This lets you replace a built-in importer, for example to handle ``.py`` files differently. + +.. important:: + + Workers import Dags with the same importers as the Dag processor. Install the package that provides your + importer, and set the same importer configuration, on the Dag processor, on the workers, and anywhere + else Dags are parsed, such as where you run ``airflow dags`` commands. + +Writing a custom importer +------------------------- + +An importer subclasses :class:`~airflow.sdk.importers.AbstractDagImporter`, typed with the class of Dag +definition it handles, and implements: + +``can_handle(definition)`` + Whether the importer can import a definition, a path, or a file name. +``list_dag_definitions(bundle, *, safe_mode)`` + Yield the Dag definitions in the bundle that the importer handles. Review Comment: Could this say what `bundle` actually is? Only the Dag processor manager passes the configured bundle. In the parse process, on workers and in the CLI, `DagBag` passes a `LocalDagBundle` built from the bundle's name and path, and for listing its `path` is the single file being parsed ([`_find_definitions`](https://github.com/apache/airflow/blob/ef9feefd83fcd1cf7b18c245d424fd2eb5634211/airflow-core/src/airflow/dag_processing/dagbag.py#L307-L313)). So an importer that walks `bundle.path` with `rglob` or `os.walk` lists nothing in the parse process (both return empty for a file path), and one that reads attributes of, say, a git bundle gets an `AttributeError`. Something like "rely only on `bundle.name` and `bundle.path`, and `path` may be a single file; `find_file_dag_definitions` handles both" would cover it. ########## airflow-core/docs/administration-and-deployment/dag-importers.rst: ########## @@ -0,0 +1,154 @@ + .. Licensed to the Apache Software Foundation (ASF) under one + or more contributor license agreements. See the NOTICE file + distributed with this work for additional information + regarding copyright ownership. The ASF licenses this file + to you under the Apache License, Version 2.0 (the + "License"); you may not use this file except in compliance + with the License. You may obtain a copy of the License at + + .. http://www.apache.org/licenses/LICENSE-2.0 + + .. Unless required by applicable law or agreed to in writing, + software distributed under the License is distributed on an + "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + KIND, either express or implied. See the License for the + specific language governing permissions and limitations + under the License. + +.. _dag-importers: + +Dag Importers +============= + +|experimental| + +.. versionadded:: 3.4.0 + +.. warning:: + + The Dag importer interface may still change in a minor release until Airflow 3.5.0, while it is + being stabilized. If you maintain a custom importer, check the release notes when you upgrade. + +A Dag importer turns the contents of a :doc:`Dag bundle <dag-bundles>` into Dags. Importers let Airflow +load Dags from formats other than Python files, such as YAML or JSON pipeline definitions, without changes +to the Dag processor, the scheduler, or the workers. + +Airflow ships with two importers: + +**airflow.sdk.importers.PythonDagImporter** (``.py`` and ``.pyc``) + Loads a Python file as a module and collects the Dags it defines, as described in + :ref:`concepts-dag-loading`. + +**airflow.sdk.importers.ZipImporter** (``.zip``) + Imports each Python file in a zip archive, as described in :ref:`concepts-packaging-dags`. + +Each :doc:`language SDK <../authoring-and-scheduling/language-sdks/index>` with a configured coordinator +also registers an importer for its bundle format. + +How Airflow uses importers +-------------------------- + +An importer handles a Dag definition: a unit of Dag source, such as a Python file or a member of a zip +archive. + +* The Dag processor asks each importer to list the Dag definitions it can handle in a bundle, and queues the + files that hold them. A file holding several definitions, such as a zip archive, is parsed as one. +* A parsing process imports each definition with the importer that listed it. The resulting Dags go through + the same validation and :doc:`cluster policies <cluster-policies>` as Dags defined in Python, and the + importer's errors are shown as import errors. +* The source code the importer returns for a definition is shown in the **Code** tab of the UI. +* A worker uses the same importer to load the Dag before it runs a task. + +A Dag definition must be a file in the bundle, or nested in one, as a zip member is. The Dag processor +ignores any other definition with a warning. + +Configuring importers +--------------------- + +Register importers in :ref:`config:dag_processor__dag_importer_configs`. Each entry supplies the +``classpath`` of the importer and, optionally, the ``extensions`` it handles and the ``kwargs`` to create it +with: + +.. code-block:: ini + + [dag_processor] + dag_importer_configs = [ + { + "classpath": "my_company.importers.YamlDagImporter", + "extensions": [".yaml", ".yml"] + } + ] + +Without ``extensions``, the importer handles the extensions its class declares. An entry with only a +classpath can be given as a string, for example ``["my_company.importers.YamlDagImporter"]``. + +To register an importer for one bundle only, add an ``importers`` key, in the same format, to that bundle's +entry in :ref:`config:dag_processor__dag_bundle_config_list`: + +.. code-block:: ini + + [dag_processor] + dag_bundle_config_list = [ + { + "name": "pipelines", + "classpath": "airflow.dag_processing.bundles.local.LocalDagBundle", + "kwargs": {"path": "/opt/airflow/pipelines"}, + "importers": ["my_company.importers.YamlDagImporter"] + } + ] + +Each extension is handled by one importer. Importers are registered in this order, and a later importer +takes over an extension from an earlier one, with a warning in the logs: + +#. the built-in importers; +#. the importers of the language SDK coordinators; +#. ``dag_importer_configs``; +#. the bundle's ``importers``. + +This lets you replace a built-in importer, for example to handle ``.py`` files differently. + +.. important:: + + Workers import Dags with the same importers as the Dag processor. Install the package that provides your + importer, and set the same importer configuration, on the Dag processor, on the workers, and anywhere + else Dags are parsed, such as where you run ``airflow dags`` commands. + +Writing a custom importer +------------------------- + +An importer subclasses :class:`~airflow.sdk.importers.AbstractDagImporter`, typed with the class of Dag +definition it handles, and implements: + +``can_handle(definition)`` + Whether the importer can import a definition, a path, or a file name. +``list_dag_definitions(bundle, *, safe_mode)`` + Yield the Dag definitions in the bundle that the importer handles. +``import_definition(definition, bundle)`` + Import the Dags of one definition and return them in a :class:`~airflow.sdk.importers.DagImportResult`. +``get_source_code(definition)`` + Return the source of a definition and its language, as a :class:`~airflow.sdk.importers.DagSourceCode`. + +For file-based formats, :func:`~airflow.sdk.importers.find_file_dag_definitions` walks the bundle the same +way the built-in importers do: it applies :ref:`.airflowignore <concepts:airflowignore>` and yields a +:class:`~airflow.sdk.importers.FilesystemDagDefinition` for each file with one of the given extensions. + +Keep the following in mind when you write an importer: + +* Report a definition that fails to import as a :class:`~airflow.sdk.importers.DagImportError` in the + result instead of raising, so it is shown as an import error. In ``list_dag_definitions``, yield a + ``DagImportError`` for a file that cannot be read, so the rest of the bundle is still listed. +* Declare ``supported_extensions`` as an attribute: Airflow replaces it with the ``extensions`` of the + importer's configuration entry. +* ``safe_mode`` is the value of :ref:`config:core__dag_discovery_safe_mode`. To skip files that clearly hold no Review Comment: A worker builds its `DagBag` with `safe_mode=False` ([task_runner.py](https://github.com/apache/airflow/blob/ef9feefd83fcd1cf7b18c245d424fd2eb5634211/task-sdk/src/airflow/sdk/execution_time/task_runner.py#L1069-L1071)), whatever `dag_discovery_safe_mode` says. An importer that filters on `might_contain_dag` sees a different value there, so could the sentence mention the worker case? ########## airflow-core/docs/administration-and-deployment/dag-importers.rst: ########## @@ -0,0 +1,154 @@ + .. Licensed to the Apache Software Foundation (ASF) under one + or more contributor license agreements. See the NOTICE file + distributed with this work for additional information + regarding copyright ownership. The ASF licenses this file + to you under the Apache License, Version 2.0 (the + "License"); you may not use this file except in compliance + with the License. You may obtain a copy of the License at + + .. http://www.apache.org/licenses/LICENSE-2.0 + + .. Unless required by applicable law or agreed to in writing, + software distributed under the License is distributed on an + "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + KIND, either express or implied. See the License for the + specific language governing permissions and limitations + under the License. + +.. _dag-importers: + +Dag Importers +============= + +|experimental| + +.. versionadded:: 3.4.0 + +.. warning:: + + The Dag importer interface may still change in a minor release until Airflow 3.5.0, while it is + being stabilized. If you maintain a custom importer, check the release notes when you upgrade. + +A Dag importer turns the contents of a :doc:`Dag bundle <dag-bundles>` into Dags. Importers let Airflow +load Dags from formats other than Python files, such as YAML or JSON pipeline definitions, without changes +to the Dag processor, the scheduler, or the workers. + +Airflow ships with two importers: + +**airflow.sdk.importers.PythonDagImporter** (``.py`` and ``.pyc``) + Loads a Python file as a module and collects the Dags it defines, as described in + :ref:`concepts-dag-loading`. + +**airflow.sdk.importers.ZipImporter** (``.zip``) + Imports each Python file in a zip archive, as described in :ref:`concepts-packaging-dags`. + +Each :doc:`language SDK <../authoring-and-scheduling/language-sdks/index>` with a configured coordinator +also registers an importer for its bundle format. + +How Airflow uses importers +-------------------------- + +An importer handles a Dag definition: a unit of Dag source, such as a Python file or a member of a zip +archive. + +* The Dag processor asks each importer to list the Dag definitions it can handle in a bundle, and queues the + files that hold them. A file holding several definitions, such as a zip archive, is parsed as one. +* A parsing process imports each definition with the importer that listed it. The resulting Dags go through + the same validation and :doc:`cluster policies <cluster-policies>` as Dags defined in Python, and the + importer's errors are shown as import errors. +* The source code the importer returns for a definition is shown in the **Code** tab of the UI. +* A worker uses the same importer to load the Dag before it runs a task. + +A Dag definition must be a file in the bundle, or nested in one, as a zip member is. The Dag processor +ignores any other definition with a warning. + +Configuring importers +--------------------- + +Register importers in :ref:`config:dag_processor__dag_importer_configs`. Each entry supplies the +``classpath`` of the importer and, optionally, the ``extensions`` it handles and the ``kwargs`` to create it +with: + +.. code-block:: ini + + [dag_processor] + dag_importer_configs = [ + { + "classpath": "my_company.importers.YamlDagImporter", + "extensions": [".yaml", ".yml"] + } + ] + +Without ``extensions``, the importer handles the extensions its class declares. An entry with only a +classpath can be given as a string, for example ``["my_company.importers.YamlDagImporter"]``. + +To register an importer for one bundle only, add an ``importers`` key, in the same format, to that bundle's +entry in :ref:`config:dag_processor__dag_bundle_config_list`: + +.. code-block:: ini + + [dag_processor] + dag_bundle_config_list = [ + { + "name": "pipelines", + "classpath": "airflow.dag_processing.bundles.local.LocalDagBundle", + "kwargs": {"path": "/opt/airflow/pipelines"}, + "importers": ["my_company.importers.YamlDagImporter"] + } + ] + +Each extension is handled by one importer. Importers are registered in this order, and a later importer +takes over an extension from an earlier one, with a warning in the logs: + +#. the built-in importers; +#. the importers of the language SDK coordinators; +#. ``dag_importer_configs``; +#. the bundle's ``importers``. + +This lets you replace a built-in importer, for example to handle ``.py`` files differently. Review Comment: Replacing the `.py` importer doesn't reach `.py` members of zip archives. `ZipImporter` builds its own `PythonDagImporter` unless it's given `internal_importers` ([zip_importer.py](https://github.com/apache/airflow/blob/ef9feefd83fcd1cf7b18c245d424fd2eb5634211/task-sdk/src/airflow/sdk/importers/zip_importer.py#L133-L136)), and extension ownership only filters top-level files. Could this sentence say that, or point at the `internal_importers` kwarg? ########## airflow-core/docs/administration-and-deployment/dag-importers.rst: ########## @@ -0,0 +1,154 @@ + .. Licensed to the Apache Software Foundation (ASF) under one + or more contributor license agreements. See the NOTICE file + distributed with this work for additional information + regarding copyright ownership. The ASF licenses this file + to you under the Apache License, Version 2.0 (the + "License"); you may not use this file except in compliance + with the License. You may obtain a copy of the License at + + .. http://www.apache.org/licenses/LICENSE-2.0 + + .. Unless required by applicable law or agreed to in writing, + software distributed under the License is distributed on an + "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + KIND, either express or implied. See the License for the + specific language governing permissions and limitations + under the License. + +.. _dag-importers: + +Dag Importers +============= + +|experimental| + +.. versionadded:: 3.4.0 + +.. warning:: + + The Dag importer interface may still change in a minor release until Airflow 3.5.0, while it is Review Comment: Is 3.5.0 an agreed target? Read literally, this promises the interface stops changing at 3.5.0. `release-process.rst` defines experimental with no end date, and the config option and the `AbstractDagImporter` docstring use the plain marker. Unless that's been decided I'd drop the version here and in the two other places this PR repeats it (public-airflow-interface.rst and task-sdk/docs/api.rst). ########## airflow-core/docs/administration-and-deployment/dag-importers.rst: ########## @@ -0,0 +1,154 @@ + .. Licensed to the Apache Software Foundation (ASF) under one + or more contributor license agreements. See the NOTICE file + distributed with this work for additional information + regarding copyright ownership. The ASF licenses this file + to you under the Apache License, Version 2.0 (the + "License"); you may not use this file except in compliance + with the License. You may obtain a copy of the License at + + .. http://www.apache.org/licenses/LICENSE-2.0 + + .. Unless required by applicable law or agreed to in writing, + software distributed under the License is distributed on an + "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + KIND, either express or implied. See the License for the + specific language governing permissions and limitations + under the License. + +.. _dag-importers: + +Dag Importers +============= + +|experimental| + +.. versionadded:: 3.4.0 + +.. warning:: + + The Dag importer interface may still change in a minor release until Airflow 3.5.0, while it is + being stabilized. If you maintain a custom importer, check the release notes when you upgrade. + +A Dag importer turns the contents of a :doc:`Dag bundle <dag-bundles>` into Dags. Importers let Airflow +load Dags from formats other than Python files, such as YAML or JSON pipeline definitions, without changes +to the Dag processor, the scheduler, or the workers. + +Airflow ships with two importers: + +**airflow.sdk.importers.PythonDagImporter** (``.py`` and ``.pyc``) + Loads a Python file as a module and collects the Dags it defines, as described in + :ref:`concepts-dag-loading`. + +**airflow.sdk.importers.ZipImporter** (``.zip``) + Imports each Python file in a zip archive, as described in :ref:`concepts-packaging-dags`. + +Each :doc:`language SDK <../authoring-and-scheduling/language-sdks/index>` with a configured coordinator +also registers an importer for its bundle format. + +How Airflow uses importers +-------------------------- + +An importer handles a Dag definition: a unit of Dag source, such as a Python file or a member of a zip +archive. + +* The Dag processor asks each importer to list the Dag definitions it can handle in a bundle, and queues the + files that hold them. A file holding several definitions, such as a zip archive, is parsed as one. +* A parsing process imports each definition with the importer that listed it. The resulting Dags go through + the same validation and :doc:`cluster policies <cluster-policies>` as Dags defined in Python, and the + importer's errors are shown as import errors. +* The source code the importer returns for a definition is shown in the **Code** tab of the UI. +* A worker uses the same importer to load the Dag before it runs a task. + +A Dag definition must be a file in the bundle, or nested in one, as a zip member is. The Dag processor +ignores any other definition with a warning. + +Configuring importers +--------------------- + +Register importers in :ref:`config:dag_processor__dag_importer_configs`. Each entry supplies the +``classpath`` of the importer and, optionally, the ``extensions`` it handles and the ``kwargs`` to create it +with: + +.. code-block:: ini + + [dag_processor] + dag_importer_configs = [ + { + "classpath": "my_company.importers.YamlDagImporter", + "extensions": [".yaml", ".yml"] + } + ] + +Without ``extensions``, the importer handles the extensions its class declares. An entry with only a +classpath can be given as a string, for example ``["my_company.importers.YamlDagImporter"]``. + +To register an importer for one bundle only, add an ``importers`` key, in the same format, to that bundle's +entry in :ref:`config:dag_processor__dag_bundle_config_list`: + +.. code-block:: ini + + [dag_processor] + dag_bundle_config_list = [ + { + "name": "pipelines", + "classpath": "airflow.dag_processing.bundles.local.LocalDagBundle", + "kwargs": {"path": "/opt/airflow/pipelines"}, + "importers": ["my_company.importers.YamlDagImporter"] + } + ] + +Each extension is handled by one importer. Importers are registered in this order, and a later importer +takes over an extension from an earlier one, with a warning in the logs: + +#. the built-in importers; +#. the importers of the language SDK coordinators; +#. ``dag_importer_configs``; +#. the bundle's ``importers``. + +This lets you replace a built-in importer, for example to handle ``.py`` files differently. + +.. important:: + + Workers import Dags with the same importers as the Dag processor. Install the package that provides your + importer, and set the same importer configuration, on the Dag processor, on the workers, and anywhere + else Dags are parsed, such as where you run ``airflow dags`` commands. + +Writing a custom importer +------------------------- + +An importer subclasses :class:`~airflow.sdk.importers.AbstractDagImporter`, typed with the class of Dag +definition it handles, and implements: + +``can_handle(definition)`` + Whether the importer can import a definition, a path, or a file name. +``list_dag_definitions(bundle, *, safe_mode)`` + Yield the Dag definitions in the bundle that the importer handles. +``import_definition(definition, bundle)`` + Import the Dags of one definition and return them in a :class:`~airflow.sdk.importers.DagImportResult`. +``get_source_code(definition)`` + Return the source of a definition and its language, as a :class:`~airflow.sdk.importers.DagSourceCode`. + +For file-based formats, :func:`~airflow.sdk.importers.find_file_dag_definitions` walks the bundle the same +way the built-in importers do: it applies :ref:`.airflowignore <concepts:airflowignore>` and yields a +:class:`~airflow.sdk.importers.FilesystemDagDefinition` for each file with one of the given extensions. + +Keep the following in mind when you write an importer: + +* Report a definition that fails to import as a :class:`~airflow.sdk.importers.DagImportError` in the + result instead of raising, so it is shown as an import error. In ``list_dag_definitions``, yield a + ``DagImportError`` for a file that cannot be read, so the rest of the bundle is still listed. +* Declare ``supported_extensions`` as an attribute: Airflow replaces it with the ``extensions`` of the + importer's configuration entry. +* ``safe_mode`` is the value of :ref:`config:core__dag_discovery_safe_mode`. To skip files that clearly hold no + Dag before a parsing process is started for them, override + :meth:`~airflow.sdk.importers.AbstractDagImporter.might_contain_dag` and call it in + ``list_dag_definitions``. +* Airflow creates one importer instance per bundle and reuses it for every definition. The Dag processor + creates it before it forks the parsing processes, so keep ``__init__`` cheap and do not open connections or + start threads there. +* The built-in ``dagbag_import_timeout`` applies to Python modules only; a slow importer is bounded by + :ref:`config:dag_processor__dag_file_processor_timeout`. Review Comment: Listing isn't bounded by either timeout. The manager calls `list_dag_definitions` in its own loop ([`_find_files_in_bundle`](https://github.com/apache/airflow/blob/ef9feefd83fcd1cf7b18c245d424fd2eb5634211/airflow-core/src/airflow/dag_processing/manager.py#L1016-L1025)) with only a try/except around it, and `dag_file_processor_timeout` is checked against the parse child only. Since the `might_contain_dag` bullet above encourages reading files during listing, a slow listing stalls the Dag processor's main loop for every bundle. Can this say that listing runs in the Dag processor itself with no timeout and has to stay cheap, and that only `import_definition` is bounded? Also `dagbag_import_timeout` isn't Python-only anymore: native language SDK parsing honours it too ([lang_sdk_processor.py](https://github.com/apache/airflow/blob/ef9feefd83fcd1cf7b18c245d424fd2eb5634211/airflow-core/src/airflow/dag_processing/lang_sdk_processor.py#L104-L109)). -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
