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

Commit bb55c46

Browse files
committed
made merge_rows into an instance method
1 parent ab30b02 commit bb55c46

5 files changed

Lines changed: 26 additions & 33 deletions

File tree

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

Lines changed: 12 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,7 @@
1515

1616
from __future__ import annotations
1717

18-
from typing import Callable, Sequence, TYPE_CHECKING
18+
from typing import Sequence, TYPE_CHECKING
1919

2020
import time
2121

@@ -145,20 +145,20 @@ def _read_rows_attempt(self) -> CrossSync.Iterable[Row]:
145145
)
146146
except _RowSetComplete:
147147
# if we've already seen all the rows, we're done
148-
return self.merge_rows(None, self._operation_metric, self._predicate)
148+
return self.merge_rows(None)
149149
# revise the limit based on number of rows already yielded
150150
if self._remaining_count is not None:
151151
self.request.rows_limit = self._remaining_count
152152
if self._remaining_count == 0:
153-
return self.merge_rows(None, self._operation_metric, self._predicate)
153+
return self.merge_rows(None)
154154
# create and return a new row merger
155155
gapic_stream = self.target.client._gapic_client.read_rows(
156156
self.request,
157157
timeout=next(self.attempt_timeout_gen),
158158
retry=None,
159159
)
160160
chunked_stream = self.chunk_stream(gapic_stream)
161-
return self.merge_rows(chunked_stream, self._operation_metric, self._predicate)
161+
return self.merge_rows(chunked_stream)
162162

