Skip to content

Commit b7203cd

Browse files
author
Peter Wicks
authored
Fix MaxID logic for GCSToBigQueryOperator (#26768)
1 parent 9a6fc73 commit b7203cd

2 files changed

Lines changed: 15 additions & 13 deletions

File tree

airflow/providers/google/cloud/transfers/gcs_to_bigquery.py

Lines changed: 11 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -319,15 +319,15 @@ def execute(self, context: Context):
319319
location=self.location,
320320
use_legacy_sql=False,
321321
)
322-
row = list(bq_hook.get_job(job_id=job_id, location=self.location).result())
323-
if row:
324-
max_id = row[0] if row[0] else 0
325-
self.log.info(
326-
'Loaded BQ data with max %s.%s=%s',
327-
self.destination_project_dataset_table,
328-
self.max_id_key,
329-
max_id,
330-
)
331-
return max_id
332-
else:
322+
result = bq_hook.get_job(job_id=job_id, location=self.location).result()
323+
row = next(iter(result), None)
324+
if row is None:
333325
raise RuntimeError(f"The {select_command} returned no rows!")
326+
max_id = row[0]
327+
self.log.info(
328+
'Loaded BQ data with max %s.%s=%s',
329+
self.destination_project_dataset_table,
330+
self.max_id_key,
331+
max_id,
332+
)
333+
return max_id

tests/providers/google/cloud/transfers/test_gcs_to_bigquery.py

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,8 @@
2020
import unittest
2121
from unittest import mock
2222

23+
from google.cloud.bigquery.table import Row
24+
2325
from airflow.providers.google.cloud.transfers.gcs_to_bigquery import GCSToBigQueryOperator
2426

2527
TASK_ID = 'test-gcs-to-bq-operator'
@@ -43,11 +45,11 @@ def test_execute_explicit_project(self, bq_hook):
4345
max_id_key=MAX_ID_KEY,
4446
)
4547

46-
bq_hook.return_value.get_job.return_value.result.return_value = ('1',)
48+
bq_hook.return_value.get_job.return_value.result.return_value = [Row(('100',), {'f0_': 0})]
4749

4850
result = operator.execute(None)
4951

50-
assert result == '1'
52+
assert result == '100'
5153

5254
bq_hook.return_value.run_query.assert_called_once_with(
5355
sql="SELECT MAX(id) FROM `test-project.dataset.table`",

0 commit comments

Comments
 (0)