Skip to content

Commit ca2d9ba

Browse files
leossantosLee-W
andauthored
Apply suggestions from code review
Co-authored-by: Wei Lee <weilee.rx@gmail.com>
1 parent 411ef4c commit ca2d9ba

3 files changed

Lines changed: 14 additions & 14 deletions

File tree

airflow/models/dag.py

Lines changed: 12 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -4070,9 +4070,9 @@ def dags_needing_dagruns(cls, session: Session) -> tuple[Query, dict[str, tuple[
40704070
you should ensure that any scheduling decisions are made in a single transaction -- as soon as the
40714071
transaction is committed it will be unlocked.
40724072
4073-
For dataset-triggered scheduling, DAGs that have ``DatasetDagRunQueue`` rows but no matching
4073+
For dataset-triggered scheduling, Dags that have ``DatasetDagRunQueue`` rows but no matching
40744074
``SerializedDagModel`` row are omitted from the returned ``dataset_triggered_dag_info`` until
4075-
serialization exists; queue rows are **not** deleted here so the scheduler can re-evaluate on a
4075+
serialization exists; DDRQs are **not** deleted here so the scheduler can re-evaluate on a
40764076
later run.
40774077
"""
40784078
from airflow.models.serialized_dag import SerializedDagModel
@@ -4100,16 +4100,16 @@ def dag_ready(dag_id: str, cond: BaseDataset, statuses: dict) -> bool | None:
41004100
select(SerializedDagModel).where(SerializedDagModel.dag_id.in_(dag_statuses.keys()))
41014101
).all()
41024102
ser_dag_ids = {s.dag_id for s in ser_dags}
4103-
missing_from_serialized = set(by_dag.keys()) - ser_dag_ids
4104-
if missing_from_serialized:
4105-
log.debug(
4106-
"DAGs in DDRQ but missing SerializedDagModel "
4107-
"(skippingcondition cannot be evaluated): %s",
4108-
sorted(missing_from_serialized),
4109-
)
4110-
for dag_id in missing_from_serialized:
4111-
del by_dag[dag_id]
4112-
del dag_statuses[dag_id]
4103+
missing_from_serialized = set(by_dag.keys()) - ser_dag_ids
4104+
if missing_from_serialized:
4105+
log.debug(
4106+
"Dags have queued dataset events (DDRQs), but are not found in the serialized_dag table."
4107+
" — skipping Dag run creation: %s",
4108+
sorted(missing_from_serialized),
4109+
)
4110+
for dag_id in missing_from_serialized:
4111+
del by_dag[dag_id]
4112+
del dag_statuses[dag_id]
41134113
for ser_dag in ser_dags:
41144114
dag_id = ser_dag.dag_id
41154115
statuses = dag_statuses[dag_id]

newsfragments/63546.bugfix.rst

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1 +1 @@
1-
Fix premature dataset-triggered DagRuns when ``SerializedDagModel`` was missing while ``DatasetDagRunQueue`` still had rows for that DAG; queue entries are kept for the next evaluation.
1+
Fix premature dataset-triggered DagRuns when ``SerializedDagModel`` was missing while ``DatasetDagRunQueue`` still had rows for that Dag; queue entries are kept for the next evaluation.

tests/models/test_dag.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3066,7 +3066,7 @@ def test_dags_needing_dagruns_datasets(self, dag_maker, session):
30663066
assert dag_models == [dag_model]
30673067

30683068
def test_dags_needing_dagruns_skips_ddrq_when_serialized_dag_missing(self, session, caplog):
3069-
"""DDRQ rows for a dag_id without SerializedDagModel must be skipped (no dataset_triggered info).
3069+
"""DDRQ rows for a Dag without SerializedDagModel must be skipped (no dataset_triggered info).
30703070
30713071
Rows must remain in ``dataset_dag_run_queue`` so the scheduler can re-evaluate on a later
30723072
heartbeat once ``SerializedDagModel`` exists (``dags_needing_dagruns`` only drops them from

0 commit comments

Comments
 (0)