|
36 | 36 | import com.google.cloud.spanner.TransactionContext; |
37 | 37 | import com.google.cloud.spanner.TransactionRunner; |
38 | 38 | import com.google.cloud.spanner.Value; |
| 39 | +import java.time.Duration; |
| 40 | +import java.time.Instant; |
39 | 41 | import java.util.Collections; |
40 | 42 | import java.util.Map; |
41 | 43 | import org.apache.beam.sdk.io.gcp.spanner.changestreams.model.PartitionMetadata; |
@@ -238,14 +240,39 @@ public void testInTransactionContextUpdateToFinished() { |
238 | 240 | @Test |
239 | 241 | public void testInTransactionContextUpdateWatermark() { |
240 | 242 | ArgumentCaptor<Mutation> mutation = ArgumentCaptor.forClass(Mutation.class); |
241 | | - doNothing().when(transaction).buffer(mutation.capture()); |
242 | | - assertNull(inTransactionContext.updateWatermark(PARTITION_TOKEN, WATERMARK)); |
| 243 | + when(transaction.readRow(any(), any(), any())) |
| 244 | + .thenReturn( |
| 245 | + Struct.newBuilder() |
| 246 | + .set(PartitionMetadataAdminDao.COLUMN_WATERMARK) |
| 247 | + .to(WATERMARK) |
| 248 | + .build()); |
| 249 | + Instant largerWatermark = WATERMARK.toSqlTimestamp().toInstant().plus(Duration.ofSeconds(1)); |
| 250 | + assertNull( |
| 251 | + inTransactionContext.updateWatermark( |
| 252 | + PARTITION_TOKEN, |
| 253 | + Timestamp.ofTimeSecondsAndNanos( |
| 254 | + largerWatermark.getEpochSecond(), largerWatermark.getNano()))); |
| 255 | + verify(transaction).buffer(mutation.capture()); |
243 | 256 | Map<String, Value> mutationValueMap = mutation.getValue().asMap(); |
244 | 257 | assertEquals( |
245 | 258 | PARTITION_TOKEN, |
246 | 259 | mutationValueMap.get(PartitionMetadataAdminDao.COLUMN_PARTITION_TOKEN).getString()); |
247 | 260 | assertEquals( |
248 | | - WATERMARK, mutationValueMap.get(PartitionMetadataAdminDao.COLUMN_WATERMARK).getTimestamp()); |
| 261 | + Timestamp.ofTimeSecondsAndNanos( |
| 262 | + largerWatermark.getEpochSecond(), largerWatermark.getNano()), |
| 263 | + mutationValueMap.get(PartitionMetadataAdminDao.COLUMN_WATERMARK).getTimestamp()); |
| 264 | + } |
| 265 | + |
| 266 | + @Test |
| 267 | + public void testInTransactionContextDoNotUpdateWatermark() { |
| 268 | + when(transaction.readRow(any(), any(), any())) |
| 269 | + .thenReturn( |
| 270 | + Struct.newBuilder() |
| 271 | + .set(PartitionMetadataAdminDao.COLUMN_WATERMARK) |
| 272 | + .to(WATERMARK) |
| 273 | + .build()); |
| 274 | + assertNull(inTransactionContext.updateWatermark(PARTITION_TOKEN, WATERMARK)); |
| 275 | + verify(transaction, times(0)).buffer(any(Mutation.class)); |
249 | 276 | } |
250 | 277 |
|
251 | 278 | @Test |
|
0 commit comments