kaxil opened a new pull request, #73454:
URL: https://github.com/apache/airflow/pull/73454
When a triggerer's heartbeat lapses for longer than
`triggerer_health_check_threshold`, `Trigger.assign_unassigned` hands its
triggers to another triggerer. The old triggerer is often still alive: the
runner-silence watchdog deliberately skips heartbeats when the runner
subprocess stops responding, and a long GC pause, a stopped container or a
network partition do the same. When it recovers,
`TriggerRunnerSupervisor.update_triggers` sees those triggers missing from its
assignment list and queues them for cancellation with the user-action message,
so `run_trigger` calls `on_kill()` on every one of them. Nothing in that path
distinguishes "the row was deleted because a user cleared the task" from "the
row was reassigned".
For any trigger that cancels remote work from `on_kill()` this cancels a job
the new owner has just started polling. The batch triggers in the Anthropic and
Common AI providers, and the EMR, Redshift Data, Databricks, Dataproc and
BigQuery triggers all do. The `BaseTrigger.on_kill` docstring and the deferring
docs both promise it does not fire on redistribution, so this is a contract
violation, present since `on_kill()` shipped in 3.3.0 (#65590). The
`cancel_triggers` docstring claimed redistribution "goes through a separate
path"; there was no such path.
## Solution
Before queueing a cancel, the supervisor asks the DB which of the departed
ids still have a trigger row and who owns them now
(`Trigger.fetch_assignments`, behind an overridable `fetch_trigger_assignments`
hook that owns its own session so `update_triggers` stays DB-free for
subclasses without a metadata DB). A row that survives after leaving our list
was reassigned, so those ids go to a new `releasing_triggers` set and reach the
runner as `TriggerStateSync.to_release`; the runner cancels those asyncio tasks
with a distinct sentinel, and `run_trigger` drops the trigger without
`on_kill()`. Only a row that is gone, which is what a clear or
mark-success/failed produces, takes the existing user-action path. The
classification runs once per departed id, since ids already queued are excluded
from the next loop's difference. The new field defaults to an empty set, so
existing `TriggerStateSync` producers are unaffected.
**Why in the supervisor rather than the runner.** The runner subprocess has
no DB access; the supervisor already owns the assignment query and the per-loop
`create_session()` in `build_trigger_workloads`, so the one extra `SELECT id
... WHERE id IN (...)` sits next to it, and the runner keeps its existing
message-based contract.
## Reproduction
Two triggerers started with `--queues afl297` and
`triggerer.queues_enabled=True` (so the standalone triggerer stays out of it),
a task on that queue deferring to a trigger whose `on_kill()` appends a line to
a file, then `SIGSTOP` on the owning triggerer and its runner child until the
trigger row moved, then `SIGCONT`.
| | Old owner after `SIGCONT` | `on_kill()` marker |
|---|---|---|
| main | `Trigger cancelled by user action, invoking on_kill` | written by
the old owner's runner while the row already belonged to the new owner and the
task was still `deferred` |
| this PR | `Triggers were reassigned to another triggerer, releasing them
without on_kill new_owners={3: 8}` and, in the runner, `Trigger reassigned to
another triggerer, dropping it without on_kill` | not written |
With this PR, marking the deferred task failed through the REST API still
produced the marker from the new owner's runner within seconds, so the
user-action path is intact.
## Known gap
Row existence proves the row was reassigned, not that the new owner has
started it. If a user clears the task in the window after the old owner
classified the trigger as reassigned and before the new owner fetched its
workload, neither side calls `on_kill()`: the old owner released, and the new
owner skips a trigger whose task instance is gone. This is the same outcome
main has when the old owner has actually died, so `on_kill()` on a user action
was already best-effort while a failover is in flight. The fix errs toward
leaving remote work running rather than cancelling work another triggerer is
polling.
---
* Read the **[Pull Request
Guidelines](https://github.com/apache/airflow/blob/main/contributing-docs/05_pull_requests.rst#pull-request-guidelines)**
for more information. Note: commit author/co-author name and email in commits
become permanently public when merged.
* For fundamental code changes, an Airflow Improvement Proposal
([AIP](https://cwiki.apache.org/confluence/display/AIRFLOW/Airflow+Improvement+Proposals))
is needed.
* When adding dependency, check compliance with the [ASF 3rd Party License
Policy](https://www.apache.org/legal/resolved.html#category-x).
* For significant user-facing changes create newsfragment:
`{pr_number}.significant.rst`, in
[airflow-core/newsfragments](https://github.com/apache/airflow/tree/main/airflow-core/newsfragments).
You can add this file in a follow-up commit after the PR is created so you
know the PR number.
--
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]