SameerMesiah97 commented on code in PR #73728:
URL: https://github.com/apache/airflow/pull/73728#discussion_r4108837785
##########
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:
I think this can instantiate the same configured importer more than once
when it supports multiple extensions. For example, if one spec is registered
for both `.yaml` and `.yml`, `list(self._extension_specs.values())` contains
that spec twice. The first call to `_materialise_spec` removes both entries,
but the second spec is still in the snapshot. It then creates another importer
and appends it to `_ordered_importers`, even though there are no extensions
left for it to take over.
That means `list_dag_definitions()` could call the importer twice and yield
duplicate definitions. Could we materialise each spec only once and add a test
for an importer registered with multiple extensions?
--
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]