From 1b953213783d3f876377e6cea8e5a3d65a5ad4cc Mon Sep 17 00:00:00 2001 From: Kalin Stoyanov Date: Wed, 24 Jun 2026 16:49:14 +0300 Subject: [PATCH 1/3] Serialize dag if it is not serialized even if it does not have task instances --- .../src/airflow/models/serialized_dag.py | 2 +- .../tests/unit/models/test_serialized_dag.py | 21 +++++++++++++++++++ 2 files changed, 22 insertions(+), 1 deletion(-) diff --git a/airflow-core/src/airflow/models/serialized_dag.py b/airflow-core/src/airflow/models/serialized_dag.py index f3f5b78309680..425841543657c 100644 --- a/airflow-core/src/airflow/models/serialized_dag.py +++ b/airflow-core/src/airflow/models/serialized_dag.py @@ -723,7 +723,7 @@ def write_dag( ) ) - if dag_version and not has_task_instances: + if dag_version and not has_task_instances and serialized_dag_hash is not None: # This is for dynamic DAGs that the hashes changes often. We should update # the serialized dag, the dag_version and the dag_code instead of a new version # if the dag_version is not associated with any task instances diff --git a/airflow-core/tests/unit/models/test_serialized_dag.py b/airflow-core/tests/unit/models/test_serialized_dag.py index fd3568f34619c..e81b347c9fd04 100644 --- a/airflow-core/tests/unit/models/test_serialized_dag.py +++ b/airflow-core/tests/unit/models/test_serialized_dag.py @@ -502,6 +502,27 @@ def test_new_dag_versions_are_created_if_there_is_a_dagrun(self, dag_maker, sess assert session.scalar(select(func.count()).select_from(DagVersion)) == 2 assert session.scalar(select(func.count()).select_from(SDM)) == 2 + def test_new_dag_version_is_created_when_version_exists_but_serialized_dag_row_missing( + self, dag_maker, session + ): + with dag_maker("dag1") as dag: + PythonOperator(task_id="task1", python_callable=lambda: None) + assert session.scalar(select(func.count()).select_from(SDM)) == 1 + assert session.scalar(select(func.count()).select_from(DagVersion)) == 1 + + # Simulate the broken state: DagVersion present, serialized_dag row gone. + session.execute(delete(SDM).where(SDM.dag_id == dag.dag_id)) + session.flush() + assert session.scalar(select(func.count()).select_from(SDM)) == 0 + assert session.scalar(select(func.count()).select_from(DagVersion)) == 1 + + SDM.write_dag(LazyDeserializedDAG.from_dag(dag), bundle_name="dag_maker") + + # A new version and serialized row must have been created. + assert session.scalar(select(func.count()).select_from(DagVersion)) == 2 + assert session.scalar(select(func.count()).select_from(SDM)) == 1 + assert SDM.get(dag.dag_id, session=session) is not None + def test_example_dag_sorting_serialised_dag(self, session): """ This test asserts if different dag ids -- simple or complex, can be sorted From 432c1d9ad676a40e92125a9a106bad780a3b3d76 Mon Sep 17 00:00:00 2001 From: Kalin Stoyanov Date: Thu, 25 Jun 2026 14:45:07 +0300 Subject: [PATCH 2/3] update unit test --- .../tests/unit/models/test_serialized_dag.py | 34 +++++-------------- 1 file changed, 9 insertions(+), 25 deletions(-) diff --git a/airflow-core/tests/unit/models/test_serialized_dag.py b/airflow-core/tests/unit/models/test_serialized_dag.py index e81b347c9fd04..c6c4319272c46 100644 --- a/airflow-core/tests/unit/models/test_serialized_dag.py +++ b/airflow-core/tests/unit/models/test_serialized_dag.py @@ -502,27 +502,6 @@ def test_new_dag_versions_are_created_if_there_is_a_dagrun(self, dag_maker, sess assert session.scalar(select(func.count()).select_from(DagVersion)) == 2 assert session.scalar(select(func.count()).select_from(SDM)) == 2 - def test_new_dag_version_is_created_when_version_exists_but_serialized_dag_row_missing( - self, dag_maker, session - ): - with dag_maker("dag1") as dag: - PythonOperator(task_id="task1", python_callable=lambda: None) - assert session.scalar(select(func.count()).select_from(SDM)) == 1 - assert session.scalar(select(func.count()).select_from(DagVersion)) == 1 - - # Simulate the broken state: DagVersion present, serialized_dag row gone. - session.execute(delete(SDM).where(SDM.dag_id == dag.dag_id)) - session.flush() - assert session.scalar(select(func.count()).select_from(SDM)) == 0 - assert session.scalar(select(func.count()).select_from(DagVersion)) == 1 - - SDM.write_dag(LazyDeserializedDAG.from_dag(dag), bundle_name="dag_maker") - - # A new version and serialized row must have been created. - assert session.scalar(select(func.count()).select_from(DagVersion)) == 2 - assert session.scalar(select(func.count()).select_from(SDM)) == 1 - assert SDM.get(dag.dag_id, session=session) is not None - def test_example_dag_sorting_serialised_dag(self, session): """ This test asserts if different dag ids -- simple or complex, can be sorted @@ -896,10 +875,12 @@ def __init__(self, *, task_id: str, **kwargs): # Hashes should be identical assert hash_1 == hash_2, "Hashes should be identical when dicts are sorted consistently" - def test_dynamic_dag_update_preserves_null_check(self, dag_maker, session): + def test_new_dag_version_is_created_when_version_exists_but_serialized_dag_row_missing( + self, dag_maker, session + ): """ Test that dynamic DAG update gracefully handles case where SerializedDagModel doesn't exist. - This preserves the null-check fix from PR #56422 and tests the direct UPDATE path. + This adds the serialized dag entry """ with dag_maker(dag_id="test_missing_serdag", serialized=True, session=session) as dag: EmptyOperator(task_id="task1") @@ -927,7 +908,7 @@ def test_dynamic_dag_update_preserves_null_check(self, dag_maker, session): # Verify no SerializedDagModel exists assert SDM.get("test_missing_serdag", session=session) is None - # Try to update - should return False gracefully (not crash) + # Try to update - should create the serialized dag result = SDM.write_dag( dag=lazy_dag, bundle_name="test_bundle", @@ -935,8 +916,11 @@ def test_dynamic_dag_update_preserves_null_check(self, dag_maker, session): min_update_interval=None, session=session, ) + session.commit() - assert result is False # Should return False when SerializedDagModel is missing + serialized_dag = session.scalar(select(SDM).where(SDM.dag_id == "test_missing_serdag").limit(1)) + assert serialized_dag is not None + assert result is True def test_dynamic_dag_update_success(self, dag_maker, session): """ From eb448d00fc52c23fa5d36669e26d06b4483c81d7 Mon Sep 17 00:00:00 2001 From: Kalin Stoyanov Date: Wed, 5 Aug 2026 10:52:42 +0300 Subject: [PATCH 3/3] match serialized dag entry to latest dag version during prefetch --- .../src/airflow/models/serialized_dag.py | 42 ++-- .../tests/unit/models/test_serialized_dag.py | 216 +++++++++++++++++- 2 files changed, 231 insertions(+), 27 deletions(-) diff --git a/airflow-core/src/airflow/models/serialized_dag.py b/airflow-core/src/airflow/models/serialized_dag.py index 425841543657c..ac049abf12390 100644 --- a/airflow-core/src/airflow/models/serialized_dag.py +++ b/airflow-core/src/airflow/models/serialized_dag.py @@ -547,28 +547,6 @@ def _prefetch_dag_write_metadata( if not dag_id_list: return {} - # Fetch the serialized_dag (last_updated, dag_hash) of the latest DagVersion per dag_id, - # ordering by version_number so it stays consistent with the DagVersion picked by dv_subq. - sd_subq = ( - select( - cls.dag_id.label("dag_id"), - cls.last_updated.label("last_updated"), - cls.dag_hash.label("dag_hash"), - func.row_number() - .over(partition_by=cls.dag_id, order_by=DagVersion.version_number.desc()) - .label("rn"), - ) - .join(DagVersion, cls.dag_version_id == DagVersion.id) - .where(cls.dag_id.in_(dag_id_list)) - .subquery() - ) - sd_rows = session.execute( - select(sd_subq.c.dag_id, sd_subq.c.last_updated, sd_subq.c.dag_hash).where(sd_subq.c.rn == 1) - ).all() - sd_by_dag_id: dict[str, tuple[datetime, str]] = { - row.dag_id: (row.last_updated, row.dag_hash) for row in sd_rows - } - # Fetch latest DagVersion per dag_id, ordering by version_number to match write_dag. dv_subq = ( select( @@ -586,6 +564,23 @@ def _prefetch_dag_write_metadata( ).all() dv_by_dag_id: dict[str, DagVersion] = {dv.dag_id: dv for dv in dag_versions} + # Fetch the serialized_dag (last_updated, dag_hash) of the latest DagVersion per dag_id, + # outer join with dv_subq so None is set when latest dag version has no serialized entry. + sd_subq = ( + select( + dv_subq.c.dag_id, + cls.last_updated.label("last_updated"), + cls.dag_hash.label("dag_hash"), + ) + .where(dv_subq.c.rn == 1) + .outerjoin(cls, cls.dag_version_id == dv_subq.c.id) + .subquery() + ) + sd_rows = session.execute(select(sd_subq.c.dag_id, sd_subq.c.last_updated, sd_subq.c.dag_hash)).all() + sd_by_dag_id: dict[str, tuple[datetime, str]] = { + row.dag_id: (row.last_updated, row.dag_hash) for row in sd_rows + } + return { dag_id: DagWriteMetadata( last_updated=sd_by_dag_id[dag_id][0] if dag_id in sd_by_dag_id else None, @@ -727,6 +722,9 @@ def write_dag( # This is for dynamic DAGs that the hashes changes often. We should update # the serialized dag, the dag_version and the dag_code instead of a new version # if the dag_version is not associated with any task instances + # exception is when the *latest* dag version has no corresponding serialized dag + # which is denoted by serialized_dag_hash == None, then fall through to write + new_serialized_dag = cls(dag) # Use direct UPDATE to avoid loading the full serialized DAG diff --git a/airflow-core/tests/unit/models/test_serialized_dag.py b/airflow-core/tests/unit/models/test_serialized_dag.py index c6c4319272c46..c79f8c1f1a5b0 100644 --- a/airflow-core/tests/unit/models/test_serialized_dag.py +++ b/airflow-core/tests/unit/models/test_serialized_dag.py @@ -34,7 +34,7 @@ from airflow.models.dag import DagModel from airflow.models.dag_version import DagVersion from airflow.models.deadline_alert import DeadlineAlert as DAM -from airflow.models.serialized_dag import SerializedDagModel as SDM +from airflow.models.serialized_dag import DagWriteMetadata, SerializedDagModel as SDM from airflow.providers.standard.operators.bash import BashOperator from airflow.providers.standard.operators.empty import EmptyOperator from airflow.providers.standard.operators.python import PythonOperator @@ -688,6 +688,36 @@ def test_prefetch_dag_write_metadata_returns_latest_version(self, dag_maker, ses assert metadata.dag_version is not None assert metadata.dag_version.version_number == 2 + def test_prefetch_dag_write_metadata_latest_version_without_sdm_entry(self, dag_maker, session): + """When the latest DagVersion has no SDM entry, last_updated and dag_hash are None.""" + with dag_maker("prefetch_no_sdm_dag") as dag: + EmptyOperator(task_id="task1") + + v1 = session.scalar(select(DagVersion).where(DagVersion.dag_id == dag.dag_id)) + assert v1 is not None + assert v1.version_number == 1 + + # v1 has an SDM entry (written by dag_maker); create v2 directly without one. + v2 = DagVersion.write_dag( + dag_id=dag.dag_id, + bundle_name="dag_maker", + bundle_version="v2", + session=session, + ) + session.flush() + + assert v2.version_number == 2 + + result = SDM._prefetch_dag_write_metadata([dag.dag_id], session=session) + metadata = result[dag.dag_id] + + assert metadata.dag_version is not None + assert metadata.dag_version.version_number == 2 + assert metadata.dag_version.dag_id == dag.dag_id + assert metadata.dag_version.bundle_version == "v2" + assert metadata.last_updated is None + assert metadata.dag_hash is None + def test_new_dag_version_created_when_bundle_name_changes_and_hash_unchanged(self, dag_maker, session): """Test that new dag_version is created if bundle_name changes but DAG is unchanged.""" # Create and write initial DAG @@ -879,8 +909,7 @@ def test_new_dag_version_is_created_when_version_exists_but_serialized_dag_row_m self, dag_maker, session ): """ - Test that dynamic DAG update gracefully handles case where SerializedDagModel doesn't exist. - This adds the serialized dag entry + Test that dynamic DAG update creates a SerializedDagModel if it doesn't exist. """ with dag_maker(dag_id="test_missing_serdag", serialized=True, session=session) as dag: EmptyOperator(task_id="task1") @@ -908,7 +937,7 @@ def test_new_dag_version_is_created_when_version_exists_but_serialized_dag_row_m # Verify no SerializedDagModel exists assert SDM.get("test_missing_serdag", session=session) is None - # Try to update - should create the serialized dag + # Try to update - should create a new serialized dag row under the latest DagVersion result = SDM.write_dag( dag=lazy_dag, bundle_name="test_bundle", @@ -918,9 +947,186 @@ def test_new_dag_version_is_created_when_version_exists_but_serialized_dag_row_m ) session.commit() - serialized_dag = session.scalar(select(SDM).where(SDM.dag_id == "test_missing_serdag").limit(1)) + assert result is True + latest_version = session.scalar( + select(DagVersion) + .where(DagVersion.dag_id == "test_missing_serdag") + .order_by(DagVersion.version_number.desc()) + .limit(1) + ) + assert latest_version is not None + serialized_dag = SDM.get("test_missing_serdag", session=session) assert serialized_dag is not None + assert serialized_dag.dag_version_id == latest_version.id + + def test_write_dag_returns_false_when_update_finds_no_serialized_row(self, dag_maker, session): + """ + Test that write_dag returns False when the UPDATE path finds no serialized row. + + This exercises the rowcount == 0 branch: _prefetched carries a non-None dag_hash + (so the dynamic-update UPDATE path is taken), but the SDM row has been deleted + between prefetch and execution, so the UPDATE affects 0 rows. + """ + with dag_maker(dag_id="test_rowcount_zero", serialized=True, session=session) as dag: + EmptyOperator(task_id="task1") + + lazy_dag = LazyDeserializedDAG.from_dag(dag) + SDM.write_dag( + dag=lazy_dag, + bundle_name="test_bundle", + bundle_version=None, + session=session, + ) + session.commit() + + dag_version = session.scalar( + select(DagVersion) + .where(DagVersion.dag_id == "test_rowcount_zero") + .order_by(DagVersion.version_number.desc()) + .limit(1) + ) + assert dag_version is not None + + # Build prefetched metadata that looks like the SDM row exists (non-None hash), + # then delete the actual row so the UPDATE finds nothing. + stale_hash = session.scalar(select(SDM.dag_hash).where(SDM.dag_id == "test_rowcount_zero")) + assert stale_hash is not None + session.execute(delete(SDM).where(SDM.dag_id == "test_rowcount_zero")) + session.commit() + + prefetched = DagWriteMetadata( + last_updated=None, + dag_hash=stale_hash, + dag_version=dag_version, + ) + # Change the dag so the hash differs, forcing the dynamic-update branch. + from airflow.sdk.definitions.dag import DAG as SdkDAG + + new_dag_obj = SdkDAG(dag_id="test_rowcount_zero", schedule=None) + with new_dag_obj: + EmptyOperator(task_id="task1") + EmptyOperator(task_id="task2") + new_lazy = LazyDeserializedDAG.from_dag(new_dag_obj) + + result = SDM.write_dag( + dag=new_lazy, + bundle_name="test_bundle", + bundle_version=None, + min_update_interval=None, + session=session, + _prefetched=prefetched, + ) + + assert result is False + assert SDM.get("test_rowcount_zero", session=session) is None + + def test_latest_dag_version_with_no_serialized_row_and_changed_hash_heals(self, dag_maker, session): + """ + v1 has a serialized row, v2 is a bare DagVersion (no serialized row), and + the incoming Dag hash differs from v1's. The v2 dag should be serialized successfully + """ + with dag_maker(dag_id="test_v2_bare_changed_hash", serialized=True, session=session) as dag: + EmptyOperator(task_id="task1") + + lazy_dag = LazyDeserializedDAG.from_dag(dag) + SDM.write_dag(dag=lazy_dag, bundle_name="test_bundle", bundle_version=None, session=session) + session.commit() + + # Create a bare v2 DagVersion with no corresponding serialized row. + DagVersion.write_dag(dag_id="test_v2_bare_changed_hash", bundle_name="test_bundle", session=session) + session.commit() + + v2 = session.scalar( + select(DagVersion) + .where(DagVersion.dag_id == "test_v2_bare_changed_hash") + .order_by(DagVersion.version_number.desc()) + .limit(1) + ) + assert v2 is not None + assert v2.version_number == 2 + + # A changed dag produces a different hash — bypasses the unchanged-hash short-circuit. + from airflow.sdk.definitions.dag import DAG as SdkDAG + + changed_dag = SdkDAG(dag_id="test_v2_bare_changed_hash", schedule=None) + with changed_dag: + EmptyOperator(task_id="task1") + EmptyOperator(task_id="task2") + changed_lazy = LazyDeserializedDAG.from_dag(changed_dag) + + result = SDM.write_dag( + dag=changed_lazy, + bundle_name="test_bundle", + bundle_version=None, + min_update_interval=None, + session=session, + ) + session.commit() + assert result is True + # The code falls through to DagVersion.write_dag, creating a new v3 with the serialized + # row attached to it — v2 stays bare. Assert the latest version now has a serialized row. + latest = session.scalar( + select(DagVersion) + .where(DagVersion.dag_id == "test_v2_bare_changed_hash") + .order_by(DagVersion.version_number.desc()) + .limit(1) + ) + assert latest is not None + assert latest.version_number == 3 + assert session.scalar(select(SDM).where(SDM.dag_version_id == latest.id)) is not None + + def test_latest_dag_version_with_no_serialized_row_and_unchanged_hash_heals(self, dag_maker, session): + """ + v1 has a serialized row, v2 is a bare DagVersion (no serialized row), and + the incoming Dag hash matches v1's. The v2 dag should be serialized successfully + """ + with dag_maker(dag_id="test_v2_bare_unchanged_hash", serialized=True, session=session) as dag: + EmptyOperator(task_id="task1") + + lazy_dag = LazyDeserializedDAG.from_dag(dag) + SDM.write_dag(dag=lazy_dag, bundle_name="test_bundle", bundle_version=None, session=session) + session.commit() + + v1_hash = session.scalar(select(SDM.dag_hash).where(SDM.dag_id == "test_v2_bare_unchanged_hash")) + assert v1_hash is not None + + # Create a bare v2 DagVersion with no corresponding serialized row. + DagVersion.write_dag(dag_id="test_v2_bare_unchanged_hash", bundle_name="test_bundle", session=session) + session.commit() + + v2 = session.scalar( + select(DagVersion) + .where(DagVersion.dag_id == "test_v2_bare_unchanged_hash") + .order_by(DagVersion.version_number.desc()) + .limit(1) + ) + assert v2 is not None + assert v2.version_number == 2 + + # Same dag — same hash as v1. Without the fix the prefetch returns v1's hash (inner + # join excludes v2), the unchanged-hash guard fires, and v2 never gets a serialized row. + result = SDM.write_dag( + dag=lazy_dag, + bundle_name="test_bundle", + bundle_version=None, + min_update_interval=None, + session=session, + ) + session.commit() + + assert result is True + # The code falls through to DagVersion.write_dag, creating a new v3 with the serialized + # row attached to it — v2 stays bare. Assert the latest version now has a serialized row. + latest = session.scalar( + select(DagVersion) + .where(DagVersion.dag_id == "test_v2_bare_unchanged_hash") + .order_by(DagVersion.version_number.desc()) + .limit(1) + ) + assert latest is not None + assert latest.version_number == 3 + assert session.scalar(select(SDM).where(SDM.dag_version_id == latest.id)) is not None def test_dynamic_dag_update_success(self, dag_maker, session): """