dilnazanlid commented on code in PR #73728:
URL: https://github.com/apache/airflow/pull/73728#discussion_r4116956003
##########
task-sdk/src/airflow/sdk/importers/base.py:
##########
@@ -483,6 +485,45 @@ def supported_extensions(self) -> list[str]:
"""Return all registered file extensions."""
return sorted(set(self._extension_importers) |
set(self._extension_specs))
+ def list_dag_definitions(
+ self,
+ bundle: BaseDagBundle,
+ *,
+ safe_mode: bool = True,
+ ) -> Iterator[tuple[AbstractDagImporter[Any], DagDefinition |
DagImportError]]:
+ """
+ List DAG definitions in a bundle across all registered importers.
+
+ Each item is paired with the importer that listed it, which is the one
to import it
+ with: a composite importer lists definitions no extension lookup would
route back to
+ it, such as the members of an archive. A :class:`DagImportError` item
is a
+ discovery-time failure rather than a source to import.
+ """
+ for spec in list(self._extension_specs.values()):
+ self._materialise_spec(spec)
+
+ for importer in self._ordered_importers:
+ for item in importer.list_dag_definitions(bundle,
safe_mode=safe_mode):
+ # A file whose extension was claimed by a later registration
belongs to that
+ # importer; yielding it here as well would import it twice.
+ if (
+ isinstance(item, FilesystemDagDefinition)
+ and self._extension_importers.get(item.suffix, importer)
is not importer
+ ):
+ continue
+ yield importer, item
+
+ def _materialise_spec(self, spec: _ImporterSpec) ->
AbstractDagImporter[Any]:
+ """Instantiate a configured spec and take over every extension it was
registered for."""
+ importer = self._instantiate_spec(spec)
+ for ext, registered in list(self._extension_specs.items()):
+ if registered is spec:
+ self._extension_importers[ext] = importer
+ del self._extension_specs[ext]
+ if importer not in self._ordered_importers:
+ self._ordered_importers.append(importer)
+ return importer
Review Comment:
Good catch, thanks. `list_dag_definitions` now takes one pending spec at a
time, so a spec registered for several extensions is instantiated once and
takes over all of its extensions.
The duplicate matters for composite importers. For plain files the suffix
check already dropped the second copy, but zip members were yielded twice.
##########
task-sdk/src/airflow/sdk/importers/python_importer.py:
##########
@@ -147,18 +147,18 @@ def list_dag_definitions(
safe_mode: bool = True,
) -> Iterator[FileDagDefinition | DagImportError]:
"""
- List Python DAG files in a bundle matching supported extensions.
+ List Python DAG files under the bundle path matching supported
extensions.
A lightweight content sniff (``might_contain_dag``) is applied here so
files that
clearly hold no DAG never become definitions -- keeping the discovered
set (and the
eventual parse-process count) close to the number of real DAG files.
Zip members are
discovered by :class:`..zip_importer.ZipImporter`, not here.
"""
- if not bundle.path.is_dir():
- return
for definition in find_file_dag_definitions(bundle.path,
self.supported_extensions):
if self.might_contain_dag(definition, safe_mode):
yield definition
+ else:
+ log.debug("Skipping %r: no Airflow DAG markers found",
definition)
Review Comment:
Done
--
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]