ukiband opened a new issue, #71923:
URL: https://github.com/apache/airflow/issues/71923

   ### Apache Airflow version
   
   3.1.8
   
   ### What happened?
   
   During a rolling restart of Celery workers (deployment), a sensor running in 
`mode="reschedule"` was killed right around a poke boundary. Two things then 
happened concurrently:
   
   1. The task supervisor's final "reschedule" state update reached the 
execution API server, which inserts a `task_reschedule` row referencing the 
current `task_instance.id` (`TIRescheduleStatePayload` branch in 
`airflow/api_fastapi/execution_api/routes/task_instances.py`).
   2. The executor reported the task as failed, and the scheduler ran 
`handle_failure` for the same TI. Since the TI is retry-eligible and in RUNNING 
state, `fetch_handle_failure_context` calls `prepare_db_for_next_try` 
(`airflow/models/taskinstance.py`), which — in the scheduler's open transaction 
— records TI history, executes `DELETE FROM task_reschedule WHERE ti_id = <old 
id>`, and assigns a new uuid7 to `ti.id`. The new id is flushed later by 
`TaskInstance.save_to_db`.
   
   Between the DELETE and the UPDATE, `fetch_handle_failure_context` invokes 
listener hooks (`on_task_instance_failed`) **inside the same open 
transaction**. In our deployment the OpenLineage provider's listener forks a 
subprocess for its emission and blocks the caller up to `[openlineage] 
execution_timeout` (default 10s), then terminates it with up to 3 more seconds 
of grace (`_fork_execute` / `_terminate_with_wait` in the provider's 
`plugins/listener.py`). The subprocess hit the timeout, so the scheduler's 
transaction stayed open for ~11 seconds between the DELETE and the id-change 
UPDATE.
   
   The API server's insert committed inside that window. Under READ COMMITTED 
the scheduler's DELETE had only removed rows committed before it ran, so at 
flush time a `task_reschedule` row referencing the old id existed, and the 
id-change UPDATE failed:
   
   ```
   psycopg2.errors.ForeignKeyViolation: update or delete on table 
"task_instance" violates
   foreign key constraint "task_reschedule_ti_fkey" on table "task_reschedule"
   DETAIL:  Key (id)=(01a01dc5-da82-72df-a563-bd2d3ab1cbfb) is still referenced 
from table "task_reschedule".
   
   [SQL: UPDATE task_instance SET id=%(id)s::UUID, end_date=%(end_date)s, 
duration=%(duration)s,
    state=%(state)s, updated_at=%(updated_at)s WHERE task_instance.id = 
%(task_instance_id)s::UUID]
   [parameters: {'id': '01a01ded-b73d-7ba1-b080-d57875580320',
    'end_date': datetime.datetime(2026, 8, 20, 6, 48, 33, 339896, 
tzinfo=Timezone('UTC')),
    'duration': 1014.393528, 'state': <TaskInstanceState.UP_FOR_RETRY: 
'up_for_retry'>,
    'updated_at': datetime.datetime(2026, 8, 20, 6, 48, 44, 511734, 
tzinfo=Timezone('UTC')),
    'task_instance_id': '01a01dc5-da82-72df-a563-bd2d3ab1cbfb'}]
   ```
   
   The `IntegrityError` propagates out of `process_executor_events` -> 
`ti.handle_failure` -> `TaskInstance.save_to_db` -> `session.flush()`, is not 
handled anywhere in `_run_scheduler_loop`, and the scheduler process exits. 
(The same conflict does not repeat for the same TI after restart — the 
offending row is then committed and visible, so the DELETE removes it — but 
overall recovery took hours, see the impact note below.)
   
   Forensic timeline reconstructed from one occurrence (times UTC):
   
   | time | event | evidence |
   |---|---|---|
   | 06:31:38.946 | sensor try starts | `start_date` in the UP_FOR_RETRY log 
line |
   | ~06:48:33 | worker killed during rolling restart; poke ends; supervisor's 
