shahar1 commented on code in PR #69586:
URL: https://github.com/apache/airflow/pull/69586#discussion_r4090619258
##########
providers/google/src/airflow/providers/google/cloud/triggers/dataflow.py:
##########
@@ -336,9 +396,48 @@ def serialize(self) -> tuple[str, dict[str, Any]]:
"expected_terminal_state": self.expected_terminal_state,
"impersonation_chain": self.impersonation_chain,
"cancel_timeout": self.cancel_timeout,
+ "cancel_on_kill": self.cancel_on_kill,
+ "drain_pipeline": self.drain_pipeline,
},
)
+ async def on_kill(self) -> None:
+ """Stop the Dataflow job when the user acts on the deferred task."""
+ if not self.cancel_on_kill or not self.job_id or not self.project_id:
Review Comment:
Same as the `TemplateJobStartTrigger` comment: dropping the `project_id`
check lets `cancel_job` fall back to the connection's project.
```suggestion
if not self.cancel_on_kill or not self.job_id:
```
##########
providers/google/src/airflow/providers/google/cloud/triggers/dataflow.py:
##########
@@ -95,9 +108,48 @@ def serialize(self) -> tuple[str, dict[str, Any]]:
"poll_sleep": self.poll_sleep,
"impersonation_chain": self.impersonation_chain,
"cancel_timeout": self.cancel_timeout,
+ "cancel_on_kill": self.cancel_on_kill,
+ "drain_pipeline": self.drain_pipeline,
},
)
+ async def on_kill(self) -> None:
+ """Stop the Dataflow job when the user acts on the deferred task."""
+ if not self.cancel_on_kill or not self.job_id or not self.project_id:
Review Comment:
`DataflowHook.cancel_job` already falls back to the connection's default
project (`@GoogleBaseHook.fallback_to_default_project_id`), and every operator
passes `project_id=None` by default. With `not self.project_id` here, a
deferred kill silently does nothing for users who rely on the connection's
project. The BigQuery and Dataproc triggers check only `job_id`. See the review
body.
```suggestion
if not self.cancel_on_kill or not self.job_id:
```
##########
providers/google/src/airflow/providers/google/cloud/triggers/dataflow.py:
##########
@@ -95,9 +108,48 @@ def serialize(self) -> tuple[str, dict[str, Any]]:
"poll_sleep": self.poll_sleep,
"impersonation_chain": self.impersonation_chain,
"cancel_timeout": self.cancel_timeout,
+ "cancel_on_kill": self.cancel_on_kill,
+ "drain_pipeline": self.drain_pipeline,
},
)
+ async def on_kill(self) -> None:
+ """Stop the Dataflow job when the user acts on the deferred task."""
+ if not self.cancel_on_kill or not self.job_id or not self.project_id:
+ return
+ self.log.info(
+ "Stopping Dataflow job. Project ID: %s, Location: %s, Job ID: %s,
drain: %s",
+ self.project_id,
+ self.location,
+ self.job_id,
+ self.drain_pipeline,
+ )
+ try:
+ # Build the synchronous hook and cancel inside the worker thread:
the hook resolves the
+ # connection eagerly during construction, which must not run in
the triggerer's event loop.
+ await sync_to_async(self._stop_job)()
+ self.log.info("Dataflow job %s stopped.", self.job_id)
+ except Exception:
+ self.log.exception(
+ "Failed to stop Dataflow job %s. The job may still be
running.",
+ self.job_id,
+ )
+
+ def _stop_job(self) -> None:
+ """Cancel or drain the Dataflow job through the synchronous hook (runs
off the event loop)."""
+ hook = DataflowHook(
+ gcp_conn_id=self.gcp_conn_id,
+ impersonation_chain=self.impersonation_chain,
+ drain_pipeline=self.drain_pipeline,
+ cancel_timeout=self.cancel_timeout,
Review Comment:
With `cancel_timeout` set, `cancel()` polls `_wait_for_states` after sending
the request. Off the main thread its SIGALRM `timeout` does nothing, and the
poll runs on asgiref's shared single-thread executor, so it blocks other
triggers' `sync_to_async` calls and can make queued `on_kill` cancels get
dropped. See the review body. `None` sends the request and returns:
```suggestion
cancel_timeout=None,
```
##########
providers/google/src/airflow/providers/google/cloud/triggers/dataflow.py:
##########
@@ -336,9 +396,48 @@ def serialize(self) -> tuple[str, dict[str, Any]]:
"expected_terminal_state": self.expected_terminal_state,
"impersonation_chain": self.impersonation_chain,
"cancel_timeout": self.cancel_timeout,
+ "cancel_on_kill": self.cancel_on_kill,
+ "drain_pipeline": self.drain_pipeline,
},
)
+ async def on_kill(self) -> None:
+ """Stop the Dataflow job when the user acts on the deferred task."""
+ if not self.cancel_on_kill or not self.job_id or not self.project_id:
+ return
+ self.log.info(
+ "Stopping Dataflow job. Project ID: %s, Location: %s, Job ID: %s,
drain: %s",
+ self.project_id,
+ self.location,
+ self.job_id,
+ self.drain_pipeline,
+ )
+ try:
+ # Build the synchronous hook and cancel inside the worker thread:
the hook resolves the
+ # connection eagerly during construction, which must not run in
the triggerer's event loop.
+ await sync_to_async(self._stop_job)()
+ self.log.info("Dataflow job %s stopped.", self.job_id)
+ except Exception:
+ self.log.exception(
+ "Failed to stop Dataflow job %s. The job may still be
running.",
+ self.job_id,
+ )
+
+ def _stop_job(self) -> None:
+ """Cancel or drain the Dataflow job through the synchronous hook (runs
off the event loop)."""
+ hook = DataflowHook(
+ gcp_conn_id=self.gcp_conn_id,
+ impersonation_chain=self.impersonation_chain,
+ drain_pipeline=self.drain_pipeline,
+ cancel_timeout=self.cancel_timeout,
Review Comment:
Same as the `TemplateJobStartTrigger` comment: this only needs to send the
request, not wait for it.
```suggestion
cancel_timeout=None,
```
##########
providers/google/src/airflow/providers/google/cloud/operators/dataflow.py:
##########
@@ -298,6 +298,9 @@ class
DataflowTemplatedJobStartOperator(GoogleCloudBaseOperator):
https://cloud.google.com/dataflow/docs/templates/executing-templates
:param deferrable: Run operator in the deferrable mode.
+ :param cancel_on_kill: If True (default), cancel the Dataflow job when the
task is killed,
+ both while the operator is running and, for a deferred task, while it
waits in the
+ triggerer.
Review Comment:
`BaseTrigger.on_kill` exists only in Airflow 3.3+, and this provider
supports `apache-airflow>=2.11.0`. On older versions a deferred kill still
leaves the job running. Please add "(Airflow 3.3+)" here and at the other two
operators and the YAML `on_kill` docstring, or add the pre-3.3 fallback the
BigQuery and Dataproc triggers use.
--
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]