dilnazanlid commented on code in PR #74020:
URL: https://github.com/apache/airflow/pull/74020#discussion_r4164450813
##########
airflow-core/src/airflow/dag_processing/manager.py:
##########
@@ -970,56 +979,43 @@ def _refresh_dag_bundles(self, known_files: dict[str,
set[DagFileInfo]]):
self._resort_file_queue()
self._add_new_files_to_queue(known_files=known_files)
- def _find_files_in_bundle(self, bundle: BaseDagBundle) -> list[Path]:
- """Get relative paths for dag files from bundle dir."""
- # Build up a list of Python files that could contain DAGs
- self.log.info("Searching for files in %s at %s", bundle.name,
bundle.path)
- rel_paths = [
- Path(x).relative_to(bundle.path)
- for x in list_py_file_paths(bundle.path,
safe_mode=self.dag_discovery_safe_mode)
- ]
+ def _find_files_in_bundle(self, bundle: BaseDagBundle) -> set[DagFileInfo]:
+ """
+ List the files to parse in a bundle through its importers.
+
+ A file holding several Dag definitions (a zip archive, for instance)
is parsed as one.
+ """
+ self.log.info("Searching for Dag definitions in %s at %s",
bundle.name, bundle.path)
+ registry = get_importer_registry(bundle.name)
Review Comment:
There is the ADR to move the whole process model to the importers in #73457
but I included `warm_importers()` method into the registry to have it here. As
the importers are experimental in 3.4. and there are changes planned in the ADR
- I think it should be okay.
cc @uranusjr @jason810496
##########
airflow-core/src/airflow/dag_processing/manager.py:
##########
@@ -970,56 +979,43 @@ def _refresh_dag_bundles(self, known_files: dict[str,
set[DagFileInfo]]):
self._resort_file_queue()
self._add_new_files_to_queue(known_files=known_files)
- def _find_files_in_bundle(self, bundle: BaseDagBundle) -> list[Path]:
- """Get relative paths for dag files from bundle dir."""
- # Build up a list of Python files that could contain DAGs
- self.log.info("Searching for files in %s at %s", bundle.name,
bundle.path)
- rel_paths = [
- Path(x).relative_to(bundle.path)
- for x in list_py_file_paths(bundle.path,
safe_mode=self.dag_discovery_safe_mode)
- ]
+ def _find_files_in_bundle(self, bundle: BaseDagBundle) -> set[DagFileInfo]:
+ """
+ List the files to parse in a bundle through its importers.
+
+ A file holding several Dag definitions (a zip archive, for instance)
is parsed as one.
+ """
+ self.log.info("Searching for Dag definitions in %s at %s",
bundle.name, bundle.path)
+ registry = get_importer_registry(bundle.name)
+ definition_locs: defaultdict[Path, set[str]] = defaultdict(set)
Review Comment:
Thanks, I adde them in the PR description. Also, found another user visible
change - zip archives without explicit `.zip` extension (e.g. `.egg`) are no
longer discovered so they may be ignored unless explicitly added into the
importer config `extensions` field for ZipImporter.
##########
airflow-core/src/airflow/dag_processing/manager.py:
##########
@@ -970,56 +979,43 @@ def _refresh_dag_bundles(self, known_files: dict[str,
set[DagFileInfo]]):
self._resort_file_queue()
self._add_new_files_to_queue(known_files=known_files)
- def _find_files_in_bundle(self, bundle: BaseDagBundle) -> list[Path]:
- """Get relative paths for dag files from bundle dir."""
- # Build up a list of Python files that could contain DAGs
- self.log.info("Searching for files in %s at %s", bundle.name,
bundle.path)
- rel_paths = [
- Path(x).relative_to(bundle.path)
- for x in list_py_file_paths(bundle.path,
safe_mode=self.dag_discovery_safe_mode)
- ]
+ def _find_files_in_bundle(self, bundle: BaseDagBundle) -> set[DagFileInfo]:
+ """
+ List the files to parse in a bundle through its importers.
+
+ A file holding several Dag definitions (a zip archive, for instance)
is parsed as one.
+ """
+ self.log.info("Searching for Dag definitions in %s at %s",
bundle.name, bundle.path)
+ registry = get_importer_registry(bundle.name)
+ definition_locs: defaultdict[Path, set[str]] = defaultdict(set)
+ for _, item in registry.list_dag_definitions(bundle,
safe_mode=self.dag_discovery_safe_mode):
Review Comment:
Thanks, now python importer `list_dag_definitions` guards each file the way
ZipImporter does - unreadable file is yielded as a `DagImportError`, so it's
still queued, the error surfaces, while the rest of the bundle is listed. I
feel like it is suitable as the whole `might_contain_dag` method is defined on
the importer level - it should be responsible for using.
--
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]