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

Commit 64ce0a4

Browse files
committed
improved retry instrumentation
1 parent b91da1c commit 64ce0a4

5 files changed

Lines changed: 17 additions & 80 deletions

File tree

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@
2121
import google.cloud.bigtable_v2.types.bigtable as types_pb
2222
import google.cloud.bigtable.data.exceptions as bt_exceptions
2323
from google.cloud.bigtable.data._helpers import _attempt_timeout_generator
24-
from google.cloud.bigtable.data._helpers import _tracked_exception_factory
24+
from google.cloud.bigtable.data._helpers import _retry_exception_factory
2525

2626
# mutate_rows requests are limited to this number of mutations
2727
from google.cloud.bigtable.data.mutations import _MUTATE_ROWS_REQUEST_MUTATION_LIMIT
@@ -106,7 +106,8 @@ def __init__(
106106
self.is_retryable,
107107
metric.backoff_generator,
108108
operation_timeout,
109-
exception_factory=_tracked_exception_factory(self._operation_metric),
109+
exception_factory=self._operation_metric.track_terminal_error(_retry_exception_factory),
110+
on_error=self._operation_metric.track_retryable_error(),
110111
)
111112
# initialize state
112113
self.timeout_generator = _attempt_timeout_generator(

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,7 @@
3232
from google.cloud.bigtable.data.exceptions import _RowSetComplete
3333
from google.cloud.bigtable.data.exceptions import _ResetRow
3434
from google.cloud.bigtable.data._helpers import _attempt_timeout_generator
35-
from google.cloud.bigtable.data._helpers import _tracked_exception_factory
35+
from google.cloud.bigtable.data._helpers import _retry_exception_factory
3636

3737
from google.api_core import retry as retries
3838

@@ -123,7 +123,8 @@ def start_operation(self) -> CrossSync.Iterable[Row]:
123123
self._predicate,
124124
self._operation_metric.backoff_generator,
125125
self.operation_timeout,
126-
exception_factory=_tracked_exception_factory(self._operation_metric),
126+
exception_factory=self._operation_metric.track_terminal_error(_retry_exception_factory),
127+
on_error=self._operation_metric.track_retryable_error(),
127128
)
128129

129130
def _read_rows_attempt(self) -> CrossSync.Iterable[Row]:

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

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -73,7 +73,6 @@
7373
from google.cloud.bigtable.data._helpers import _WarmedInstanceKey
7474
from google.cloud.bigtable.data._helpers import _CONCURRENCY_LIMIT
7575
from google.cloud.bigtable.data._helpers import _retry_exception_factory
76-
from google.cloud.bigtable.data._helpers import _tracked_exception_factory
7776
from google.cloud.bigtable.data._helpers import _validate_timeouts
7877
from google.cloud.bigtable.data._helpers import _get_error_type
7978
from google.cloud.bigtable.data._helpers import _get_retryable_errors
@@ -1392,7 +1391,8 @@ async def execute_rpc():
13921391
predicate,
13931392
operation_metric.backoff_generator,
13941393
operation_timeout,
1395-
exception_factory=_tracked_exception_factory(operation_metric),
1394+
exception_factory=operation_metric.track_terminal_error(_retry_exception_factory),
1395+
on_error=operation_metric.track_retryable_error(),
13961396
)
13971397

13981398
@CrossSync.convert(replace_symbols={"MutationsBatcherAsync": "MutationsBatcher"})
@@ -1526,7 +1526,8 @@ async def mutate_row(
15261526
predicate,
15271527
operation_metric.backoff_generator,
15281528
operation_timeout,
1529-
exception_factory=_retry_exception_factory,
1529+
exception_factory=operation_metric.track_terminal_error(_retry_exception_factory),
1530+
on_error=operation_metric.track_retryable_error(),
15301531
)
15311532

15321533
@CrossSync.convert

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

Lines changed: 7 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -80,16 +80,6 @@ def wrapper(self, continuation, client_call_details, request):
8080

8181
return wrapper
8282

