anxkhn commented on code in PR #69118:
URL: https://github.com/apache/airflow/pull/69118#discussion_r4050534060
##########
providers/google/src/airflow/providers/google/cloud/sensors/bigquery.py:
##########
@@ -426,4 +435,18 @@ def poke(self, context: Context) -> bool:
table =
hook.get_client(project_id=self.project_id).get_table(table_ref)
except NotFound as err:
raise ValueError(f"Table {table_uri} not found") from err
- return table.streaming_buffer is None
+
+ if table.streaming_buffer is not None:
+ self._consecutive_empty = 0
+ return False
+
+ self._consecutive_empty += 1
+ self.log.info(
+ "Streaming buffer reported empty (%s/%s confirmations) for table:
%s",
+ self._consecutive_empty,
+ self.empty_confirmations,
+ table_uri,
+ )
+ if self._consecutive_empty >= self.empty_confirmations:
+ return True
+ return False
Review Comment:
Addressed in 1a8aeabe6d. In reschedule mode, once the first empty
observation is seen, `poke()` now completes the bounded confirmation sequence
in the same process before returning. This prevents operator reconstruction
from resetting the counter. A regression test exercises
`BaseSensorOperator.execute()` with `mode="reschedule"` and verifies two polls
occur without raising `AirflowRescheduleException`.\n\n---\nDrafted-by:
OpenCode (GPT-5.6) (no human review before posting)
--
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]