VladaZakharova commented on code in PR #69118:
URL: https://github.com/apache/airflow/pull/69118#discussion_r4131756308
##########
providers/google/src/airflow/providers/google/cloud/sensors/bigquery.py:
##########
@@ -422,8 +432,28 @@ def poke(self, context: Context) -> bool:
)
hook = BigQueryHook(gcp_conn_id=self.gcp_conn_id,
impersonation_chain=self.impersonation_chain)
- try:
- 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
+ client = hook.get_client(project_id=self.project_id)
Review Comment:
Keeping the confirmation loop inside poke() means mode="reschedule" occupies
a worker and can bypass the normal sensor timeout check when the final poll
succeeds. Could the empty-state/timestamp be persisted across actual
reschedules, or could this mode use the trigger implementation instead? Please
also add an execute-level timeout test; the current direct poke() test does not
cover BaseSensorOperator semantics.
##########
providers/google/tests/system/google/cloud/bigquery/example_bigquery_streaming_buffer_sensor.py:
##########
@@ -109,19 +108,6 @@ def streaming_insert(ds: str | None = None) -> None:
rows=[{"value": 100, "ds": ds}],
fail_on_error=True,
)
Review Comment:
The updated system-test DAG inserts rows, runs the synchronous sensor,
performs an update, and only then runs the deferrable sensor.
Therefore, only the synchronous sensor sees the “immediately after streaming
insert” condition. Add another streaming insert or a separate table immediately
before the deferrable senso
--
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]