diff --git a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/firestore/FirestoreV1WriteFn.java b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/firestore/FirestoreV1WriteFn.java index 70c2b91ffbfd..6bbb00e76f2d 100644 --- a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/firestore/FirestoreV1WriteFn.java +++ b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/firestore/FirestoreV1WriteFn.java @@ -51,6 +51,7 @@ import org.apache.beam.sdk.transforms.display.DisplayData; import org.apache.beam.sdk.transforms.windowing.BoundedWindow; import org.apache.beam.sdk.util.BackOffUtils; +import org.apache.beam.sdk.util.Preconditions; import org.apache.beam.sdk.values.KV; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; @@ -87,7 +88,6 @@ static final class BatchWriteFnWithSummary extends BaseBatchWriteFn context, - Instant timestamp, List> writeFailures, Runnable logMessage) { throw new FailedWritesException( @@ -125,12 +125,11 @@ static final class BatchWriteFnWithDeadLetterQueue extends BaseBatchWriteFn context, - Instant timestamp, List> writeFailures, Runnable logMessage) { logMessage.run(); for (KV kv : writeFailures) { - context.output(kv.getKey(), timestamp, kv.getValue()); + context.output(kv.getKey(), kv.getValue().maxTimestamp(), kv.getValue()); } } @@ -274,7 +273,6 @@ public void processElement(ProcessContext context, BoundedWindow window) throws getWriteType(write), getName(write)); handleWriteFailures( contextAdapter, - clock.instant(), ImmutableList.of( KV.of( new WriteFailure( @@ -466,7 +464,7 @@ private DoFlushStatus doFlush( if (okCount == writesCount) { handleWriteSummary( context, - end, + Preconditions.checkArgumentNotNull(okWindow).maxTimestamp(), KV.of(new WriteSuccessSummary(okCount, okBytes), coerceNonNull(okWindow)), () -> LOG.debug( @@ -481,7 +479,6 @@ private DoFlushStatus doFlush( int finalOkCount = okCount; handleWriteFailures( context, - end, ImmutableList.copyOf(nonRetryableWrites), () -> LOG.warn( @@ -506,7 +503,7 @@ private DoFlushStatus doFlush( if (okCount > 0) { handleWriteSummary( context, - end, + Preconditions.checkArgumentNotNull(okWindow).maxTimestamp(), KV.of(new WriteSuccessSummary(okCount, okBytes), coerceNonNull(okWindow)), logMessage); } else { @@ -542,7 +539,6 @@ private enum DoFlushStatus { abstract void handleWriteFailures( ContextAdapter context, - Instant timestamp, List> writeFailures, Runnable logMessage);