reschedule request in flight | deployment timeline |
   | 06:48:33.340 | scheduler enters `fetch_handle_failure_context` (it stamps 
`end_date = utcnow()` at entry; that exact value ends up in the failing 
UPDATE's parameters). `prepare_db_for_next_try` runs immediately: the 
replacement uuid7 embeds timestamp 06:48:33.341, so the DELETE + id 
reassignment happened here, at the start of the window | `end_date` / 
`duration` params above; uuid7 timestamp of the new id |
   | 06:48:43.355 | `OpenLineage process with pid NN expired and will be 
terminated by listener` — i.e. the listener waited its full `execution_timeout` 
(default 10s: 33.34 + 10.0 ≈ 43.35) inside the transaction | scheduler log |
   | 06:48:44.511 | `Marking task as UP_FOR_RETRY` | scheduler log |
   | 06:48:44.533 | flush -> ForeignKeyViolation -> scheduler exits | scheduler 
log |
   
   Supporting detail: both UUIDs are uuid7, so they embed timestamps — the old 
id decodes to 06:05:00.930 (when the TI was created for this run) and the new 
id to 06:48:33.341, the same millisecond as the `end_date` stamp. The new id 
was assigned ~11.2s before the failing flush (06:48:44.533), which shows the 
transaction really did sit open across the listener wait with the DELETE 
already executed.
   
   OpenLineage is only the amplifier here: any listener that is slow in 
`on_task_instance_failed` stretches the race window — normally much shorter, 
just the gap between adjacent statements in the same function — to seconds, 
because listeners run between the `task_reschedule` DELETE and the id-change 
flush, inside the open transaction. The underlying race (reschedule-insert vs. 
retry-prep id change) exists without any listener, just with a much smaller 
window.
   
   ### What you think should happen instead?
   
   The scheduler should never crash on this. Ideas, not mutually exclusive:
   
   1. Re-issue (or move) the `TaskReschedule` delete so it is adjacent to the 
id-change flush — e.g. delete again right before `save_to_db` flushes a changed 
`ti.id`, leaving no multi-second gap for a concurrent insert.
   2. Make the two writers serialize: the execution API's reschedule branch 
inserts `task_reschedule` after reading the TI without a conflicting lock; 
taking a lock on the TI row (or rechecking TI state under lock) in one or both 
paths would force ordering.
   3. Defense in depth: handle the integrity conflict in the scheduler's 
failure-handling path without terminating the scheduler — the exact recovery 
(retrying vs. deferring the event) needs care around session state and event 
idempotency, so we leave the mechanism open.
   4. More generally, consider not invoking listener hooks while the 
failure-handling transaction is open — a slow listener currently extends DB 
transaction lifetime in the scheduler's critical path.
   
   ### How to reproduce
   
   This is an intermittent timing race; the listener stall makes it practical 
to hit, but the steps below make the collision likely, not certain:
   
   1. Airflow 3.1.x, PostgreSQL metadata DB, CeleryExecutor.
   2. A DAG with a sensor in `mode="reschedule"` (short `poke_interval`, e.g. 
60s) and `retries >= 1`.
   3. Register any listener whose `on_task_instance_failed` blocks ~10s (the 
OpenLineage provider pointing at a slow/unreachable HTTP endpoint with default 
`execution_timeout` reproduces this naturally).
   4. Kill the Celery worker (SIGKILL the pod/process) around a poke boundary, 
so that the executor reports failure while the supervisor's reschedule state 
update may still be in flight to the API server.
   5. When the reschedule update commits inside the listener stall, the 
scheduler dies with the ForeignKeyViolation above. Killing the worker does not 
guarantee the reschedule request is actually in flight at that moment, so 
several attempts may be needed. A deterministic reproduction would need to 
control the timing of the API server's reschedule commit relative to the 
listener stall (e.g. delaying the API server's transaction); we have not built 
that yet. Without step 3 the window is normally much shorter (no deliberate 
wait between the DELETE and the flush), which matches our observation that only 
the environment with the OpenLineage listener enabled ever hits it.
   
   ### Operating System
   
   Debian 12 (bookworm), Python 3.10.20 — container image derived from the 
official apache/airflow:3.1.8-python3.10 image (glibc 2.36 confirmed in the 
running container)
   
   ### Versions of Apache Airflow Providers
   
   apache-airflow-providers-openlineage==2.11.0
   apache-airflow-providers-celery==3.17.0
   apache-airflow-providers-postgres==6.6.0
   
   ### Deployment
   
   Official Apache Airflow Helm Chart
   
   ### Deployment details
   
   - Official Apache Airflow Helm chart `airflow-1.19.0`, custom image based on 
the official apache/airflow:3.1.8-python3.10 image
   - CeleryExecutor; external PostgreSQL 18.3 metadata DB
   - ~36 Celery workers; crash observed during rolling worker restarts 
triggered by an image change
   - OpenLineage provider enabled with HTTP transport (`timeout: 30`), listener 
`execution_timeout` left at default (10s)
   
   ### Anything else?
   
   Operational impact observed in our incident: after one rolling restart of 
the workers, the scheduler container restarted 11 times over ~3.2 hours before 
stabilizing (kubelet restart count and container timestamps). The first 
captured crash is the FK traceback above; the later restarts were not 
individually root-caused. The final two container instances' logs contain no FK 
error and end with a SIGTERM-initiated graceful exit (code 0) a few minutes 
into their life, while the scheduler was working through the backlog of 
task-failure events (one instance processed 113 failure events in its 5 
minutes). This is consistent with liveness-probe restarts — every failed TI's 
handling runs listener hooks (up to `execution_timeout` each) inside the 
scheduler loop, which can starve the scheduler heartbeat under a large failure 
backlog — but we did not retain the pod events to confirm this, and will attach 
pod events / probe logs if we capture a recurrence. An otherwise identical 
environment
  without the OpenLineage listener went through the same rolling restart with 
no scheduler crash.
   
   - Code inspection of `main` at commit 
00f1457ebfb650406e7cb29325ce3f1d252ce86a (2026-08-21) suggests the same 
structure is present (`prepare_db_for_next_try` deleting `task_reschedule` rows 
before reassigning `ti.id`, and listener hooks invoked inside 
`fetch_handle_failure_context` before the flush); we have not reproduced this 
on `main`.
   - Searched existing issues for `task_reschedule_ti_fkey` / this UPDATE 
pattern and found no report of this scenario.
   
   ### Are you willing to submit PR?
   
   - [x] Yes I am willing to submit a PR! (maintainer guidance on the preferred 
fix direction above would be appreciated)
   
   ### Code of Conduct
   
   - [x] I agree to follow this project's [Code of 
Conduct](https://github.com/apache/airflow/blob/main/CODE_OF_CONDUCT.md)
   


-- 
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]

Reply via email to