This is an automated email from the ASF dual-hosted git repository.
jrmccluskey pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new 7872437f3fd Fix partial insertError metric status bug (#40356)
7872437f3fd is described below
commit 7872437f3fd8a96922c7a9031f20b153824abe9b
Author: Jack McCluskey <[email protected]>
AuthorDate: Wed Sep 30 13:50:42 2026 -0400
Fix partial insertError metric status bug (#40356)
* Fix partial insertError metric status bug
* Update bigquery_tools.py
streamline comment
---
sdks/python/apache_beam/io/gcp/bigquery_tools.py | 5 +++-
.../apache_beam/io/gcp/bigquery_tools_test.py | 30 ++++++++++++++++++++++
2 files changed, 34 insertions(+), 1 deletion(-)
diff --git a/sdks/python/apache_beam/io/gcp/bigquery_tools.py
b/sdks/python/apache_beam/io/gcp/bigquery_tools.py
index 5cc3441e171..abe204c237a 100644
--- a/sdks/python/apache_beam/io/gcp/bigquery_tools.py
+++ b/sdks/python/apache_beam/io/gcp/bigquery_tools.py
@@ -780,7 +780,10 @@ class BigQueryWrapper(object):
service_call_metric.call('ok')
else:
for insert_error in errors:
- service_call_metric.call(insert_error['errors'][0])
+ # Record the BigQuery reason (e.g. 'invalid') of the first error for
+ # each failed row.
+ row_errors = insert_error.get('errors') or [{}]
+ service_call_metric.call(row_errors[0].get('reason') or 'unknown')
except (ClientError, GoogleAPICallError) as e:
# e.code contains the numeric http status code.
service_call_metric.call(e.code)
diff --git a/sdks/python/apache_beam/io/gcp/bigquery_tools_test.py
b/sdks/python/apache_beam/io/gcp/bigquery_tools_test.py
index 078c4216094..c2a5c73f675 100644
--- a/sdks/python/apache_beam/io/gcp/bigquery_tools_test.py
+++ b/sdks/python/apache_beam/io/gcp/bigquery_tools_test.py
@@ -575,6 +575,36 @@ class TestBigQueryWrapper(unittest.TestCase):
self.verify_write_call_metric(
"my_project", "my_dataset", "my_table", "ok", 1)
+ @unittest.skipIf(ClientError is None, 'GCP dependencies are not installed')
+ def test_insert_rows_sets_metric_on_row_errors(self):
+ MetricsEnvironment.process_wide_container().reset()
+ client = mock.Mock()
+
+ def row_error(index, *reasons):
+ return {
+ 'index': index,
+ 'errors': [{
+ 'reason': r, 'message': 'msg'
+ } for r in reasons],
+ }
+
+ client.insert_rows_json.return_value = [
+ row_error(0, 'invalid', 'stopped'),
+ row_error(1, 'invalid'),
+ row_error(2, 'stopped'),
+ ]
+ wrapper = beam.io.gcp.bigquery_tools.BigQueryWrapper(client)
+ success, errors = wrapper.insert_rows(
+ "my_project", "my_dataset", "my_table", [{'a': 1}] * 3)
+
+ self.assertFalse(success)
+ self.assertEqual(client.insert_rows_json.return_value, errors)
+ # One metric per failed row, labelled with the reason of its first error.
+ self.verify_write_call_metric(
+ "my_project", "my_dataset", "my_table", "invalid", 2)
+ self.verify_write_call_metric(
+ "my_project", "my_dataset", "my_table", "stopped", 1)
+
@unittest.skipIf(ClientError is None, 'GCP dependencies are not installed')
def test_start_query_job_priority_configuration(self):
client = mock.Mock()