From a2a4b1a2b7596351e14e3c8f66da9d1a64407945 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Thu, 18 Sep 2025 10:55:41 +0200 Subject: [PATCH 1/4] fix output timestamp to be based on input window, not walltime. --- .../sdk/io/gcp/firestore/FirestoreV1WriteFn.java | 14 ++------------ 1 file changed, 2 insertions(+), 12 deletions(-) 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..abed10ad80fb 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 @@ -87,7 +87,6 @@ static final class BatchWriteFnWithSummary extends BaseBatchWriteFn context, - Instant timestamp, List> writeFailures, Runnable logMessage) { throw new FailedWritesException( @@ -97,11 +96,10 @@ void handleWriteFailures( @Override void handleWriteSummary( ContextAdapter context, - Instant timestamp, KV tuple, Runnable logMessage) { logMessage.run(); - context.output(tuple.getKey(), timestamp, tuple.getValue()); + context.output(tuple.getKey(), tuple.getValue().maxTimestamp(), tuple.getValue()); } } @@ -125,19 +123,17 @@ 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()); } } @Override void handleWriteSummary( ContextAdapter context, - Instant timestamp, KV tuple, Runnable logMessage) { logMessage.run(); @@ -274,7 +270,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 +461,6 @@ private DoFlushStatus doFlush( if (okCount == writesCount) { handleWriteSummary( context, - end, KV.of(new WriteSuccessSummary(okCount, okBytes), coerceNonNull(okWindow)), () -> LOG.debug( @@ -481,7 +475,6 @@ private DoFlushStatus doFlush( int finalOkCount = okCount; handleWriteFailures( context, - end, ImmutableList.copyOf(nonRetryableWrites), () -> LOG.warn( @@ -506,7 +499,6 @@ private DoFlushStatus doFlush( if (okCount > 0) { handleWriteSummary( context, - end, KV.of(new WriteSuccessSummary(okCount, okBytes), coerceNonNull(okWindow)), logMessage); } else { @@ -542,13 +534,11 @@ private enum DoFlushStatus { abstract void handleWriteFailures( ContextAdapter context, - Instant timestamp, List> writeFailures, Runnable logMessage); abstract void handleWriteSummary( ContextAdapter context, - Instant timestamp, KV tuple, Runnable logMessage); From 653d0853a867373c97557fecb0a1228834299d67 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Fri, 19 Sep 2025 09:45:15 +0200 Subject: [PATCH 2/4] fix output timestamp to be based on input window, not walltime. --- .../sdk/io/gcp/firestore/FirestoreV1WriteFn.java | 16 ++++++++++++++-- 1 file changed, 14 insertions(+), 2 deletions(-) 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 abed10ad80fb..57c19db8a5eb 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,6 +88,7 @@ static final class BatchWriteFnWithSummary extends BaseBatchWriteFn context, + Instant timestamp, List> writeFailures, Runnable logMessage) { throw new FailedWritesException( @@ -96,10 +98,11 @@ void handleWriteFailures( @Override void handleWriteSummary( ContextAdapter context, + Instant timestamp, KV tuple, Runnable logMessage) { logMessage.run(); - context.output(tuple.getKey(), tuple.getValue().maxTimestamp(), tuple.getValue()); + context.output(tuple.getKey(), timestamp, tuple.getValue()); } } @@ -123,17 +126,19 @@ static final class BatchWriteFnWithDeadLetterQueue extends BaseBatchWriteFn context, + Instant timestamp, List> writeFailures, Runnable logMessage) { logMessage.run(); for (KV kv : writeFailures) { - context.output(kv.getKey(), kv.getValue().maxTimestamp(), kv.getValue()); + context.output(kv.getKey(), timestamp, kv.getValue()); } } @Override void handleWriteSummary( ContextAdapter context, + Instant timestamp, KV tuple, Runnable logMessage) { logMessage.run(); @@ -270,6 +275,7 @@ public void processElement(ProcessContext context, BoundedWindow window) throws getWriteType(write), getName(write)); handleWriteFailures( contextAdapter, + clock.instant(), ImmutableList.of( KV.of( new WriteFailure( @@ -461,6 +467,7 @@ private DoFlushStatus doFlush( if (okCount == writesCount) { handleWriteSummary( context, + Preconditions.checkArgumentNotNull(okWindow).maxTimestamp(), KV.of(new WriteSuccessSummary(okCount, okBytes), coerceNonNull(okWindow)), () -> LOG.debug( @@ -475,6 +482,8 @@ private DoFlushStatus doFlush( int finalOkCount = okCount; handleWriteFailures( context, + Preconditions.checkArgumentNotNull(okWindow).maxTimestamp() + , ImmutableList.copyOf(nonRetryableWrites), () -> LOG.warn( @@ -499,6 +508,7 @@ private DoFlushStatus doFlush( if (okCount > 0) { handleWriteSummary( context, + Preconditions.checkArgumentNotNull(okWindow).maxTimestamp(), KV.of(new WriteSuccessSummary(okCount, okBytes), coerceNonNull(okWindow)), logMessage); } else { @@ -534,11 +544,13 @@ private enum DoFlushStatus { abstract void handleWriteFailures( ContextAdapter context, + Instant timestamp, List> writeFailures, Runnable logMessage); abstract void handleWriteSummary( ContextAdapter context, + Instant timestamp, KV tuple, Runnable logMessage); From 2e2da95e44a696e1ab20f6fa23babd98fcdc3204 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Fri, 19 Sep 2025 09:46:18 +0200 Subject: [PATCH 3/4] spotless --- .../apache/beam/sdk/io/gcp/firestore/FirestoreV1WriteFn.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) 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 57c19db8a5eb..8983216fc78e 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 @@ -482,8 +482,7 @@ private DoFlushStatus doFlush( int finalOkCount = okCount; handleWriteFailures( context, - Preconditions.checkArgumentNotNull(okWindow).maxTimestamp() - , + Preconditions.checkArgumentNotNull(okWindow).maxTimestamp(), ImmutableList.copyOf(nonRetryableWrites), () -> LOG.warn( From 5caf34d578a1171bfd743b48744df3c232410502 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Fri, 19 Sep 2025 13:16:02 +0200 Subject: [PATCH 4/4] fix preconditions --- .../beam/sdk/io/gcp/firestore/FirestoreV1WriteFn.java | 7 +------ 1 file changed, 1 insertion(+), 6 deletions(-) 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 8983216fc78e..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 @@ -88,7 +88,6 @@ static final class BatchWriteFnWithSummary extends BaseBatchWriteFn context, - Instant timestamp, List> writeFailures, Runnable logMessage) { throw new FailedWritesException( @@ -126,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()); } } @@ -275,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( @@ -482,7 +479,6 @@ private DoFlushStatus doFlush( int finalOkCount = okCount; handleWriteFailures( context, - Preconditions.checkArgumentNotNull(okWindow).maxTimestamp(), ImmutableList.copyOf(nonRetryableWrites), () -> LOG.warn( @@ -543,7 +539,6 @@ private enum DoFlushStatus { abstract void handleWriteFailures( ContextAdapter context, - Instant timestamp, List> writeFailures, Runnable logMessage);