This is an automated email from the ASF dual-hosted git repository.
damccorm pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new af99c7ce5d4 [Bigtable] Fix Bigtable segment truncation when open end
key is startKey + null byte (#39842) (#39843)
af99c7ce5d4 is described below
commit af99c7ce5d4b2516c54309e8a01cd9298c6f18f9
Author: Aditya Narayan <[email protected]>
AuthorDate: Mon Aug 24 23:12:10 2026 +0530
[Bigtable] Fix Bigtable segment truncation when open end key is startKey +
null byte (#39842) (#39843)
When truncateRequest splits an open end-key range where endKey == lastKey +
"\0",
Bigtable server rejects the next request with INVALID_ARGUMENT (start_key
must be
less than end_key) because start_key_open is normalized to lastKey + "\0".
Skip such exhausted open ranges during segment truncation.
Fixes #39842
---
.../sdk/io/gcp/bigtable/BigtableServiceImpl.java | 8 +++
.../io/gcp/bigtable/BigtableServiceImplTest.java | 81 ++++++++++++++++++++++
2 files changed, 89 insertions(+)
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableServiceImpl.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableServiceImpl.java
index f7aa50a7437..9eaf441a4bf 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableServiceImpl.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableServiceImpl.java
@@ -481,6 +481,14 @@ class BigtableServiceImpl implements BigtableService {
segment.addRowRanges(newRange.build());
} else {
// Row is split, remove all read rowKeys and split RowSet at last
buffered Row
+ if (rowRange.getEndKeyCase() == RowRange.EndKeyCase.END_KEY_OPEN
+ && !rowRange.getEndKeyOpen().isEmpty()) {
+ ByteString lastKeyWithNull =
lastKey.concat(ByteString.copyFrom(new byte[] {0}));
+ if (ByteStringComparator.INSTANCE.compare(lastKeyWithNull,
rowRange.getEndKeyOpen())
+ >= 0) {
+ continue;
+ }
+ }
segment.addRowRanges(newRange.setStartKeyOpen(lastKey).build());
}
}
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableServiceImplTest.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableServiceImplTest.java
index bec5a5470cb..dc756fcbe58 100644
---
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableServiceImplTest.java
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableServiceImplTest.java
@@ -763,6 +763,87 @@ public class BigtableServiceImplTest {
Mockito.verify(mockCallMetric, Mockito.times(3)).call("ok");
}
+ /**
+ * This test ensures that when a range has an open end key that is equal to
the start key plus a
+ * null byte (e.g. [k, k\0)), and the buffer hits the byte limit on key k,
truncateRequest
+ * correctly detects that the range is exhausted and does not create an
invalid range (k, k\0).
+ */
+ @Test
+ public void testReadRangeWithNullByteEndKeyAtByteLimit() throws IOException {
+ ByteString startKey = ByteString.copyFromUtf8("exact_key");
+ ByteString endKey = startKey.concat(ByteString.copyFrom(new byte[] {0}));
+ RowRange mockRowRange =
+
RowRange.newBuilder().setStartKeyClosed(startKey).setEndKeyOpen(endKey).build();
+
+ long segmentByteLimit = DEFAULT_ROW_SIZE / 2;
+
+ byte[] largeMemory = new byte[(int) DEFAULT_ROW_SIZE];
+ Row expectedRow =
+ Row.newBuilder()
+ .setKey(startKey)
+ .addFamilies(
+ Family.newBuilder()
+ .setName("Family")
+ .addColumns(
+ Column.newBuilder()
+
.setQualifier(ByteString.copyFromUtf8("LargeMemoryRow"))
+ .addCells(
+ Cell.newBuilder()
+ .setValue(ByteString.copyFrom(largeMemory))
+
.setTimestampMicros(System.currentTimeMillis())
+ .build())
+ .build())
+ .build())
+ .build();
+
+ List<List<Row>> expectedResults =
+ ImmutableList.of(ImmutableList.of(expectedRow), ImmutableList.of());
+
+ ServerStreamingCallable<Query, Row> mockCallable =
Mockito.mock(ServerStreamingCallable.class);
+
+ StreamController mockController = Mockito.mock(StreamController.class);
+ doAnswer(
+ new Answer<Void>() {
+ @Override
+ public Void answer(InvocationOnMock invocation) throws Throwable
{
+ cancelled.set(true);
+ return null;
+ }
+ })
+ .when(mockController)
+ .cancel();
+
+ doAnswer(new MultipleAnswer<Row>(expectedResults, mockController))
+ .when(mockCallable)
+ .call(any(Query.class), any(ResponseObserver.class),
any(ApiCallContext.class));
+
when(mockStub.createReadRowsCallable(any(RowAdapter.class))).thenReturn(mockCallable);
+ ServerStreamingCallable<Query, Row> callable =
+ mockStub.createReadRowsCallable(new
BigtableServiceImpl.BigtableRowProtoAdapter());
+
when(mockBigtableDataClient.readRowsCallable(any(RowAdapter.class))).thenReturn(callable);
+
+ BigtableService.Reader underTest =
+ new BigtableServiceImpl.BigtableSegmentReaderImpl(
+ mockBigtableDataClient,
+ bigtableDataSettings.getProjectId(),
+ bigtableDataSettings.getInstanceId(),
+ TABLE_ID,
+ RowSet.newBuilder().addRowRanges(mockRowRange).build(),
+ RowFilter.getDefaultInstance(),
+ SEGMENT_SIZE,
+ segmentByteLimit,
+ mockCallMetric);
+
+ List<Row> actualResults = new ArrayList<>();
+ Assert.assertTrue(underTest.start());
+ do {
+ actualResults.add(underTest.getCurrentRow());
+ } while (underTest.advance());
+
+ Assert.assertEquals(ImmutableList.of(expectedRow), actualResults);
+ Mockito.verify(mockCallable, Mockito.times(1))
+ .call(any(Query.class), any(ResponseObserver.class),
any(ApiCallContext.class));
+ }
+
/**
* This test ensures the Exception handling inside of the scanHandler. This
test will check if a
* StatusRuntimeException was thrown.