Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -107,14 +107,26 @@ async def add_aggregate_row(
@CrossSync.convert
async def delete_rows(self):
if self.rows:
request = {
"table_name": self.target.table_name,
"entries": [
{"row_key": row, "mutations": [{"delete_from_row": {}}]}
for row in self.rows
],
}
await self.target.client._gapic_client.mutate_rows(request)
# Chunk deletions to 5,000 rows. While Bigtable officially supports up to
# 100,000 mutations per MutateRows RPC, sending massive batches may hit
# the default gRPC 4MB client payload size limit due to metadata
# serialization overhead. Keeping chunks at 5,000 ensures we stay safely
# under 4MB and minimizes transient network timeouts on live connections.
chunk_size = 5000
for i in range(0, len(self.rows), chunk_size):
chunk = self.rows[i : i + chunk_size]
Comment thread
parthea marked this conversation as resolved.
Outdated
request = {
**self.target._request_path,
"entries": [
{"row_key": row, "mutations": [{"delete_from_row": {}}]}
for row in chunk
],
}
# Await and consume the gRPC stream to guarantee execution
stream = await self.target.client._gapic_client.mutate_rows(request)
async for response in stream:
pass
Comment thread
parthea marked this conversation as resolved.
Outdated


@CrossSync.convert
async def retrieve_cell_value(self, target, row_key):
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -94,14 +94,19 @@ def add_aggregate_row(

def delete_rows(self):
if self.rows:
request = {
"table_name": self.target.table_name,
"entries": [
{"row_key": row, "mutations": [{"delete_from_row": {}}]}
for row in self.rows
],
}
self.target.client._gapic_client.mutate_rows(request)
chunk_size = 5000
for i in range(0, len(self.rows), chunk_size):
chunk = self.rows[i : i + chunk_size]
request = {
**self.target._request_path,
"entries": [
{"row_key": row, "mutations": [{"delete_from_row": {}}]}
for row in chunk
],
}
stream = self.target.client._gapic_client.mutate_rows(request)
for response in stream:
pass

def retrieve_cell_value(self, target, row_key):
"""Helper to read an individual row"""
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2998,15 +2998,11 @@ def test_execute_query_with_params(self, client, execute_query_mock, prepare_moc
def test_execute_query_with_view_parameters(
self, client, execute_query_mock, prepare_mock
):
values = [
*chunked_responses(2, str_val("test2"), int_val(9), token=b"r2"),
]
values = [*chunked_responses(2, str_val("test2"), int_val(9), token=b"r2")]
execute_query_mock.return_value = self._make_gapic_stream(values)
query_str = f"SELECT a, b FROM {self.TABLE_NAME} WHERE user_id = VIEW_PARAMETERS('user_id')"
result = client.execute_query(
query_str,
self.INSTANCE_NAME,
view_parameters={"user_id": "alice"},
query_str, self.INSTANCE_NAME, view_parameters={"user_id": "alice"}
)
results = [r for r in result]
assert len(results) == 1
Expand All @@ -3015,7 +3011,6 @@ def test_execute_query_with_view_parameters(
assert execute_query_mock.call_count == 1
assert prepare_mock.call_count == 1
assert prepare_mock.call_args[1]["request"]["query"] == query_str

request = execute_query_mock.call_args[0][0]
assert "user_id" in request.view_parameters
assert request.view_parameters["user_id"].string_value == "alice"
Expand Down
Loading