From 31e9214b681efea30307b3ee4604d856b4e4c59b Mon Sep 17 00:00:00 2001 From: Michael Petro Date: Tue, 7 Feb 2023 14:27:05 -0500 Subject: [PATCH 1/4] dag processing manager, dag serialization delete, only filter on dag folder when dag processing is in a standalone processor --- airflow/dag_processing/manager.py | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/airflow/dag_processing/manager.py b/airflow/dag_processing/manager.py index 2246e7f3a7262..38fd4048ffcaf 100644 --- a/airflow/dag_processing/manager.py +++ b/airflow/dag_processing/manager.py @@ -764,9 +764,13 @@ def _refresh_dag_dir(self): else: dag_filelocs.append(fileloc) + processor_subdir = None + if self.standalone_dag_processor: + processor_subdir = self.get_dag_directory() + SerializedDagModel.remove_deleted_dags( alive_dag_filelocs=dag_filelocs, - processor_subdir=self.get_dag_directory(), + processor_subdir=processor_subdir, ) DagModel.deactivate_deleted_dags(self._file_paths) From 78ed2d7be316be71a9e27a2af8b36323131ec4d3 Mon Sep 17 00:00:00 2001 From: Michael Petro Date: Wed, 8 Feb 2023 09:31:15 -0500 Subject: [PATCH 2/4] Remove redundant change --- airflow/dag_processing/manager.py | 6 +----- 1 file changed, 1 insertion(+), 5 deletions(-) diff --git a/airflow/dag_processing/manager.py b/airflow/dag_processing/manager.py index 38fd4048ffcaf..2246e7f3a7262 100644 --- a/airflow/dag_processing/manager.py +++ b/airflow/dag_processing/manager.py @@ -764,13 +764,9 @@ def _refresh_dag_dir(self): else: dag_filelocs.append(fileloc) - processor_subdir = None - if self.standalone_dag_processor: - processor_subdir = self.get_dag_directory() - SerializedDagModel.remove_deleted_dags( alive_dag_filelocs=dag_filelocs, - processor_subdir=processor_subdir, + processor_subdir=self.get_dag_directory(), ) DagModel.deactivate_deleted_dags(self._file_paths) From 50dc0c78fa1fa017d96f4b3f409188ee82fe26e3 Mon Sep 17 00:00:00 2001 From: Michael Petro Date: Wed, 8 Feb 2023 09:31:45 -0500 Subject: [PATCH 3/4] remove deleted dags from serialized dags, fix None check --- airflow/models/serialized_dag.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/airflow/models/serialized_dag.py b/airflow/models/serialized_dag.py index 3e1ec0c70cce3..53e5e2ccbd223 100644 --- a/airflow/models/serialized_dag.py +++ b/airflow/models/serialized_dag.py @@ -251,7 +251,7 @@ def remove_deleted_dags( cls.fileloc_hash.notin_(alive_fileloc_hashes), cls.fileloc.notin_(alive_dag_filelocs), or_( - cls.processor_subdir is None, + cls.processor_subdir.is_(None), cls.processor_subdir == processor_subdir, ), ) From 91fe276abd3d5feffec3b403b860a0261c41e582 Mon Sep 17 00:00:00 2001 From: Michael Petro Date: Wed, 8 Feb 2023 09:52:44 -0500 Subject: [PATCH 4/4] serialized_dag tests, add processor_subdir to remove_deleted_dags test --- tests/models/test_serialized_dag.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/models/test_serialized_dag.py b/tests/models/test_serialized_dag.py index d4e7d4e92316f..b425cd8f658cd 100644 --- a/tests/models/test_serialized_dag.py +++ b/tests/models/test_serialized_dag.py @@ -169,7 +169,7 @@ def test_remove_dags_by_filepath(self): # remove repeated files for those DAGs that define multiple dags in the same file (set comprehension) example_dag_files = list({dag.fileloc for dag in filtered_example_dags_list}) example_dag_files.remove(dag_removed_by_file.fileloc) - SDM.remove_deleted_dags(example_dag_files) + SDM.remove_deleted_dags(example_dag_files, processor_subdir="/tmp/test") assert not SDM.has_dag(dag_removed_by_file.dag_id) def test_bulk_sync_to_db(self):