Skip to content
This repository was archived by the owner on Apr 1, 2026. It is now read-only.

Commit 2e0e402

Browse files
committed
updated sync files
1 parent 253284f commit 2e0e402

8 files changed

Lines changed: 135 additions & 383 deletions

File tree

google/cloud/bigtable/data/_sync_autogen/_mutate_rows.py

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -84,7 +84,10 @@ def __init__(
8484
self.is_retryable,
8585
metric.backoff_generator,
8686
operation_timeout,
87-
exception_factory=_retry_exception_factory,
87+
exception_factory=self._operation_metric.track_terminal_error(
88+
_retry_exception_factory
89+
),
90+
on_error=self._operation_metric.track_retryable_error,
8891
)
8992
self.timeout_generator = _attempt_timeout_generator(
9093
attempt_timeout, operation_timeout

google/cloud/bigtable/data/_sync_autogen/_read_rows.py

Lines changed: 9 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919
from __future__ import annotations
2020
from typing import Sequence, TYPE_CHECKING
2121
import time
22+
from grpc import StatusCode
2223
from google.cloud.bigtable_v2.types import ReadRowsRequest as ReadRowsRequestPB
2324
from google.cloud.bigtable_v2.types import ReadRowsResponse as ReadRowsResponsePB
2425
from google.cloud.bigtable_v2.types import RowSet as RowSetPB
@@ -107,7 +108,10 @@ def start_operation(self) -> CrossSync._Sync_Impl.Iterable[Row]:
107108
self._predicate,
108109
self._operation_metric.backoff_generator,
109110
self.operation_timeout,
110-
exception_factory=_retry_exception_factory,
111+
exception_factory=self._operation_metric.track_terminal_error(
112+
_retry_exception_factory
113+
),
114+
on_error=self._operation_metric.track_retryable_error,
111115
)
112116

