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

   ### Apache Airflow version
   
   3.2.2 (also present on main / google provider 22.6.0)
   
   ### What happened
   
   When a running triggerer job's heartbeat goes stale past `[triggerer] 
triggerer_health_check_threshold` (30s default) — e.g. because the 
TriggerRunner subprocess event loop is temporarily blocked on a busy/undersized 
triggerer, or because the new `runner_health_check_threshold` watchdog 
deliberately skips heartbeats when the runner is silent — a sibling triggerer's 
`assign_unassigned()` steals all of its trigger rows. The original triggerer, 
still alive, then diffs `load_triggers()`, sees the rows reassigned away, and 
calls `task.cancel()` on every one of its in-flight trigger coroutines.
   
   `DataprocSubmitTrigger.run()` (and `DataprocBatchTrigger`, 
`DataprocClusterTrigger` — same pattern) catches the `asyncio.CancelledError`, 
awaits `safe_to_cancel()` (supervisor comms call via greenback), finds the TI 
still `DEFERRED`, skips `cancel_job` — and then **does not re-raise**. The 
coroutine exits cleanly having emitted zero events.
   
   `cleanup_finished_triggers()` only recognizes a cancellation when 
`task.result()` raises `CancelledError`. A swallowed `CancelledError` makes the 
exit look like a crash:
   
   ```
   Trigger exited without sending an event. Dependent tasks will be failed.
   ```
   
   → `Trigger.submit_failure()` → every dependent deferred task is rescheduled 
with `next_method=__fail__` and fails on the worker with `TaskDeferralError: 
Trigger failure`, while the Dataproc job keeps running to completion.
   
   So a transient triggerer stall is amplified into a mass failure of all 
deferred tasks it hosted — the exact opposite of what trigger migration is 
supposed to achieve (transparent failover).
   
   This is the same bug shape fixed for `BigQueryInsertJobTrigger` in 
https://github.com/apache/airflow/pull/63730 ("Fix BigQueryInsertJobTrigger not 
propagating CancelledError"): the `except CancelledError` block must re-raise 
after cleanup so the framework can tell "cancelled during migration" from 
"crashed".
   
   ### How to reproduce
   
   1. Two triggerers, several `DataprocSubmitJobOperator(deferrable=True)` 
tasks deferred.
   2. Stall one triggerer job's heartbeat > `triggerer_health_check_threshold` 
(e.g. suspend the process, or saturate the runner event loop so 
`runner_health_check_threshold` trips heartbeat suppression).
   3. Sibling steals the trigger rows; original triggerer cancels its 
coroutines.
   4. All affected tasks fail with `TaskDeferralError: Trigger failure` instead 
of transparently migrating.
   
   Observed in production at scale on Cloud Composer 3 
(`composer-3-airflow-3.2.2-build.2`): bursts of 12+ and 3+ deferred tasks 
failing within a second, triggerer log showing 12x `Trigger exited without 
sending an event` + `Got response for unknown request frame` warnings + `Task 
cancelling ... coro=<greenback_shim()>` entries immediately beforehand.
   
   ### Suggested fix
   
   In 
`providers/google/src/airflow/providers/google/cloud/triggers/dataproc.py`, add 
a bare `raise` at the end of the `except asyncio.CancelledError:` blocks (all 
three trigger classes), so cancellation propagates and 
`cleanup_finished_triggers` takes the expected `except (CancelledError, ...) -> 
del + continue` path — no `submit_failure`, sibling triggerer's copy resumes 
the task normally.
   
   Possibly also worth auditing other google provider triggers (BigQuery 
excepted — fixed in #63730) and other providers for the same 
swallowed-CancelledError pattern.
   


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