Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 1 addition & 2 deletions airflow-core/src/airflow/dag_processing/manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,6 @@
from airflow.dag_processing.bundles.manager import DagBundlesManager
from airflow.dag_processing.collection import update_dag_parsing_results_in_db
from airflow.dag_processing.processor import DagFileParsingResult, DagFileProcessorProcess
from airflow.exceptions import AirflowException
from airflow.models.asset import remove_references_to_deleted_dags
from airflow.models.dag import DagModel
from airflow.models.dagbag import DagPriorityParsingRequest
Expand Down Expand Up @@ -853,7 +852,7 @@ def _refresh_dag_bundles(self, known_files: dict[str, set[DagFileInfo]]):
try:
bundle.initialize()
any_refreshed = True
except AirflowException as e:
except Exception as e:
self.log.exception("Error initializing bundle %s: %s", bundle.name, e)
continue
try:
Expand Down
33 changes: 33 additions & 0 deletions airflow-core/tests/unit/dag_processing/test_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -3251,6 +3251,39 @@ def test_refresh_dag_bundles_get_bundle_state_failure_skips_bundle(self):

bundle.refresh.assert_not_called()

def test_refresh_dag_bundles_initialize_non_airflow_exception_skips_bundle(self):
"""
A bundle whose initialize() raises a non AirflowException must be skipped, not
left to propagate and abort refresh for every other bundle.
"""
manager = DagFileProcessorManager(max_runs=1)
failing_bundle = self._make_refresh_bundle()
failing_bundle.name = "failing_bundle"
failing_bundle.is_initialized = False
failing_bundle.initialize.side_effect = FileNotFoundError("Repository path not found")

healthy_bundle = self._make_refresh_bundle()
healthy_bundle.name = "healthy_bundle"

manager._dag_bundles = [failing_bundle, healthy_bundle]
manager._force_refresh_bundles = set()

with (
mock.patch.object(
manager, "get_bundle_state", return_value=BundleState(last_refreshed=None, version=None)
),
mock.patch.object(manager, "update_bundle_state"),
mock.patch.object(manager, "_find_files_in_bundle", return_value=[]),
mock.patch.object(manager, "deactivate_deleted_dags"),
mock.patch.object(manager, "clear_orphaned_import_errors"),
mock.patch.object(manager, "handle_removed_files"),
mock.patch.object(manager, "_resort_file_queue"),
mock.patch.object(manager, "_add_new_files_to_queue"),
):
manager._refresh_dag_bundles({})

healthy_bundle.refresh.assert_called_once()

def test_refresh_dag_bundles_update_bundle_state_failure_still_scans_files(self):
"""A failure in update_bundle_state() logs but does not skip file scanning.

Expand Down