diff --git a/CHANGES.txt b/CHANGES.txt index d9caa34bb..c56709a25 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,5 +1,6 @@ 0.5.0 ----- + * Prefer exact-range CDC checkpoint on load to avoid restart rewind (CASSSIDECAR-486) * Wire CDC configs in configs table to SidecarCdcOptions/SidecarStatePersister (CASSSIDECAR-483) * Implement durable operational job tracker (CASSSIDECAR-374) * Remove filesystem path from Http response (CASSSIDECAR-477) diff --git a/server/src/main/java/org/apache/cassandra/sidecar/db/CdcDatabaseAccessor.java b/server/src/main/java/org/apache/cassandra/sidecar/db/CdcDatabaseAccessor.java index 31c6b046e..7b187f9c5 100644 --- a/server/src/main/java/org/apache/cassandra/sidecar/db/CdcDatabaseAccessor.java +++ b/server/src/main/java/org/apache/cassandra/sidecar/db/CdcDatabaseAccessor.java @@ -19,6 +19,7 @@ package org.apache.cassandra.sidecar.db; +import java.math.BigInteger; import java.nio.ByteBuffer; import java.util.Arrays; import java.util.List; @@ -144,10 +145,23 @@ public Stream loadStateForRange(String jobId, TokenRange range) Arrays.stream(splits).mapToObj(Integer::toString).collect(Collectors.joining(","))); Stream futures = Arrays.stream(splits) .mapToObj(split -> selectCdcRange(jobId, split)); - Stream rows = await(futures); - return rows.filter(row -> !row.isNull(0) && !row.isNull(1) && !row.isNull(2)) - .filter(row -> TokenSplitUtil.overlaps(range, row.getVarint(0), row.getVarint(1))) - .map(row -> ByteBufUtils.getArray(row.getBytes(2))); + List overlappingRows = await(futures) + .filter(row -> !row.isNull(0) && !row.isNull(1) && !row.isNull(2)) + .filter(row -> TokenSplitUtil.overlaps(range, row.getVarint(0), row.getVarint(1))) + .collect(Collectors.toList()); + + List exactMatches = overlappingRows.stream() + .filter(row -> exactRangeMatch(range, row.getVarint(0), row.getVarint(1))) + .collect(Collectors.toList()); + + List selectedRows = exactMatches.isEmpty() ? overlappingRows : exactMatches; + + return selectedRows.stream().map(row -> ByteBufUtils.getArray(row.getBytes(2))); + } + + static boolean exactRangeMatch(TokenRange range, BigInteger start, BigInteger end) + { + return range.lowerEndpoint().equals(start) && range.upperEndpoint().equals(end); } @NotNull diff --git a/server/src/test/java/org/apache/cassandra/sidecar/db/CdcDatabaseAccessorTests.java b/server/src/test/java/org/apache/cassandra/sidecar/db/CdcDatabaseAccessorTests.java index 2db66cb24..19d85c582 100644 --- a/server/src/test/java/org/apache/cassandra/sidecar/db/CdcDatabaseAccessorTests.java +++ b/server/src/test/java/org/apache/cassandra/sidecar/db/CdcDatabaseAccessorTests.java @@ -226,6 +226,51 @@ void testExpand(int numNodes) } } + /** + * Regression for restart rewind: when an exact (start, end) row exists, stale overlapping rows + * from older range boundaries must not participate in load (they would otherwise win via min-merge). + */ + @Test + void testExactMatchIgnoresStaleOverlappingRanges() + { + Partitioner partitioner = Partitioner.Murmur3Partitioner; + MockCdcStateV2 datastore = new MockCdcStateV2(); + String jobId = UUID.randomUUID().toString(); + SidecarSchema mockSidecarSchema = mock(SidecarSchema.class); + CdcStatesSchema mockCdcStatesSchema = mock(CdcStatesSchema.class, RETURNS_DEEP_STUBS); + when(mockSidecarSchema.tableSchema(CdcStatesSchema.class)).thenReturn(mockCdcStatesSchema); + // 8 = total storage buckets on the ring (splits 0..7), not the number of rows returned on load. + TokenSplitUtil tokenSplitUtil = new TokenSplitUtil(8); + Provider tokenSplitUtilProvider = () -> tokenSplitUtil; + + CdcDatabaseAccessor db = new CdcDatabaseAccessor(mockSidecarSchema, + getMockCQLSessionProvider(datastore, mockCdcStatesSchema), + tokenSplitUtilProvider, + getMockInstanceMetaDataFetcher()); + + List tokens = TokenSplitUtil.splitTokens(8, partitioner); + // ownedRange spans storage buckets 3 and 4: (t[3], t[5]] crosses two of the eight ring slices. + BigInteger lower = tokens.get(3); + BigInteger upper = tokens.get(5); + // staleAdjacentRange overlaps ownedRange but has a different (start, end) clustering key. + BigInteger staleLower = tokens.get(4); + TokenRange ownedRange = TokenRange.openClosed(lower, upper); + TokenRange staleAdjacentRange = TokenRange.openClosed(staleLower, upper); + + ByteBuffer currentState = randomBytes(100); + ByteBuffer staleState = randomBytes(200); + + await(db.storeStateAsync(jobId, staleAdjacentRange, staleState, System.currentTimeMillis()).stream()); + await(db.storeStateAsync(jobId, ownedRange, currentState, System.currentTimeMillis()).stream()); + + List loaded = db.loadStateForRange(jobId, ownedRange).collect(Collectors.toList()); + // persist duplicates the exact-match row into every overlapping storage bucket (3 and 4 here). + assertThat(loaded).hasSize(2); + // both blobs are the current checkpoint; the stale adjacent row is excluded by exact-match-first load. + loaded.forEach(bytes -> assertByteBufferEquals(currentState, bytes)); + assertThat(loaded).noneMatch(bytes -> Arrays.equals(bytes, toByteArray(staleState))); + } + @Test void testOverlaps() {