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):