163163
@CrossSync.convert()
164164
async def chunk_stream(
@@ -217,9 +217,7 @@ async def chunk_stream(
217217
replace_symbols={"__aiter__": "__iter__", "__anext__": "__next__"},
218218
)
219219
async def merge_rows(
220-
chunks: CrossSync.Iterable[ReadRowsResponsePB.CellChunk] | None,
221-
operation_metric: ActiveOperationMetric,
222-
retryable_predicate: Callable[[Exception], bool],
220+
self, chunks: CrossSync.Iterable[ReadRowsResponsePB.CellChunk] | None
223221
) -> CrossSync.Iterable[Row]:
224222
"""
225223
Merge chunks into rows
@@ -231,7 +229,7 @@ async def merge_rows(
231229
"""
232230
try:
233231
if chunks is None:
234-
operation_metric.end_with_success()
232+
self._operation_metric.end_with_success()
235233
return
236234
it = chunks.__aiter__()
237235
# For each row
@@ -240,7 +238,7 @@ async def merge_rows(
240238
c = await it.__anext__()
241239
except CrossSync.StopIteration:
242240
# stream complete
243-
operation_metric.end_with_success()
241+
self._operation_metric.end_with_success()
244242
return
245243
row_key = c.row_key
246244

@@ -321,8 +319,8 @@ async def merge_rows(
321319
yield Row(row_key, cells)
322320
# most metric operations use setters, but this one updates
323321
# the value directly to avoid extra overhead
324-
if operation_metric.active_attempt is not None:
325-
operation_metric.active_attempt.application_blocking_time_ns += ( # type: ignore
322+
if self._operation_metric.active_attempt is not None:
323+
self._operation_metric.active_attempt.application_blocking_time_ns += ( # type: ignore
326324
time.monotonic_ns() - block_time
327325
) * 1000
328326
break
@@ -342,11 +340,11 @@ async def merge_rows(
342340
except CrossSync.StopIteration:
343341
raise InvalidChunk("premature end of stream")
344342
except Exception as generic_exception:
345-
if not retryable_predicate(generic_exception):
346-
operation_metric.end_attempt_with_status(generic_exception)
343+
if not self._predicate(generic_exception):
344+
self._operation_metric.end_attempt_with_status(generic_exception)
347345
raise generic_exception
348346
else:
349-
operation_metric.end_with_success()
347+
self._operation_metric.end_with_success()
350348

351349
@staticmethod
352350
def _revise_request_rowset(

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

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -87,7 +87,6 @@
8787
from google.cloud.bigtable.data.row_filters import RowFilterChain
8888
from google.cloud.bigtable.data._metrics import BigtableClientSideMetricsController
8989
from google.cloud.bigtable.data._metrics import OperationType
90-
from google.cloud.bigtable.data._metrics.handlers._stdout import _StdoutMetricsHandler
9190

9291
from google.cloud.bigtable.data._cross_sync import CrossSync
9392

@@ -944,7 +943,7 @@ def __init__(
944943

945944
self._metrics = BigtableClientSideMetricsController(
946945
client._metrics_interceptor,
947-
handlers=[_StdoutMetricsHandler()],
946+
handlers=[],
948947
project_id=self.client.project,
949948
instance_id=instance_id,
950949
table_id=table_id,

google/cloud/bigtable/data/_metrics/data_model.py

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -396,7 +396,6 @@ def _handle_error(message: str) -> None:
396396
"""
397397
full_message = f"Error in Bigtable Metrics: {message}"
398398
LOGGER.warning(full_message)
399-
raise RuntimeError(full_message)
400399

401400
def __enter__(self):
402401
"""

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

Lines changed: 12 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,7 @@
1717
# This file is automatically generated by CrossSync. Do not edit manually.
1818

1919
from __future__ import annotations
20-
from typing import Callable, Sequence, TYPE_CHECKING
20+
from typing import Sequence, TYPE_CHECKING
2121
import time
2222
from google.cloud.bigtable_v2.types import ReadRowsRequest as ReadRowsRequestPB
2323
from google.cloud.bigtable_v2.types import ReadRowsResponse as ReadRowsResponsePB
@@ -126,16 +126,16 @@ def _read_rows_attempt(self) -> CrossSync._Sync_Impl.Iterable[Row]:
126126
last_seen_row_key=self._last_yielded_row_key,
127127
)
128128
except _RowSetComplete:
129-
return self.merge_rows(None, self._operation_metric, self._predicate)
129+
return self.merge_rows(None)
130130
if self._remaining_count is not None:
131131
self.request.rows_limit = self._remaining_count
132132
if self._remaining_count == 0:
133-
return self.merge_rows(None, self._operation_metric, self._predicate)
133+
return self.merge_rows(None)
134134
gapic_stream = self.target.client._gapic_client.read_rows(
135135
self.request, timeout=next(self.attempt_timeout_gen), retry=None
136136
)
137137
chunked_stream = self.chunk_stream(gapic_stream)
138-
return self.merge_rows(chunked_stream, self._operation_metric, self._predicate)
138+
return self.merge_rows(chunked_stream)
139139

140140
def chunk_stream(
141141
self,
@@ -182,9 +182,7 @@ def chunk_stream(
182182

183183
@staticmethod
184184
def merge_rows(
185-
chunks: CrossSync._Sync_Impl.Iterable[ReadRowsResponsePB.CellChunk] | None,
186-
operation_metric: ActiveOperationMetric,
187-
retryable_predicate: Callable[[Exception], bool],
185+
self, chunks: CrossSync._Sync_Impl.Iterable[ReadRowsResponsePB.CellChunk] | None
188186
) -> CrossSync._Sync_Impl.Iterable[Row]:
189187
"""Merge chunks into rows
190188
@@ -194,14 +192,14 @@ def merge_rows(
194192
Row: the next row in the stream"""
195193
try:
196194
if chunks is None:
197-
operation_metric.end_with_success()
195+
self._operation_metric.end_with_success()
198196
return
199197
it = chunks.__iter__()
200198
while True:
201199
try:
202200
c = it.__next__()
203201
except CrossSync._Sync_Impl.StopIteration:
204-
operation_metric.end_with_success()
202+
self._operation_metric.end_with_success()
205203
return
206204
row_key = c.row_key
207205
if not row_key:
@@ -268,8 +266,8 @@ def merge_rows(
268266
if c.commit_row:
269267
block_time = time.monotonic_ns()
270268
yield Row(row_key, cells)
271-
if operation_metric.active_attempt is not None:
272-
operation_metric.active_attempt.application_blocking_time_ns += (
269+
if self._operation_metric.active_attempt is not None:
270+
self._operation_metric.active_attempt.application_blocking_time_ns += (
273271
time.monotonic_ns() - block_time
274272
) * 1000
275273
break
@@ -289,11 +287,11 @@ def merge_rows(
289287
except CrossSync._Sync_Impl.StopIteration:
290288
raise InvalidChunk("premature end of stream")
291289
except Exception as generic_exception:
292-
if not retryable_predicate(generic_exception):
293-
operation_metric.end_attempt_with_status(generic_exception)
290+
if not self._predicate(generic_exception):
291+
self._operation_metric.end_attempt_with_status(generic_exception)
294292
raise generic_exception
295293
else:
296-
operation_metric.end_with_success()
294+
self._operation_metric.end_with_success()
297295

298296
@staticmethod
299297
def _revise_request_rowset(row_set: RowSetPB, last_seen_row_key: bytes) -> RowSetPB:

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

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -74,7 +74,6 @@
7474
from google.cloud.bigtable.data.row_filters import RowFilterChain
7575
from google.cloud.bigtable.data._metrics import BigtableClientSideMetricsController
7676
from google.cloud.bigtable.data._metrics import OperationType
77-
from google.cloud.bigtable.data._metrics.handlers._stdout import _StdoutMetricsHandler
7877
from google.cloud.bigtable.data._cross_sync import CrossSync
7978
from typing import Iterable
8079
from grpc import insecure_channel
@@ -735,7 +734,7 @@ def __init__(
735734
)
736735
self._metrics = BigtableClientSideMetricsController(
737736
client._metrics_interceptor,
738-
handlers=[_StdoutMetricsHandler()],
737+
handlers=[],
739738
project_id=self.client.project,
740739
instance_id=instance_id,
741740
table_id=table_id,

0 commit comments

Comments
 (0)