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]