83-
84-
def _end_attempt(operation, exc, metadata):
85-
"""Helper to add metadata and exception to an operation"""
86-
if metadata is not None:
87-
operation.add_response_metadata(metadata)
88-
if exc is not None:
89-
# end attempt. If it succeeded, let higher levels decide when to end operation
90-
operation.end_attempt_with_status(exc)
91-
92-
9383
@CrossSync.convert
9484
async def _get_metadata(source) -> dict[str, str|bytes] | None:
9585
"""Helper to extract metadata from a call or RpcError"""
@@ -129,7 +119,6 @@ def register_operation(self, operation):
129119
When registered, the operation will receive metadata updates:
130120
- start_attempt if attempt not started when rpc is being sent
131121
- add_response_metadata after call is complete
132-
- end_attempt_with_status if attempt receives an error
133122
134123
The interceptor will register itself as a handeler for the operation,
135124
so it can unregister the operation when it is complete
@@ -149,23 +138,17 @@ def on_operation_cancelled(self, op):
149138
async def intercept_unary_unary(
150139
self, operation, continuation, client_call_details, request
151140
):
152-
encountered_status: Exception | StatusCode | None = None
153141
metadata = None
154142
try:
155143
call = await continuation(client_call_details, request)
156144
metadata = await _get_metadata(call)
157-
if CrossSync.is_async:
158-
encountered_status = await call.code()
159-
elif isinstance(call, Exception):
160-
# sync unary calls return exception objects without raising
161-
encountered_status = call
162145
return call
163146
except Exception as rpc_error:
164147
metadata = await _get_metadata(rpc_error)
165-
encountered_status = rpc_error
166148
raise rpc_error
167149
finally:
168-
_end_attempt(operation, encountered_status, metadata)
150+
if metadata is not None:
151+
operation.add_response_metadata(metadata)
169152

170153
@CrossSync.convert
171154
@_with_operation_from_metadata
@@ -195,11 +178,13 @@ async def response_wrapper(call):
195178
finally:
196179
if call is not None:
197180
metadata = await _get_metadata(encountered_exc or call)
198-
_end_attempt(operation, encountered_exc, metadata)
181+
if metadata is not None:
182+
operation.add_response_metadata(metadata)
199183

200184
try:
201185
return response_wrapper(await continuation(client_call_details, request))
202186
except Exception as rpc_error:
203-
# handle errors while intializing stream
204-
_end_attempt(operation, rpc_error, await _get_metadata(rpc_error))
187+
metadata = await _get_metadata(rpc_error)
188+
if metadata is not None:
189+
operation.add_response_metadata(metadata)
205190
raise rpc_error

google/cloud/bigtable/data/_helpers.py

Lines changed: 0 additions & 51 deletions
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,6 @@
3131
import grpc
3232
from google.cloud.bigtable.data._async.client import _DataApiTargetAsync
3333
from google.cloud.bigtable.data._sync_autogen.client import _DataApiTarget
34-
from google.cloud.bigtable.data._metrics.data_model import ActiveOperationMetric
3534

3635
"""
3736
Helper functions used in various places in the library.
@@ -121,56 +120,6 @@ def _retry_exception_factory(
121120
return source_exc, cause_exc
122121

123122

124-
def _tracked_exception_factory(
125-
operation: "ActiveOperationMetric",
126-
) -> Callable[
127-
[list[Exception], RetryFailureReason, float | None],
128-
tuple[Exception, Exception | None],
129-
]:
130-
"""
131-
wraps and extends _retry_exception_factory to add client-side metrics tracking.
132-
133-
When the rpc raises a terminal error, record any discovered metadata and finalize
134-
the associated operation
135-
136-
Used by streaming rpcs, which can't always be perfectly captured by context managers
137-
for operation termination.
138-
139-
Args:
140-
exc_list: list of exceptions encountered during operation
141-
is_timeout: whether the operation failed due to timeout
142-
timeout_val: the operation timeout value in seconds, for constructing
143-
the error message
144-
operation: the operation to finalize when an exception is built
145-
Returns:
146-
tuple[Exception, Exception|None]:
147-
tuple of the exception to raise, and a cause exception if applicabl
148-
"""
149-
150-
def wrapper(
151-
exc_list: list[Exception], reason: RetryFailureReason, timeout_val: float | None
152-
) -> tuple[Exception, Exception | None]:
153-
source_exc, cause_exc = _retry_exception_factory(exc_list, reason, timeout_val)
154-
try:
155-
# record metadata from failed rpc
156-
if (
157-
isinstance(source_exc, core_exceptions.GoogleAPICallError)
158-
and source_exc.errors
159-
):
160-
rpc_error = source_exc.errors[-1]
161-
metadata = list(rpc_error.trailing_metadata()) + list(
162-
rpc_error.initial_metadata()
163-
)
164-
operation.add_response_metadata({k: v for k, v in metadata})
165-
except Exception:
166-
# ignore errors in metadata collection
167-
pass
168-
operation.end_with_status(source_exc)
169-
return source_exc, cause_exc
170-
171-
return wrapper
172-
173-
174123
def _get_timeouts(
175124
operation: float | TABLE_DEFAULT,
176125
attempt: float | None | TABLE_DEFAULT,

0 commit comments

Comments
 (0)