dilnazanlid commented on code in PR #74333: URL: https://github.com/apache/airflow/pull/74333#discussion_r4230135820
########## 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: specific description added -- 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]
