Skip to content

Cancel the Dataflow job when a user kills the deferred task - #69586

Open
steveahnahn wants to merge 2 commits into
apache:mainfrom
steveahnahn:dataflow-cancel-job-on-user-kill-deferred
Open

Cancel the Dataflow job when a user kills the deferred task#69586
steveahnahn wants to merge 2 commits into
apache:mainfrom
steveahnahn:dataflow-cancel-job-on-user-kill-deferred

Conversation

@steveahnahn

@steveahnahn steveahnahn commented Jul 8, 2026

Copy link
Copy Markdown
Contributor

Problem

  • Killing a deferred Dataflow task (mark failed, clear, mark success) leaves the Dataflow job running and billing in GCP.
  • Neither TemplateJobStartTrigger (templated/flex jobs) nor DataflowStartYamlJobTrigger (YAML jobs) has on_kill; the BigQuery and Dataproc triggers in this provider already cancel in the same situation.
  • The YAML operator's own on_kill even documents the gap: it "will not be called if a task instance is killed in a deferred state."
  • A streaming job runs until someone notices it in the Cloud console.

Change

  • Add on_kill to both TemplateJobStartTrigger and DataflowStartYamlJobTrigger: checks cancel_on_kill, job_id, project_id, then cancels via the sync DataflowHook. Identical cancel path in both.
  • drain_pipeline is honored, matching the non deferred kill path.
  • Hook is constructed and invoked inside a worker thread (sync_to_async): DataflowHook resolves its connection eagerly at construction, which must not run in the triggerer event loop. Construction inside the try means a connection error is logged, not raised.
  • Fires only on user kill, never on triggerer restart, redistribution, or timeout.

End to end verification

  1. Deferrable operator launched the job; task deferred, job JOB_STATE_RUNNING.
deferred task, job running
  1. Marked the task failed from the UI.
marking the deferred task failed
  1. Triggerer log, no exception escaping on_kill:
Trigger cancelled by user action, invoking on_kill
Stopping Dataflow job. Project ID: <redacted>, Location: us-central1, Job ID: 2026-07-08_21_43_34-418737444659630262, drain: False
  1. Job went JOB_STATE_RUNNING -> JOB_STATE_CANCELLING -> JOB_STATE_CANCELLED (gcloud dataflow jobs list). With cancel_on_kill=False the job keeps running.

The YAML trigger uses the identical DataflowHook.cancel_job path proven above, so the e2e run covers its cancel mechanism; it is additionally unit tested.

Tests

  • on_kill tests for both triggers: cancel, drain_pipeline honored, no ops (cancel_on_kill=False, missing job_id/project_id), hook errors swallowed, hook construction errors swallowed; plus operator threading of drain_pipeline into the trigger.
  • All new assertions fail without the change; 95 trigger tests pass with it; provider mypy clean.

Was generative AI tooling used to co-author this PR?
  • Yes, Claude Code (Opus 4.8)

Generated-by: Claude Code (Opus 4.8) following the guidelines

@boring-cyborg boring-cyborg Bot added area:providers provider:google Google (including GCP) related issues labels Jul 8, 2026
@steveahnahn
steveahnahn force-pushed the dataflow-cancel-job-on-user-kill-deferred branch 2 times, most recently from cadba9a to f4df591 Compare July 8, 2026 06:46
@steveahnahn
steveahnahn marked this pull request as ready for review July 8, 2026 21:02
@steveahnahn
steveahnahn requested a review from shahar1 as a code owner July 8, 2026 21:02
@steveahnahn
steveahnahn force-pushed the dataflow-cancel-job-on-user-kill-deferred branch from f4df591 to 6e7fa45 Compare July 10, 2026 23:05
@potiuk potiuk added the ready for maintainer review Set after triaging when all criteria pass. label Jul 11, 2026

@aaron-y-chen aaron-y-chen left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hi, thanks for the PR, very impressive :)

Comment thread providers/google/src/airflow/providers/google/cloud/operators/dataflow.py Outdated
@steveahnahn
steveahnahn force-pushed the dataflow-cancel-job-on-user-kill-deferred branch from 6e7fa45 to 54df2f5 Compare July 16, 2026 02:04
Killing a deferred Dataflow task (mark-failed, clear or mark-succeeded)
leaves the Dataflow job running: the trigger has no on_kill, so the
user-kill cancellation the sibling BigQuery and Dataproc triggers
already perform never happens, and a streaming pipeline keeps running
and billing until someone notices in the Cloud console. Before the
Airflow 3.3 trigger cancellation redesign the legacy path cancelled the
job from inside the trigger, so this is also a behavior regression for
deferrable users.

Implement on_kill on TemplateJobStartTrigger following the merged
BigQuery/Dataproc/EMR pattern, honoring the operator's drain_pipeline
setting so streaming jobs are drained rather than cancelled when
configured, exactly as the non-deferred kill path does. The hook is
built and invoked inside a worker thread because the synchronous
DataflowHook resolves its connection during construction, which must
not run in the triggerer's event loop.
@steveahnahn
steveahnahn force-pushed the dataflow-cancel-job-on-user-kill-deferred branch from 54df2f5 to bd3f1bb Compare July 16, 2026 02:05
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:providers provider:google Google (including GCP) related issues ready for maintainer review Set after triaging when all criteria pass.

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants