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

Commit 410cfb8

Browse files
committed
fixed mypy
1 parent 40fcbe8 commit 410cfb8

5 files changed

Lines changed: 60 additions & 70 deletions

File tree

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

Lines changed: 5 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -220,7 +220,7 @@ async def chunk_stream(
220220
)
221221
async def merge_rows(
222222
chunks: CrossSync.Iterable[ReadRowsResponsePB.CellChunk] | None,
223-
attempt_metric: ActiveAttemptMetric,
223+
attempt_metric: ActiveAttemptMetric | None,
224224
) -> CrossSync.Iterable[Row]:
225225
"""
226226
Merge chunks into rows
@@ -233,7 +233,6 @@ async def merge_rows(
233233
if chunks is None:
234234
return
235235
it = chunks.__aiter__()
236-
is_first_row = True
237236
# For each row
238237
while True:
239238
try:
@@ -316,17 +315,14 @@ async def merge_rows(
316315
Cell(value, row_key, family, qualifier, ts, list(labels))
317316
)
318317
if c.commit_row:
319-
if is_first_row:
320-
# record first row latency in metrics
321-
is_first_row = False
322-
attempt_metric.attempt_first_response()
323318
block_time = time.monotonic_ns()
324319
yield Row(row_key, cells)
325320
# most metric operations use setters, but this one updates
326321
# the value directly to avoid extra overhead
327-
attempt_metric.active_attempt.application_blocking_time_ns += ( # type: ignore
328-
time.monotonic_ns() - block_time
329-
) * 1000
322+
if attempt_metric is not None:
323+
attempt_metric.application_blocking_time_ns += ( # type: ignore
324+
time.monotonic_ns() - block_time
325+
) * 1000
330326
break
331327
c = await it.__anext__()
332328
except _ResetRow as e:

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

Lines changed: 1 addition & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,6 @@
2323
import google.cloud.bigtable.data.exceptions as bt_exceptions
2424
from google.cloud.bigtable.data._helpers import _attempt_timeout_generator
2525
from google.cloud.bigtable.data._helpers import _retry_exception_factory
26-
from google.cloud.bigtable.data._helpers import TrackedBackoffGenerator
2726
from google.cloud.bigtable.data.mutations import _MUTATE_ROWS_REQUEST_MUTATION_LIMIT
2827
from google.cloud.bigtable.data.mutations import _EntryWithProto
2928
from google.cloud.bigtable.data._cross_sync import CrossSync
@@ -80,11 +79,10 @@ def __init__(
8079
self.is_retryable = retries.if_exception_type(
8180
*retryable_exceptions, bt_exceptions._MutateRowsIncomplete
8281
)
83-
sleep_generator = TrackedBackoffGenerator(0.01, 2, 60)
8482
self._operation = lambda: CrossSync._Sync_Impl.retry_target(
8583
self._run_attempt,
8684
self.is_retryable,
87-
sleep_generator,
85+
metric.backoff_generator,
8886
operation_timeout,
8987
exception_factory=_retry_exception_factory,
9088
)
@@ -94,7 +92,6 @@ def __init__(
9492
self.mutations = [_EntryWithProto(m, m._to_pb()) for m in mutation_entries]
9593
self.remaining_indices = list(range(len(self.mutations)))
9694
self.errors: dict[int, list[Exception]] = {}
97-
metric.backoff_generator = sleep_generator
9895
self._operation_metric = metric
9996

10097
def start(self):

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

Lines changed: 13 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -30,7 +30,6 @@
3030
from google.cloud.bigtable.data.exceptions import _ResetRow
3131
from google.cloud.bigtable.data._helpers import _attempt_timeout_generator
3232
from google.cloud.bigtable.data._helpers import _retry_exception_factory
33-
from google.cloud.bigtable.data._helpers import TrackedBackoffGenerator
3433
from google.api_core import retry as retries
3534
from google.cloud.bigtable.data._cross_sync import CrossSync
3635

@@ -104,16 +103,14 @@ def start_operation(self) -> CrossSync._Sync_Impl.Iterable[Row]:
104103
105104
Yields:
106105
Row: The next row in the stream"""
107-
self._operation_metric.backoff_generator = TrackedBackoffGenerator(
108-
0.01, 60, multiplier=2
109-
)
110-
return CrossSync._Sync_Impl.retry_target_stream(
111-
self._read_rows_attempt,
112-
self._predicate,
113-
self._operation_metric.backoff_generator,
114-
self.operation_timeout,
115-
exception_factory=_retry_exception_factory,
116-
)
106+
with self._operation_metric:
107+
return CrossSync._Sync_Impl.retry_target_stream(
108+
self._read_rows_attempt,
109+
self._predicate,
110+
self._operation_metric.backoff_generator,
111+
self.operation_timeout,
112+
exception_factory=_retry_exception_factory,
113+
)
117114

