This is an automated email from the ASF dual-hosted git repository.

shahar1 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new eda8563d983 Keep polling in the BigQuery check triggers while a job is 
running (#74305)
eda8563d983 is described below

commit eda8563d98347c32f0b80db6d0b0392bb3e4cf1d
Author: Bingqin Wang <[email protected]>
AuthorDate: Tue Oct 6 00:13:36 2026 -0500

    Keep polling in the BigQuery check triggers while a job is running (#74305)
---
 .../providers/google/cloud/triggers/bigquery.py    |  34 +++---
 .../unit/google/cloud/triggers/test_bigquery.py    | 121 +++++++++++++++++++++
 2 files changed, 142 insertions(+), 13 deletions(-)

diff --git 
a/providers/google/src/airflow/providers/google/cloud/triggers/bigquery.py 
b/providers/google/src/airflow/providers/google/cloud/triggers/bigquery.py
index 244043ca5d9..359caa3a07a 100644
--- a/providers/google/src/airflow/providers/google/cloud/triggers/bigquery.py
+++ b/providers/google/src/airflow/providers/google/cloud/triggers/bigquery.py
@@ -581,18 +581,25 @@ class 
BigQueryIntervalCheckTrigger(BigQueryInsertJobTrigger):
                         }
                     )
                     return