113117
def _read_rows_attempt(self) -> CrossSync._Sync_Impl.Iterable[Row]:
@@ -268,7 +272,7 @@ def merge_rows(
268272
if self._operation_metric.active_attempt is not None:
269273
self._operation_metric.active_attempt.application_blocking_time_ns += (
270274
time.monotonic_ns() - block_time
271-
) * 1000
275+
)
272276
break
273277
c = it.__next__()
274278
except _ResetRow as e:
@@ -285,9 +289,10 @@ def merge_rows(
285289
continue
286290
except CrossSync._Sync_Impl.StopIteration:
287291
raise InvalidChunk("premature end of stream")
292+
except GeneratorExit as close_exception:
293+
self._operation_metric.end_with_status(StatusCode.CANCELLED)
294+
raise close_exception
288295
except Exception as generic_exception:
289-
if not self._predicate(generic_exception):
290-
self._operation_metric.end_attempt_with_status(generic_exception)
291296
raise generic_exception
292297
else:
293298
self._operation_metric.end_with_success()

google/cloud/bigtable/data/_sync_autogen/client.py

Lines changed: 13 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -759,7 +759,6 @@ def __init__(
759759
default_retryable_errors or ()
760760
)
761761
self._metrics = BigtableClientSideMetricsController(
762-
client._metrics_interceptor,
763762
handlers=[],
764763
project_id=self.client.project,
765764
instance_id=instance_id,
@@ -941,8 +940,9 @@ def read_row(
941940
)
942941
results_generator = row_merger.start_operation()
943942
try:
944-
return results_generator.__next__()
945-
except StopIteration:
943+
results = [a for a in results_generator]
944+
return results[0]
945+
except IndexError:
946946
return None
947947

948948
def read_rows_sharded(
@@ -1142,7 +1142,10 @@ def execute_rpc():
11421142
predicate,
11431143
operation_metric.backoff_generator,
11441144
operation_timeout,
1145-
exception_factory=_retry_exception_factory,
1145+
exception_factory=operation_metric.track_terminal_error(
1146+
_retry_exception_factory
1147+
),
1148+
on_error=operation_metric.track_retryable_error,
11461149
)
11471150

11481151
def mutations_batcher(
@@ -1265,7 +1268,10 @@ def mutate_row(
12651268
predicate,
12661269
operation_metric.backoff_generator,
12671270
operation_timeout,
1268-
exception_factory=_retry_exception_factory,
1271+
exception_factory=operation_metric.track_terminal_error(
1272+
_retry_exception_factory
1273+
),
1274+
on_error=operation_metric.track_retryable_error,
12691275
)
12701276

12711277
def bulk_mutate_rows(
@@ -1371,7 +1377,7 @@ def check_and_mutate_row(
13711377
):
13721378
false_case_mutations = [false_case_mutations]
13731379
false_case_list = [m._to_pb() for m in false_case_mutations or []]
1374-
with self._metrics.create_operation(OperationType.CHECK_AND_MUTATE):
1380+
with self._metrics.create_operation(OperationType.CHECK_AND_MUTATE) as op:
13751381
result = self.client._gapic_client.check_and_mutate_row(
13761382
request=CheckAndMutateRowRequest(
13771383
true_mutations=true_case_list,
@@ -1425,7 +1431,7 @@ def read_modify_write_row(
14251431
rules = [rules]
14261432
if not rules:
14271433
raise ValueError("rules must contain at least one item")
1428-
with self._metrics.create_operation(OperationType.READ_MODIFY_WRITE):
1434+
with self._metrics.create_operation(OperationType.READ_MODIFY_WRITE) as op:
14291435
result = self.client._gapic_client.read_modify_write_row(
14301436
request=ReadModifyWriteRowRequest(
14311437
rules=[rule._to_pb() for rule in rules],

google/cloud/bigtable/data/_sync_autogen/metrics_interceptor.py

Lines changed: 20 additions & 56 deletions
Original file line numberDiff line numberDiff line change
@@ -15,13 +15,12 @@
1515
# This file is automatically generated by CrossSync. Do not edit manually.
1616

1717
from __future__ import annotations
18+
from typing import Sequence
1819
import time
1920
from functools import wraps
20-
from google.cloud.bigtable.data._metrics.data_model import (
21-
OPERATION_INTERCEPTOR_METADATA_KEY,
22-
)
2321
from google.cloud.bigtable.data._metrics.data_model import ActiveOperationMetric
2422
from google.cloud.bigtable.data._metrics.data_model import OperationState
23+
from google.cloud.bigtable.data._metrics.data_model import OperationType
2524
from google.cloud.bigtable.data._metrics.handlers._base import MetricsHandler
2625
from grpc import UnaryUnaryClientInterceptor
2726
from grpc import UnaryStreamClientInterceptor
@@ -33,19 +32,7 @@ def _with_operation_from_metadata(func):
3332

3433
@wraps(func)
3534
def wrapper(self, continuation, client_call_details, request):
36-
found_operation_id: str | None = None
37-
try:
38-
new_metadata: list[tuple[str, str]] = []
39-
if client_call_details.metadata:
40-
for k, v in client_call_details.metadata:
41-
if k == OPERATION_INTERCEPTOR_METADATA_KEY:
42-
found_operation_id = v
43-
else:
44-
new_metadata.append((k, v))
45-
client_call_details.metadata = new_metadata
46-
except Exception:
47-
pass
48-
operation: "ActiveOperationMetric" = self.operation_map.get(found_operation_id)
35+
operation: "ActiveOperationMetric" | None = ActiveOperationMetric.get_active()
4936
if operation:
5037
if (
5138
operation.state == OperationState.CREATED
@@ -59,18 +46,13 @@ def wrapper(self, continuation, client_call_details, request):
5946
return wrapper
6047

6148

62-
def _end_attempt(operation, exc, metadata):
63-
"""Helper to add metadata and exception to an operation"""
64-
if metadata is not None:
65-
operation.add_response_metadata(metadata)
66-
if exc is not None:
67-
operation.end_attempt_with_status(exc)
68-
69-
70-
def _get_metadata(source):
49+
def _get_metadata(source) -> dict[str, str | bytes] | None:
7150
"""Helper to extract metadata from a call or RpcError"""
7251
try:
73-
return (source.trailing_metadata() or []) + (source.initial_metadata() or [])
52+
metadata: Sequence[tuple[str.str | bytes]] = (
53+
source.trailing_metadata() + source.initial_metadata()
54+
)
55+
return {k: v for (k, v) in metadata}
7456
except Exception:
7557
return None
7658

@@ -82,53 +64,31 @@ class BigtableMetricsInterceptor(
8264
An async gRPC interceptor to add client metadata and print server metadata.
8365
"""
8466

85-
def __init__(self):
86-
super().__init__()
87-
self.operation_map = {}
88-
89-
def register_operation(self, operation):
90-
"""Register an operation object to be tracked my the interceptor
91-
92-
When registered, the operation will receive metadata updates:
93-
- start_attempt if attempt not started when rpc is being sent
94-
- add_response_metadata after call is complete
95-
- end_attempt_with_status if attempt receives an error
96-
97-
The interceptor will register itself as a handeler for the operation,
98-
so it can unregister the operation when it is complete"""
99-
self.operation_map[operation.uuid] = operation
100-
operation.handlers.append(self)
101-
102-
def on_operation_complete(self, op):
103-
if op.uuid in self.operation_map:
104-
del self.operation_map[op.uuid]
105-
106-
def on_operation_cancelled(self, op):
107-
self.on_operation_complete(op)
108-
10967
@_with_operation_from_metadata
11068
def intercept_unary_unary(
11169
self, operation, continuation, client_call_details, request
11270
):
113-
encountered_exc: Exception | None = None
11471
metadata = None
11572
try:
11673
call = continuation(client_call_details, request)
11774
metadata = _get_metadata(call)
11875
return call
11976
except Exception as rpc_error:
12077
metadata = _get_metadata(rpc_error)
121-
encountered_exc = rpc_error
12278
raise rpc_error
12379
finally:
124-
_end_attempt(operation, encountered_exc, metadata)
80+
if metadata is not None:
81+
operation.add_response_metadata(metadata)
12582

12683
@_with_operation_from_metadata
12784
def intercept_unary_stream(
12885
self, operation, continuation, client_call_details, request
12986
):
13087
def response_wrapper(call):
131-
has_first_response = operation.first_response_latency is not None
88+
has_first_response = (
89+
operation.first_response_latency_ns is not None
90+
or operation.op_type != OperationType.READ_ROWS
91+
)
13292
encountered_exc = None
13393
try:
13494
for response in call:
@@ -143,10 +103,14 @@ def response_wrapper(call):
143103
raise
144104
finally:
145105
if call is not None:
146-
_end_attempt(operation, encountered_exc, _get_metadata(call))
106+
metadata = _get_metadata(encountered_exc or call)
107+
if metadata is not None:
108+
operation.add_response_metadata(metadata)
147109

148110
try:
149111
return response_wrapper(continuation(client_call_details, request))
150112
except Exception as rpc_error:
151-
_end_attempt(operation, rpc_error, _get_metadata(rpc_error))
113+
metadata = _get_metadata(rpc_error)
114+
if metadata is not None:
115+
operation.add_response_metadata(metadata)
152116
raise rpc_error

google/cloud/bigtable/data/_sync_autogen/mutations_batcher.py

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -151,7 +151,7 @@ def add_to_flow(self, mutations: RowMutationEntry | list[RowMutationEntry]):
151151
)
152152
yield mutations[start_idx:end_idx]
153153

154-
async def add_to_flow_with_metrics(
154+
def add_to_flow_with_metrics(
155155
self,
156156
mutations: RowMutationEntry | list[RowMutationEntry],
157157
metrics_controller: BigtableClientSideMetricsController,
@@ -161,9 +161,8 @@ async def add_to_flow_with_metrics(
161161
metric = metrics_controller.create_operation(OperationType.BULK_MUTATE_ROWS)
162162
flow_start_time = time.monotonic_ns()
163163
try:
164-
value = await inner_generator.__anext__()
165-
except StopAsyncIteration:
166-
metric.cancel()
164+
value = inner_generator.__next__()
165+
except CrossSync._Sync_Impl.StopIteration:
167166
return
168167
metric.flow_throttling_time_ns = time.monotonic_ns() - flow_start_time
169168
yield (value, metric)

0 commit comments

Comments
 (0)