goyaladitay11 commented on issue #59120:
URL: https://github.com/apache/airflow/issues/59120#issuecomment-5928654005
Hey, I've been reading through scheduler_job_runner.py and the issue
description. Here's what I've understood and how I'm planning to approach this:
What's happening now
The dagrun creation loop processes multiple DAGs inside a single database
transaction. When one DAG raises an exception (say a DB constraint violation or
a bad timetable evaluation), the current except Exception: ... continue handler
catches the Python-level exception — but the underlying SQLAlchemy session is
now in a broken state. The transaction is poisoned.
┌─────────────────────────────────────┐
│ Single DB Transaction │
│ │
│ ┌─ DAG_A ──── create dagrun ✅ ──┐ │
│ ├─ DAG_B ──── DB error 💥 ──────┤ │ ← exception caught, continue...
│ ├─ DAG_C ──── uses same session ┤ │ ← but session is now INVALID
│ └─ DAG_D ──── also fails ───────┘ │ ← all subsequent DAGs fail
│ │
│ session.commit() 💥 │ ← entire batch lost
└─────────────────────────────────────┘
The except Exception: continue only catches the Python exception. It doesn't
roll back the failed DB operation, so the session accumulates dirty state.
Every DAG processed after the failure either silently produces corrupt data or
raises another exception.
Why it's a real problem
In production with hundreds of DAGs, one misconfigured DAG (bad timetable,
broken serialization, constraint violation) can prevent ALL other healthy DAGs
from getting new DagRuns created. The scheduler silently stops scheduling
everything in that batch.
Proposed fix
The issue description mentions two approaches — I'm leaning toward
savepoints because they're less disruptive:
┌──────────────────────────────────────┐
│ Outer DB Transaction │
│ │
│ ┌─ SAVEPOINT ──────────────────┐ │
│ │ DAG_A ── create dagrun ✅ │ │ ← savepoint released
│ └──────────────────────────────┘ │
│ ┌─ SAVEPOINT ──────────────────┐ │
│ │ DAG_B ── DB error 💥 │ │ ← savepoint ROLLED BACK
│ └──────────────────────────────┘ │
│ ┌─ SAVEPOINT ──────────────────┐ │
│ │ DAG_C ── create dagrun ✅ │ │ ← unaffected, works fine
│ └──────────────────────────────┘ │
│ │
│ session.commit() ✅ │ ← A and C committed, B skipped
└──────────────────────────────────────┘
In SQLAlchemy, this would be session.begin_nested() which creates a
SAVEPOINT. If the operation inside fails, we roll back only that savepoint,
leaving the rest of the transaction intact.
Rough sketch of the approach:
python
for dag_model in dag_models:
try:
with session.begin_nested(): # SAVEPOINT
# create dagrun for this dag
...
except Exception:
self.log.exception(
"Failed to create DagRun for DAG '%s', skipping",
dag_model.dag_id,
)
# savepoint already rolled back, session is clean
continue
My plan
Find the exact loop in scheduler_job_runner.py where DagRuns are created in
batch
Wrap each per-DAG operation in session.begin_nested() to create a savepoint
Update the except Exception handler to log and continue cleanly
Write a test that verifies: one bad DAG in a batch doesn't prevent other
DAGs from getting their runs created
Run the scheduler test suite locally to make sure nothing else breaks
Will have a PR up soon. Let me know if the savepoint approach aligns with
what you had in mind, or if you'd prefer the smaller-transactions route instead.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]