-                elif (
-                    first_job_response_from_hook["status"] == "pending"
-                    or second_job_response_from_hook["status"] == "pending"
+                elif "error" in (
+                    first_job_response_from_hook["status"],
+                    second_job_response_from_hook["status"],
                 ):
-                    self.log.info("Query is still running...")
-                    self.log.info("Sleeping for %s seconds.", 
self.poll_interval)
-                    await asyncio.sleep(self.poll_interval)
-                else:
+                    # Report the job that failed, which is not necessarily the 
second one.
+                    failed_job_response = (
+                        first_job_response_from_hook
+                        if first_job_response_from_hook["status"] == "error"
+                        else second_job_response_from_hook
+                    )
                     yield TriggerEvent(
-                        {"status": "error", "message": 
second_job_response_from_hook["message"], "data": None}
+                        {"status": "error", "message": 
failed_job_response["message"], "data": None}
                     )
                     return
+                else:
+                    # Neither job failed, and at least one is still "pending" 
or "running".
+                    self.log.info("Query is still running...")
+                    self.log.info("Sleeping for %s seconds.", 
self.poll_interval)
+                    await asyncio.sleep(self.poll_interval)
 
         except Exception as e:
             self.log.exception("Exception occurred while checking for query 
completion")
@@ -684,15 +691,16 @@ class BigQueryValueCheckTrigger(BigQueryInsertJobTrigger):
                     hook.value_check(self.sql, self.pass_value, _records, 
self.tolerance)
                     yield TriggerEvent({"status": "success", "message": "Job 
completed", "records": _records})
                     return
-                elif response_from_hook["status"] == "pending":
-                    self.log.info("Query is still running...")
-                    self.log.info("Sleeping for %s seconds.", 
self.poll_interval)
-                    await asyncio.sleep(self.poll_interval)
-                else:
+                elif response_from_hook["status"] == "error":
                     yield TriggerEvent(
                         {"status": "error", "message": 
response_from_hook["message"], "records": None}
                     )
                     return
+                else:
+                    # The job is still "pending" or "running".
+                    self.log.info("Query is still running...")
+                    self.log.info("Sleeping for %s seconds.", 
self.poll_interval)
+                    await asyncio.sleep(self.poll_interval)
         except Exception as e:
             self.log.exception("Exception occurred while checking for query 
completion")
             yield TriggerEvent({"status": "error", "message": str(e)})
diff --git a/providers/google/tests/unit/google/cloud/triggers/test_bigquery.py 
b/providers/google/tests/unit/google/cloud/triggers/test_bigquery.py
index 25a29b8118d..2c23d3050e0 100644
--- a/providers/google/tests/unit/google/cloud/triggers/test_bigquery.py
+++ b/providers/google/tests/unit/google/cloud/triggers/test_bigquery.py
@@ -756,6 +756,99 @@ class TestBigQueryIntervalCheckTrigger:
             == actual
         )
 
+    @pytest.mark.asyncio
+    @pytest.mark.parametrize(
+        ("first_status", "second_status"),
+        [
+            pytest.param("running", "success", id="first-running"),
+            pytest.param("success", "running", id="second-running"),
+            pytest.param("running", "running", id="both-running"),
+        ],
+    )
+    @mock.patch("asyncio.sleep", new_callable=AsyncMock)
+    
@mock.patch("airflow.providers.google.cloud.hooks.bigquery.BigQueryAsyncHook.interval_check")
+    
@mock.patch("airflow.providers.google.cloud.hooks.bigquery.BigQueryAsyncHook.get_records")
+    
@mock.patch("airflow.providers.google.cloud.hooks.bigquery.BigQueryAsyncHook.get_job_output")
+    
@mock.patch("airflow.providers.google.cloud.hooks.bigquery.BigQueryAsyncHook.get_sync_hook")
+    
@mock.patch("airflow.providers.google.cloud.hooks.bigquery.BigQueryAsyncHook.get_job_status")
+    async def test_interval_check_trigger_keeps_polling_while_a_job_is_running(
+        self,
+        mock_job_status,
+        mock_sync_hook,
+        mock_get_job_output,
+        mock_get_records,
+        mock_interval_check,
+        mock_sleep,
+        first_status,
+        second_status,
+        interval_check_trigger,
+    ):
+        """A job that is still running is polled again instead of failing the 
check (#73981)."""
+        mock_job_status.side_effect = [
+            {"status": first_status, "message": f"Job {first_status}"},
+            {"status": second_status, "message": f"Job {second_status}"},
+            {"status": "success", "message": "Job completed"},
+            {"status": "success", "message": "Job completed"},
+        ]
+        mock_sync_hook.return_value.is_default_universe.return_value = True
+        mock_get_records.side_effect = lambda *_, **__: [[100]]
+
+        actual = await interval_check_trigger.run().asend(None)
+
+        assert actual == TriggerEvent(
+            {
+                "status": "success",
+                "message": "Job completed",
+                "first_row_data": [100],
+                "second_row_data": [100],
+            }
+        )
+        
mock_sleep.assert_awaited_once_with(INTERVAL_CHECK_POLLING_PERIOD_SECONDS)
+        mock_interval_check.assert_called_once()
+
+    @pytest.mark.asyncio
+    @pytest.mark.parametrize(
+        ("first_response", "second_response", "expected_message"),
+        [
+            pytest.param(
+                {"status": "error", "message": "first failed"},
+                {"status": "success", "message": "Job completed"},
+                "first failed",
+                id="first-failed-second-succeeded",
+            ),
+            pytest.param(
+                {"status": "success", "message": "Job completed"},
+                {"status": "error", "message": "second failed"},
+                "second failed",
+                id="second-failed",
+            ),
+            pytest.param(
+                {"status": "error", "message": "first failed"},
+                {"status": "running", "message": "Job running"},
+                "first failed",
+                id="first-failed-second-still-running",
+            ),
+        ],
+    )
+    @mock.patch("asyncio.sleep", new_callable=AsyncMock)
+    
@mock.patch("airflow.providers.google.cloud.hooks.bigquery.BigQueryAsyncHook.get_job_status")
+    async def test_interval_check_trigger_reports_the_failed_job(
+        self,
+        mock_job_status,
+        mock_sleep,
+        first_response,
+        second_response,
+        expected_message,
+        interval_check_trigger,
+    ):
+        """The error event carries the message of the job that failed, without 
waiting for the other one."""
+        mock_job_status.side_effect = [first_response, second_response]
+
+        actual = await interval_check_trigger.run().asend(None)
+
+        assert actual == TriggerEvent({"status": "error", "message": 
expected_message, "data": None})
+        mock_sleep.assert_not_awaited()
+
     @pytest.mark.asyncio
     
@mock.patch("airflow.providers.google.cloud.hooks.bigquery.BigQueryAsyncHook.get_job_status")
     async def test_interval_check_trigger_exception(self, mock_job_status, 
caplog, interval_check_trigger):
@@ -833,6 +926,34 @@ class TestBigQueryValueCheckTrigger:
         # Prevents error when task is destroyed while in "pending" state
         asyncio.get_event_loop().stop()
 
+    @pytest.mark.asyncio
+    @mock.patch("asyncio.sleep", new_callable=AsyncMock)
+    
@mock.patch("airflow.providers.google.cloud.hooks.bigquery.BigQueryAsyncHook.value_check")
+    
@mock.patch("airflow.providers.google.cloud.hooks.bigquery.BigQueryAsyncHook.get_records")
+    
@mock.patch("airflow.providers.google.cloud.hooks.bigquery.BigQueryAsyncHook.get_job_output")
+    
@mock.patch("airflow.providers.google.cloud.hooks.bigquery.BigQueryAsyncHook.get_job_status")
+    async def test_value_check_op_trigger_keeps_polling_while_job_is_running(
+        self,
+        mock_job_status,
+        mock_get_job_output,
+        mock_get_records,
+        mock_value_check,
+        mock_sleep,
+        value_check_trigger,
+    ):
+        """A job that is still running is polled again instead of failing the 
check (#73981)."""
+        mock_job_status.side_effect = [
+            {"status": "running", "message": "Job running"},
+            {"status": "success", "message": "Job completed"},
+        ]
+        mock_get_records.side_effect = lambda *_, **__: [[4]]
+
+        actual = await value_check_trigger.run().asend(None)
+
+        assert actual == TriggerEvent({"status": "success", "message": "Job 
completed", "records": [4]})
+        mock_sleep.assert_awaited_once_with(POLLING_PERIOD_SECONDS)
+        mock_value_check.assert_called_once()
+
     @pytest.mark.asyncio
     
@mock.patch("airflow.providers.google.cloud.hooks.bigquery.BigQueryAsyncHook.get_job_status")
     async def test_value_check_op_trigger_fail(self, mock_job_status, 
value_check_trigger):

Reply via email to