Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
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
20 changes: 20 additions & 0 deletions packages/google-cloud-bigtable/google/cloud/bigtable/batcher.py
Original file line number Diff line number Diff line change
Expand Up @@ -350,6 +350,26 @@ def _batch_completed_callback(self, future):
processed_rows = self.futures_mapping[future]
self.flow_control.release(processed_rows)
del self.futures_mapping[future]
# Surface any exception raised inside the async flush. Without this, an
# exception raised by ``_flush_rows`` (e.g. a non-retryable RPC error, a
# retry deadline, or a response-count mismatch) would be stored on the
# future and silently discarded, so the failed mutations would never be
# reported to the user -- effectively silent data loss. Per-row errors
# from a successful RPC are already recorded in ``self.exceptions`` by
# ``_flush_rows``; here the whole batch failed with a single exception,
# so record it once per row in the batch to keep the reported error
# count aligned with the number of affected mutations.
#
# A cancelled future is "done", so this callback still runs for it, but
# ``future.exception()`` would raise ``CancelledError``. Nothing here
# cancels futures today, but guard against it so the callback stays
# correct if cancellation is ever introduced.
if future.cancelled():
return
exc = future.exception()
if exc is not None:
Comment thread
sushanb marked this conversation as resolved.
for _ in range(processed_rows.rows_count):
self.exceptions.put(exc)

def _row_fits_in_batch(self, row, batch_info):
"""Checks if a row can fit in the current batch.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -213,6 +213,60 @@ def test_mutations_batcher_response_with_error_codes():
assert exc.value.exc[1].message == mocked_response[1].message


def test_mutations_batcher_asynchronous_flush_exception_is_surfaced():
"""An exception raised by the underlying ``mutate_rows`` call (e.g. a
non-retryable RPC error or a response-count mismatch) is raised inside the
async flush task. It must be captured and re-raised at ``close()`` rather
than being silently swallowed by the executor -- otherwise the failed
mutations are never reported to the user (silent data loss)."""
from google.api_core.exceptions import PermissionDenied

with mock.patch("tests.unit.v2_client.test_batcher._Table") as mocked_table:
table = mocked_table.return_value
# flush_count=2 forces the batch to flush asynchronously (through the
# executor) as soon as the second row is added
mutation_batcher = MutationsBatcher(table=table, flush_count=2)

row1 = DirectRow(row_key=b"row_key")
row1.set_cell("cf1", b"c1", b"1")
row2 = DirectRow(row_key=b"row_key")
row2.set_cell("cf1", b"c1", b"2")
table.mutate_rows.side_effect = PermissionDenied("denied")

mutation_batcher.mutate_rows([row1, row2])
with pytest.raises(MutationsBatchError) as exc:
mutation_batcher.close()
assert exc.value.message == "Errors in batch mutations."
# the whole batch (both rows) failed, so both are reported -- the error
# count stays aligned with the number of affected mutations
assert len(exc.value.exc) == 2
assert all(isinstance(e, PermissionDenied) for e in exc.value.exc)


def test_batch_completed_callback_ignores_cancelled_future():
"""A cancelled future is still "done", so the completion callback runs for
it, but ``future.exception()`` would raise ``CancelledError``. The callback
must short-circuit on a cancelled future instead of letting that propagate."""
from google.cloud.bigtable.batcher import _BatchInfo

table = _Table(TABLE_NAME)
with MutationsBatcher(table=table) as mutation_batcher:
batch_info = _BatchInfo(rows_count=2, mutations_count=2, mutations_size=0)

cancelled_future = mock.Mock()
cancelled_future.cancelled.return_value = True
cancelled_future.exception.side_effect = AssertionError(
"exception() must not be called on a cancelled future"
)
mutation_batcher.futures_mapping[cancelled_future] = batch_info

# Should not raise, should not record any exceptions
mutation_batcher._batch_completed_callback(cancelled_future)

assert cancelled_future not in mutation_batcher.futures_mapping
assert mutation_batcher.exceptions.qsize() == 0


def test_flow_control_event_is_set_when_not_blocked():
flow_control = _FlowControl()

Expand Down
Loading