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]

Reply via email to