118115
def _read_rows_attempt(self) -> CrossSync._Sync_Impl.Iterable[Row]:
119116
"""Attempt a single read_rows rpc call.
@@ -188,7 +185,7 @@ def chunk_stream(
188185
@staticmethod
189186
def merge_rows(
190187
chunks: CrossSync._Sync_Impl.Iterable[ReadRowsResponsePB.CellChunk] | None,
191-
attempt_metric: ActiveAttemptMetric,
188+
attempt_metric: ActiveAttemptMetric | None,
192189
) -> CrossSync._Sync_Impl.Iterable[Row]:
193190
"""Merge chunks into rows
194191
@@ -199,7 +196,6 @@ def merge_rows(
199196
if chunks is None:
200197
return
201198
it = chunks.__iter__()
202-
is_first_row = True
203199
while True:
204200
try:
205201
c = it.__next__()
@@ -268,14 +264,12 @@ def merge_rows(
268264
Cell(value, row_key, family, qualifier, ts, list(labels))
269265
)
270266
if c.commit_row:
271-
if is_first_row:
272-
is_first_row = False
273-
attempt_metric.attempt_first_response()
274267
block_time = time.monotonic_ns()
275268
yield Row(row_key, cells)
276-
attempt_metric.active_attempt.application_blocking_time_ns += (
277-
time.monotonic_ns() - block_time
278-
) * 1000
269+
if attempt_metric is not None:
270+
attempt_metric.application_blocking_time_ns += (
271+
time.monotonic_ns() - block_time
272+
) * 1000
279273
break
280274
c = it.__next__()
281275
except _ResetRow as e:

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

Lines changed: 32 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -66,7 +66,6 @@
6666
from google.cloud.bigtable.data._helpers import _get_retryable_errors
6767
from google.cloud.bigtable.data._helpers import _get_timeouts
6868
from google.cloud.bigtable.data._helpers import _attempt_timeout_generator
69-
from google.cloud.bigtable.data._helpers import TrackedBackoffGenerator
7069
from google.cloud.bigtable.data.mutations import Mutation, RowMutationEntry
7170
from google.cloud.bigtable.data.read_modify_write_rules import ReadModifyWriteRule
7271
from google.cloud.bigtable.data.row_filters import RowFilter
@@ -805,18 +804,17 @@ def read_rows_stream(
805804
operation_timeout, attempt_timeout, self
806805
)
807806
retryable_excs = _get_retryable_errors(retryable_errors, self)
808-
with self._metrics.create_operation(
809-
OperationType.READ_ROWS, streaming=True
810-
) as operation_metric:
811-
row_merger = CrossSync._Sync_Impl._ReadRowsOperation(
812-
query,
813-
self,
814-
operation_timeout=operation_timeout,
815-
attempt_timeout=attempt_timeout,
816-
metric=operation_metric,
817-
retryable_exceptions=retryable_excs,
818-
)
819-
return row_merger.start_operation()
807+
row_merger = CrossSync._Sync_Impl._ReadRowsOperation(
808+
query,
809+
self,
810+
operation_timeout=operation_timeout,
811+
attempt_timeout=attempt_timeout,
812+
metric=self._metrics.create_operation(
813+
OperationType.READ_ROWS, streaming=True
814+
),
815+
retryable_exceptions=retryable_excs,
816+
)
817+
return row_merger.start_operation()
820818

821819
def read_rows(
822820
self,
@@ -906,22 +904,21 @@ def read_row(
906904
operation_timeout, attempt_timeout, self
907905
)
908906
retryable_excs = _get_retryable_errors(retryable_errors, self)
909-
with self._metrics.create_operation(
910-
OperationType.READ_ROWS, streaming=False
911-
) as operation_metric:
912-
row_merger = CrossSync._Sync_Impl._ReadRowsOperation(
913-
query,
914-
self,
915-
operation_timeout=operation_timeout,
916-
attempt_timeout=attempt_timeout,
917-
metric=operation_metric,
918-
retryable_exceptions=retryable_excs,
919-
)
920-
results_generator = row_merger.start_operation()
921-
try:
922-
return results_generator.__next__()
923-
except StopIteration:
924-
return None
907+
row_merger = CrossSync._Sync_Impl._ReadRowsOperation(
908+
query,
909+
self,
910+
operation_timeout=operation_timeout,
911+
attempt_timeout=attempt_timeout,
912+
metric=self._metrics.create_operation(
913+
OperationType.READ_ROWS, streaming=False
914+
),
915+
retryable_exceptions=retryable_excs,
916+
)
917+
results_generator = row_merger.start_operation()
918+
try:
919+
return results_generator.__next__()
920+
except StopIteration:
921+
return None
925922

926923
def read_rows_sharded(
927924
self,
@@ -1101,10 +1098,9 @@ def sample_row_keys(
11011098
)
11021099
retryable_excs = _get_retryable_errors(retryable_errors, self)
11031100
predicate = retries.if_exception_type(*retryable_excs)
1104-
sleep_generator = TrackedBackoffGenerator(0.01, 2, 60)
11051101
with self._metrics.create_operation(
1106-
OperationType.SAMPLE_ROW_KEYS, backoff_generator=sleep_generator
1107-
):
1102+
OperationType.SAMPLE_ROW_KEYS
1103+
) as operation_metric:
11081104

11091105
def execute_rpc():
11101106
results = self.client._gapic_client.sample_row_keys(
@@ -1119,7 +1115,7 @@ def execute_rpc():
11191115
return CrossSync._Sync_Impl.retry_target(
11201116
execute_rpc,
11211117
predicate,
1122-
sleep_generator,
1118+
operation_metric.backoff_generator,
11231119
operation_timeout,
11241120
exception_factory=_retry_exception_factory,
11251121
)
@@ -1223,10 +1219,9 @@ def mutate_row(
12231219
)
12241220
else:
12251221
predicate = retries.if_exception_type()
1226-
sleep_generator = TrackedBackoffGenerator(0.01, 2, 60)
12271222
with self._metrics.create_operation(
1228-
OperationType.MUTATE_ROW, backoff_generator=sleep_generator
1229-
):
1223+
OperationType.MUTATE_ROW
1224+
) as operation_metric:
12301225
target = partial(
12311226
self.client._gapic_client.mutate_row,
12321227
request=MutateRowRequest(
@@ -1243,7 +1238,7 @@ def mutate_row(
12431238
return CrossSync._Sync_Impl.retry_target(
12441239
target,
12451240
predicate,
1246-
sleep_generator,
1241+
operation_metric.backoff_generator,
12471242
operation_timeout,
12481243
exception_factory=_retry_exception_factory,
12491244
)

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

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@
1515
# This file is automatically generated by CrossSync. Do not edit manually.
1616

1717
from __future__ import annotations
18+
import time
1819
from functools import wraps
1920
from google.cloud.bigtable.data._metrics.data_model import (
2021
OPERATION_INTERCEPTOR_METADATA_KEY,
@@ -76,7 +77,8 @@ def register_operation(self, operation):
7677
operation.handlers.append(self)
7778

7879
def on_operation_complete(self, op):
79-
del self.operation_map[op.uuid]
80+
if op.uuid in self.operation_map:
81+
del self.operation_map[op.uuid]
8082

8183
def on_operation_cancelled(self, op):
8284
self.on_operation_complete(op)
@@ -105,9 +107,15 @@ def intercept_unary_stream(
105107
self, operation, continuation, client_call_details, request
106108
):
107109
def response_wrapper(call):
110+
has_first_response = operation.first_response_latency is not None
108111
encountered_exc = None
109112
try:
110113
for response in call:
114+
if not has_first_response:
115+
operation.first_response_latency_ns = (
116+
time.monotonic_ns() - operation.start_time_ns
117+
)
118+
has_first_response = True
111119
yield response
112120
except Exception as e:
113121
encountered_exc = e

0 commit comments

Comments
 (0)