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, ), ) 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):