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

Commit 9fe5d4e

Browse files
committed
use new tracked_retry function
1 parent 377044a commit 9fe5d4e

11 files changed

Lines changed: 62 additions & 84 deletions

File tree

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

Lines changed: 7 additions & 8 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 _retry_exception_factory
24+
from google.cloud.bigtable.data._metrics import tracked_retry
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
@@ -101,13 +101,12 @@ def __init__(
101101
# Entry level errors
102102
bt_exceptions._MutateRowsIncomplete,
103103
)
104-
self._operation = lambda: CrossSync.retry_target(
105-
self._run_attempt,
106-
self.is_retryable,
107-
metric.backoff_generator,
108-
operation_timeout,
109-
exception_factory=metric.track_terminal_error(_retry_exception_factory),
110-
on_error=metric.track_retryable_error,
104+
self._operation = lambda: tracked_retry(
105+
retry_fn=CrossSync.retry_target,
106+
operation=metric,
107+
target=self._run_attempt,
108+
predicate=self.is_retryable,
109+
timeout=operation_timeout,
111110
)
112111
# initialize state
113112
self.timeout_generator = _attempt_timeout_generator(

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

Lines changed: 7 additions & 10 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 _retry_exception_factory
35+
from google.cloud.bigtable.data._metrics import tracked_retry
3636

3737
from google.api_core import retry as retries
3838

@@ -118,15 +118,12 @@ def start_operation(self) -> CrossSync.Iterable[Row]:
118118
Yields:
119119
Row: The next row in the stream
120120
"""
121-
return CrossSync.retry_target_stream(
122-
self._read_rows_attempt,
123-
self._predicate,
124-
self._operation_metric.backoff_generator,
125-
self.operation_timeout,
126-
exception_factory=self._operation_metric.track_terminal_error(
127-
_retry_exception_factory
128-
),
129-
on_error=self._operation_metric.track_retryable_error,
121+
return tracked_retry(
122+
retry_fn=CrossSync.retry_target_stream,
123+
operation=self._operation_metric,
124+
target=self._read_rows_attempt,
125+
predicate=self._predicate,
126+
timeout=self.operation_timeout,
130127
)
131128

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

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

Lines changed: 13 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -89,6 +89,7 @@
8989
from google.cloud.bigtable.data.row_filters import RowFilterChain
9090
from google.cloud.bigtable.data._metrics import BigtableClientSideMetricsController
9191
from google.cloud.bigtable.data._metrics import OperationType
92+
from google.cloud.bigtable.data._metrics import tracked_retry
9293

9394
from google.cloud.bigtable.data._cross_sync import CrossSync
9495

@@ -1448,15 +1449,12 @@ async def execute_rpc():
14481449
)
14491450
return [(s.row_key, s.offset_bytes) async for s in results]
14501451

1451-
return await CrossSync.retry_target(
1452-
execute_rpc,
1453-
predicate,
1454-
operation_metric.backoff_generator,
1455-
operation_timeout,
1456-
exception_factory=operation_metric.track_terminal_error(
1457-
_retry_exception_factory
1458-
),
1459-
on_error=operation_metric.track_retryable_error,
1452+
return await tracked_retry(
1453+
retry_fn=CrossSync.retry_target,
1454+
operation=operation_metric,
1455+
target=execute_rpc,
1456+
predicate=predicate,
1457+
timeout=operation_timeout,
14601458
)
14611459

14621460
@CrossSync.convert(replace_symbols={"MutationsBatcherAsync": "MutationsBatcher"})
@@ -1584,15 +1582,12 @@ async def mutate_row(
15841582
timeout=attempt_timeout,
15851583
retry=None,
15861584
)
1587-
return await CrossSync.retry_target(
1588-
target,
1589-
predicate,
1590-
operation_metric.backoff_generator,
1591-
operation_timeout,
1592-
exception_factory=operation_metric.track_terminal_error(
1593-
_retry_exception_factory
1594-
),
1595-
on_error=operation_metric.track_retryable_error,
1585+
return await tracked_retry(
1586+
retry_fn=CrossSync.retry_target,
1587+
operation=operation_metric,
1588+
target=target,
1589+
predicate=predicate,
1590+
timeout=operation_timeout,
15961591
)
15971592

15981593
@CrossSync.convert

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

Lines changed: 7 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@
2222
import google.cloud.bigtable_v2.types.bigtable as types_pb
2323
import google.cloud.bigtable.data.exceptions as bt_exceptions
2424
from google.cloud.bigtable.data._helpers import _attempt_timeout_generator
25-
from google.cloud.bigtable.data._helpers import _retry_exception_factory
25+
from google.cloud.bigtable.data._metrics import tracked_retry
2626
from google.cloud.bigtable.data.mutations import _MUTATE_ROWS_REQUEST_MUTATION_LIMIT
2727
from google.cloud.bigtable.data.mutations import _EntryWithProto
2828
from google.cloud.bigtable.data._cross_sync import CrossSync
@@ -79,13 +79,12 @@ def __init__(
7979
self.is_retryable = retries.if_exception_type(
8080
*retryable_exceptions, bt_exceptions._MutateRowsIncomplete
8181
)
82-
self._operation = lambda: CrossSync._Sync_Impl.retry_target(
83-
self._run_attempt,
84-
self.is_retryable,
85-
metric.backoff_generator,
86-
operation_timeout,
87-
exception_factory=metric.track_terminal_error(_retry_exception_factory),
88-
on_error=metric.track_retryable_error,
82+
self._operation = lambda: tracked_retry(
83+
retry_fn=CrossSync._Sync_Impl.retry_target,
84+
operation=metric,
85+
target=self._run_attempt,
86+
predicate=self.is_retryable,
87+
timeout=operation_timeout,
8988
)
9089
self.timeout_generator = _attempt_timeout_generator(
9190
attempt_timeout, operation_timeout

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

Lines changed: 7 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -30,7 +30,7 @@
3030
from google.cloud.bigtable.data.exceptions import _RowSetComplete
3131
from google.cloud.bigtable.data.exceptions import _ResetRow
3232
from google.cloud.bigtable.data._helpers import _attempt_timeout_generator
33-
from google.cloud.bigtable.data._helpers import _retry_exception_factory
33+
from google.cloud.bigtable.data._metrics import tracked_retry
3434
from google.api_core import retry as retries
3535
from google.cloud.bigtable.data._cross_sync import CrossSync
3636

@@ -103,15 +103,12 @@ def start_operation(self) -> CrossSync._Sync_Impl.Iterable[Row]:
103103
104104
Yields:
105105
Row: The next row in the stream"""
106-
return CrossSync._Sync_Impl.retry_target_stream(
107-
self._read_rows_attempt,
108-
self._predicate,
109-
self._operation_metric.backoff_generator,
110-
self.operation_timeout,
111-
exception_factory=self._operation_metric.track_terminal_error(
112-
_retry_exception_factory
113-
),
114-
on_error=self._operation_metric.track_retryable_error,
106+
return tracked_retry(
107+
retry_fn=CrossSync._Sync_Impl.retry_target_stream,
108+
operation=self._operation_metric,
109+
target=self._read_rows_attempt,
110+
predicate=self._predicate,
111+
timeout=self.operation_timeout,
115112
)
116113

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

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

Lines changed: 13 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -76,6 +76,7 @@
7676
from google.cloud.bigtable.data.row_filters import RowFilterChain
7777
from google.cloud.bigtable.data._metrics import BigtableClientSideMetricsController
7878
from google.cloud.bigtable.data._metrics import OperationType
79+
from google.cloud.bigtable.data._metrics import tracked_retry
7980
from google.cloud.bigtable.data._cross_sync import CrossSync
8081
from typing import Iterable
8182
from grpc import insecure_channel
@@ -1196,15 +1197,12 @@ def execute_rpc():
11961197
)
11971198
return [(s.row_key, s.offset_bytes) for s in results]
11981199

1199-
return CrossSync._Sync_Impl.retry_target(
1200-
execute_rpc,
1201-
predicate,
1202-
operation_metric.backoff_generator,
1203-
operation_timeout,
1204-
exception_factory=operation_metric.track_terminal_error(
1205-
_retry_exception_factory
1206-
),
1207-
on_error=operation_metric.track_retryable_error,
1200+
return tracked_retry(
1201+
retry_fn=CrossSync._Sync_Impl.retry_target,
1202+
operation=operation_metric,
1203+
target=execute_rpc,
1204+
predicate=predicate,
1205+
timeout=operation_timeout,
12081206
)
12091207

12101208
def mutations_batcher(
@@ -1322,15 +1320,12 @@ def mutate_row(
13221320
timeout=attempt_timeout,
13231321
retry=None,
13241322
)
1325-
return CrossSync._Sync_Impl.retry_target(
1326-
target,
1327-
predicate,
1328-
operation_metric.backoff_generator,
1329-
operation_timeout,
1330-
exception_factory=operation_metric.track_terminal_error(
1331-
_retry_exception_factory
1332-
),
1333-
on_error=operation_metric.track_retryable_error,
1323+
return tracked_retry(
1324+
retry_fn=CrossSync._Sync_Impl.retry_target,
1325+
operation=operation_metric,
1326+
target=target,
1327+
predicate=predicate,
1328+
timeout=operation_timeout,
13341329
)
13351330

13361331
def bulk_mutate_rows(

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

Lines changed: 0 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -58,10 +58,6 @@ def _get_metadata(source) -> dict[str, str | bytes] | None:
5858
class BigtableMetricsInterceptor(
5959
UnaryUnaryClientInterceptor, UnaryStreamClientInterceptor
6060
):
61-
"""
62-
An async gRPC interceptor to add client metadata and print server metadata.
63-
"""
64-
6561
@_with_active_operation
6662
def intercept_unary_unary(
6763
self, operation, continuation, client_call_details, request

tests/unit/data/_async/test_client.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1400,9 +1400,9 @@ async def test_customizable_retryable_errors(
14001400
predicate_builder_mock.assert_called_once_with(
14011401
*expected_retryables, *extra_retryables
14021402
)
1403-
retry_call_args = retry_fn_mock.call_args_list[0].args
1403+
retry_call_kwargs = retry_fn_mock.call_args_list[0].kwargs
14041404
# output of if_exception_type should be sent in to retry constructor
1405-
assert retry_call_args[1] is expected_predicate
1405+
assert retry_call_kwargs["predicate"] is expected_predicate
14061406

14071407
@pytest.mark.parametrize(
14081408
"fn_name,fn_args,gapic_fn",

tests/unit/data/_async/test_mutations_batcher.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1223,9 +1223,9 @@ async def test_customizable_retryable_errors(
12231223
predicate_builder_mock.assert_called_once_with(
12241224
*expected_retryables, _MutateRowsIncomplete
12251225
)
1226-
retry_call_args = retry_fn_mock.call_args_list[0].args
1226+
retry_call_kwargs = retry_fn_mock.call_args_list[0].kwargs
12271227
# output of if_exception_type should be sent in to retry constructor
1228-
assert retry_call_args[1] is expected_predicate
1228+
assert retry_call_kwargs["predicate"] is expected_predicate
12291229

12301230
@CrossSync.pytest
12311231
async def test_large_batch_write(self):

tests/unit/data/_sync_autogen/test_client.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1120,8 +1120,8 @@ def test_customizable_retryable_errors(
11201120
predicate_builder_mock.assert_called_once_with(
11211121
*expected_retryables, *extra_retryables
11221122
)
1123-
retry_call_args = retry_fn_mock.call_args_list[0].args
1124-
assert retry_call_args[1] is expected_predicate
1123+
retry_call_kwargs = retry_fn_mock.call_args_list[0].kwargs
1124+
assert retry_call_kwargs["predicate"] is expected_predicate
11251125

11261126
@pytest.mark.parametrize(
11271127
"fn_name,fn_args,gapic_fn",

0 commit comments

Comments
 (0)