This is an automated email from the ASF dual-hosted git repository.
ahmedabu98 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 544e40f5a89 [Iceberg CDC sink] Split late data and create commit
windows (#40006)
544e40f5a89 is described below
commit 544e40f5a89aef7cdb9891801ac687d022ff5478
Author: Ahmed Abualsaud <[email protected]>
AuthorDate: Mon Sep 14 20:15:16 2026 +0200
[Iceberg CDC sink] Split late data and create commit windows (#40006)
* split late data and create commit windows
* address comments
---
sdks/java/io/iceberg/build.gradle | 1 +
.../sdk/io/iceberg/cdc/sink/CommitWindows.java | 226 +++++++++++
.../sdk/io/iceberg/cdc/sink/DestinationShard.java | 61 +++
.../sdk/io/iceberg/cdc/sink/SplitLateData.java | 102 +++++
.../sdk/io/iceberg/cdc/sink/CommitWindowsTest.java | 411 +++++++++++++++++++++
5 files changed, 801 insertions(+)
diff --git a/sdks/java/io/iceberg/build.gradle
b/sdks/java/io/iceberg/build.gradle
index e2e8a12d01e..b3228004ed3 100644
--- a/sdks/java/io/iceberg/build.gradle
+++ b/sdks/java/io/iceberg/build.gradle
@@ -46,6 +46,7 @@ dependencies {
implementation library.java.vendored_guava_32_1_2_jre
implementation project(path: ":sdks:java:core", configuration: "shadow")
implementation project(path: ":model:pipeline", configuration: "shadow")
+ implementation project(path: ":sdks:java:extensions:sorter")
implementation library.java.avro
implementation library.java.slf4j_api
implementation library.java.joda_time
diff --git
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CommitWindows.java
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CommitWindows.java
new file mode 100644
index 00000000000..3696b155b00
--- /dev/null
+++
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CommitWindows.java
@@ -0,0 +1,226 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.io.iceberg.cdc.sink;
+
+import static org.apache.beam.sdk.util.Preconditions.checkStateNotNull;
+import static
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkArgument;
+
+import java.util.Map;
+import org.apache.beam.sdk.Pipeline;
+import org.apache.beam.sdk.coders.Coder;
+import org.apache.beam.sdk.coders.KvCoder;
+import org.apache.beam.sdk.coders.RowCoder;
+import org.apache.beam.sdk.extensions.sorter.BufferedExternalSorter;
+import org.apache.beam.sdk.extensions.sorter.SortValues;
+import org.apache.beam.sdk.schemas.Schema;
+import org.apache.beam.sdk.transforms.GroupByKey;
+import org.apache.beam.sdk.transforms.PTransform;
+import org.apache.beam.sdk.transforms.ParDo;
+import org.apache.beam.sdk.transforms.windowing.AfterPane;
+import org.apache.beam.sdk.transforms.windowing.AfterWatermark;
+import org.apache.beam.sdk.transforms.windowing.DefaultTrigger;
+import org.apache.beam.sdk.transforms.windowing.FixedWindows;
+import org.apache.beam.sdk.transforms.windowing.GlobalWindows;
+import org.apache.beam.sdk.transforms.windowing.Window;
+import org.apache.beam.sdk.values.KV;
+import org.apache.beam.sdk.values.PCollection;
+import org.apache.beam.sdk.values.PCollection.IsBounded;
+import org.apache.beam.sdk.values.PCollectionTuple;
+import org.apache.beam.sdk.values.PInput;
+import org.apache.beam.sdk.values.POutput;
+import org.apache.beam.sdk.values.PValue;
+import org.apache.beam.sdk.values.Row;
+import org.apache.beam.sdk.values.TupleTag;
+import org.apache.beam.sdk.values.TupleTagList;
+import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
+import org.checkerframework.checker.nullness.qual.Nullable;
+import org.joda.time.Duration;
+
+/**
+ * Windows sharded, sort-keyed records into event-time commit windows, grouped
by the {@code
+ * KV<destination, shard>} key. Each {@code (destination, shard, window)}
becomes one commit unit
+ * group and the event-time watermark is the commit barrier. Each group is
sorted by the byte sort
+ * key applied from {@link AssignCdcKeys}: one primary key's records
contiguous, in {@code (seq,
+ * kind)} order within the key.
+ *
+ * <p>Late panes are routed to the DLQ (via {@link SplitLateData}) before any
ordering happens. This
+ * is necessary because the downstream commit step skips panes if their window
token is present in
+ * an already committed snapshot. If late panes are let through, their
contents will never get to
+ * the table.
+ *
+ * <p>Windowing by input mode:
+ *
+ * <ul>
+ * <li><b>Bounded</b>: a single {@link GlobalWindows} window. One commit per
destination, no late
+ * data possible.
+ * <li><b>Unbounded</b>: event-time {@link FixedWindows} of {@code
triggeringFrequency}, firing on
+ * the watermark with late firings per element, discarding fired panes.
+ * </ul>
+ */
+final class CommitWindows
+ extends PTransform<
+ PCollection<KV<DestinationShard, KV<byte[], CdcRecord>>>,
CommitWindows.Result> {
+
+ private static final TupleTag<KV<DestinationShard, Iterable<KV<byte[],
CdcRecord>>>> ON_TIME_TAG =
+ new TupleTag<>("onTime");
+ private static final TupleTag<Row> DEAD_LETTER_TAG = new
TupleTag<>("deadLetter");
+
+ private final CdcWriteConfig config;
+ private final @Nullable Duration triggeringFrequency;
+ private final Duration allowedLateness;
+
+ CommitWindows(
+ CdcWriteConfig config, @Nullable Duration triggeringFrequency, Duration
allowedLateness) {
+ this.config = config;
+ this.triggeringFrequency = triggeringFrequency;
+ this.allowedLateness = allowedLateness;
+ }
+
+ @Override
+ public Result expand(PCollection<KV<DestinationShard, KV<byte[],
CdcRecord>>> input) {
+ Schema deadLetterSchema =
SplitLateData.deadLetterSchema(dataSchemaOf(input.getCoder()));
+
+ PCollection<KV<DestinationShard, KV<byte[], CdcRecord>>> windowed =
applyCommitWindow(input);
+
+ // Exactly one group per (destination, shard, window)
+ PCollection<KV<DestinationShard, Iterable<KV<byte[], CdcRecord>>>> grouped
=
+ windowed.apply("GroupByShardKey", GroupByKey.create());
+
+ // Late-data split before sorting anything
+ PCollectionTuple split =
+ grouped.apply(
+ "SplitLateData",
+ ParDo.of(new SplitLateData(deadLetterSchema, ON_TIME_TAG,
DEAD_LETTER_TAG))
+ .withOutputTags(ON_TIME_TAG,
TupleTagList.of(DEAD_LETTER_TAG)));
+ PCollection<KV<DestinationShard, Iterable<KV<byte[], CdcRecord>>>>
onTimeUnsorted =
+ split.get(ON_TIME_TAG).setCoder(grouped.getCoder());
+ PCollection<Row> deadLetter =
+ split.get(DEAD_LETTER_TAG).setCoder(RowCoder.of(deadLetterSchema));
+
+ // Sort each surviving group's records by the byte sort key. The secondary
key is byte[] +
+ // ByteArrayCoder, so SortValues compares the raw CdcSortKey bytes (no
coder framing):
+ // each primary key's records come out contiguous, in (seq, kind) order
within the key.
+ PCollection<KV<DestinationShard, Iterable<KV<byte[], CdcRecord>>>> sorted =
+ onTimeUnsorted.apply(
+ "SortBySeqKind",
+ SortValues.create(
+
BufferedExternalSorter.options().withMemoryMB(config.getSorterMemoryMB())));
+
+ return new Result(input.getPipeline(), sorted, deadLetter,
deadLetterSchema);
+ }
+
+ /**
+ * Applies the commit-window assignment: single {@link GlobalWindows} for
bounded, and event-time
+ * {@link FixedWindows} for unbounded.
+ */
+ private PCollection<KV<DestinationShard, KV<byte[], CdcRecord>>>
applyCommitWindow(
+ PCollection<KV<DestinationShard, KV<byte[], CdcRecord>>> input) {
+ if (input.isBounded() == IsBounded.BOUNDED) {
+ return input.apply(
+ "GlobalWindows",
+ Window.<KV<DestinationShard, KV<byte[], CdcRecord>>>into(new
GlobalWindows())
+ .triggering(DefaultTrigger.of())
+ .discardingFiredPanes());
+ }
+ return input.apply(
+ "EventTimeWindows",
+ Window.<KV<DestinationShard, KV<byte[], CdcRecord>>>into(
+ FixedWindows.of(
+ checkStateNotNull(
+ triggeringFrequency,
+ "triggeringFrequency is required for unbounded
input")))
+ .triggering(
+
AfterWatermark.pastEndOfWindow().withLateFirings(AfterPane.elementCountAtLeast(1)))
+ .withAllowedLateness(allowedLateness)
+ .discardingFiredPanes());
+ }
+
+ /** Extracts the CDC data schema carried by the input's nested {@link
CdcRecordCoder}. */
+ private static Schema dataSchemaOf(Coder<?> inputCoder) {
+ checkArgument(
+ inputCoder instanceof KvCoder,
+ "expected a KvCoder input element coder, got %s",
+ inputCoder);
+ Coder<?> valueCoder = ((KvCoder<?, ?>) inputCoder).getValueCoder();
+ checkArgument(
+ valueCoder instanceof KvCoder, "expected a KvCoder input value coder,
got %s", valueCoder);
+ Coder<?> recordCoder = ((KvCoder<?, ?>) valueCoder).getValueCoder();
+ checkArgument(
+ recordCoder instanceof CdcRecordCoder,
+ "expected a CdcRecordCoder input record coder, got %s",
+ recordCoder);
+ return ((CdcRecordCoder) recordCoder).getDataSchema();
+ }
+
+ /**
+ * The output of {@link CommitWindows}: the surviving sorted groups ready
for the delta writer,
+ * and the replayable dead-letter {@link Row}s from late panes.
+ */
+ public static final class Result implements POutput {
+
+ private final Pipeline pipeline;
+ private final PCollection<KV<DestinationShard, Iterable<KV<byte[],
CdcRecord>>>> sortedGroups;
+ private final PCollection<Row> deadLetterRows;
+ private final Schema deadLetterSchema;
+
+ private Result(
+ Pipeline pipeline,
+ PCollection<KV<DestinationShard, Iterable<KV<byte[], CdcRecord>>>>
sortedGroups,
+ PCollection<Row> deadLetterRows,
+ Schema deadLetterSchema) {
+ this.pipeline = pipeline;
+ this.sortedGroups = sortedGroups;
+ this.deadLetterRows = deadLetterRows;
+ this.deadLetterSchema = deadLetterSchema;
+ }
+
+ /** The surviving sorted groups: one per {@code (destination, shard,
window)}. */
+ public PCollection<KV<DestinationShard, Iterable<KV<byte[], CdcRecord>>>>
getSortedGroups() {
+ return sortedGroups;
+ }
+
+ /** Replayable dead-letter rows from late panes; {@link SplitLateData}
describes the shape. */
+ public PCollection<Row> getDeadLetterRows() {
+ return deadLetterRows;
+ }
+
+ /** The schema of {@link #getDeadLetterRows()}. */
+ public Schema getDeadLetterSchema() {
+ return deadLetterSchema;
+ }
+
+ @Override
+ public Pipeline getPipeline() {
+ return pipeline;
+ }
+
+ @Override
+ public Map<TupleTag<?>, PValue> expand() {
+ return ImmutableMap.<TupleTag<?>, PValue>builder()
+ .put(ON_TIME_TAG, sortedGroups)
+ .put(DEAD_LETTER_TAG, deadLetterRows)
+ .build();
+ }
+
+ @Override
+ public void finishSpecifyingOutput(
+ String transformName, PInput input, PTransform<?, ?> transform) {
+ // no-op
+ }
+ }
+}
diff --git
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/DestinationShard.java
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/DestinationShard.java
new file mode 100644
index 00000000000..199655a8a81
--- /dev/null
+++
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/DestinationShard.java
@@ -0,0 +1,61 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.io.iceberg.cdc.sink;
+
+import com.google.auto.value.AutoValue;
+import org.apache.beam.sdk.coders.Coder;
+import org.apache.beam.sdk.schemas.AutoValueSchema;
+import org.apache.beam.sdk.schemas.NoSuchSchemaException;
+import org.apache.beam.sdk.schemas.SchemaRegistry;
+import org.apache.beam.sdk.schemas.annotations.DefaultSchema;
+
+/** The shuffle key of one write group: a destination and one of its shards. */
+@DefaultSchema(AutoValueSchema.class)
+@AutoValue
+public abstract class DestinationShard {
+
+ public abstract String getDestination();
+
+ public abstract int getShard();
+
+ public static Builder builder() {
+ return new AutoValue_DestinationShard.Builder();
+ }
+
+ public static DestinationShard of(String destination, int shard) {
+ return builder().setDestination(destination).setShard(shard).build();
+ }
+
+ /** A deterministic coder, so the key can drive a {@code GroupByKey}. */
+ public static Coder<DestinationShard> coder() {
+ try {
+ return
SchemaRegistry.createDefault().getSchemaCoder(DestinationShard.class);
+ } catch (NoSuchSchemaException e) {
+ throw new RuntimeException("Could not build a coder for
DestinationShard.", e);
+ }
+ }
+
+ @AutoValue.Builder
+ public abstract static class Builder {
+ public abstract Builder setDestination(String destination);
+
+ public abstract Builder setShard(int shard);
+
+ public abstract DestinationShard build();
+ }
+}
diff --git
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/SplitLateData.java
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/SplitLateData.java
new file mode 100644
index 00000000000..52063a4202e
--- /dev/null
+++
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/SplitLateData.java
@@ -0,0 +1,102 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.io.iceberg.cdc.sink;
+
+import org.apache.beam.sdk.metrics.Counter;
+import org.apache.beam.sdk.metrics.Metrics;
+import org.apache.beam.sdk.schemas.Schema;
+import org.apache.beam.sdk.transforms.DoFn;
+import org.apache.beam.sdk.transforms.windowing.PaneInfo;
+import org.apache.beam.sdk.values.KV;
+import org.apache.beam.sdk.values.Row;
+import org.apache.beam.sdk.values.TupleTag;
+
+/**
+ * Splits late data in {@link CommitWindows}: every {@linkplain
PaneInfo.Timing#LATE late} pane is
+ * diverted to a DLQ side output; on-time and early panes pass through
unchanged.
+ *
+ * <p>Pane history is per {@code (destination, shard, window)} while the
committer's skip is per
+ * {@code (destination, window)}: a late pane on a shard with no on-time data
cannot prove its
+ * destination-window uncommitted, and a record let through into a committed
window would reach
+ * neither the table nor the dead-letter output.
+ *
+ * <p>Each dead letter nests the untouched data row under {@value #DL_RECORD}
beside {@value
+ * #DL_CHANGE_TYPE}, {@value #DL_SEQ}, and {@value #DL_DEST}; to replay,
unnest {@value #DL_RECORD}
+ * and map {@value #DL_CHANGE_TYPE}/{@value #DL_SEQ} as the sink's control
columns. Replaying is
+ * only safe while no newer change for those keys has committed; a stale
replay's equality delete
+ * removes the newer row.
+ */
+final class SplitLateData
+ extends DoFn<
+ KV<DestinationShard, Iterable<KV<byte[], CdcRecord>>>,
+ KV<DestinationShard, Iterable<KV<byte[], CdcRecord>>>> {
+
+ static final String DL_RECORD = "record";
+ static final String DL_CHANGE_TYPE = "change_type";
+ static final String DL_SEQ = "sequence_number";
+ static final String DL_DEST = "destination";
+
+ private final Counter deadLetterRecords =
+ Metrics.counter(SplitLateData.class, "deadLetterRecords");
+
+ private final Schema deadLetterSchema;
+ private final TupleTag<KV<DestinationShard, Iterable<KV<byte[],
CdcRecord>>>> onTimeTag;
+ private final TupleTag<Row> deadLetterTag;
+
+ SplitLateData(
+ Schema deadLetterSchema,
+ TupleTag<KV<DestinationShard, Iterable<KV<byte[], CdcRecord>>>>
onTimeTag,
+ TupleTag<Row> deadLetterTag) {
+ this.deadLetterSchema = deadLetterSchema;
+ this.onTimeTag = onTimeTag;
+ this.deadLetterTag = deadLetterTag;
+ }
+
+ static Schema deadLetterSchema(Schema cdcDataSchema) {
+ return Schema.builder()
+ .addRowField(DL_RECORD, cdcDataSchema)
+ .addStringField(DL_CHANGE_TYPE)
+ .addInt64Field(DL_SEQ)
+ .addStringField(DL_DEST)
+ .build();
+ }
+
+ @ProcessElement
+ public void process(
+ @Element KV<DestinationShard, Iterable<KV<byte[], CdcRecord>>> group,
+ PaneInfo pane,
+ MultiOutputReceiver out) {
+ if (pane.getTiming() == PaneInfo.Timing.LATE) {
+ String dest = group.getKey().getDestination();
+ for (KV<byte[], CdcRecord> kv : group.getValue()) {
+ CdcRecord record = kv.getValue();
+ Row deadLetter =
+ Row.withSchema(deadLetterSchema)
+ .addValue(record.getData())
+ .addValue(record.getKind().name())
+ .addValue(record.getSequenceNumber())
+ .addValue(dest)
+ .build();
+ out.get(deadLetterTag).output(deadLetter);
+ deadLetterRecords.inc();
+ }
+ } else {
+ out.get(onTimeTag).output(group);
+ }
+ }
+}
diff --git
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CommitWindowsTest.java
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CommitWindowsTest.java
new file mode 100644
index 00000000000..b4888c353ec
--- /dev/null
+++
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CommitWindowsTest.java
@@ -0,0 +1,411 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.io.iceberg.cdc.sink;
+
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.contains;
+import static org.hamcrest.Matchers.containsInAnyOrder;
+import static org.hamcrest.Matchers.equalTo;
+
+import java.util.ArrayList;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Set;
+import org.apache.beam.sdk.PipelineResult;
+import org.apache.beam.sdk.coders.ByteArrayCoder;
+import org.apache.beam.sdk.coders.KvCoder;
+import org.apache.beam.sdk.metrics.MetricNameFilter;
+import org.apache.beam.sdk.metrics.MetricResult;
+import org.apache.beam.sdk.metrics.MetricsFilter;
+import org.apache.beam.sdk.schemas.Schema;
+import org.apache.beam.sdk.testing.PAssert;
+import org.apache.beam.sdk.testing.TestPipeline;
+import org.apache.beam.sdk.testing.TestStream;
+import org.apache.beam.sdk.transforms.Create;
+import org.apache.beam.sdk.values.KV;
+import org.apache.beam.sdk.values.Row;
+import org.apache.beam.sdk.values.TimestampedValue;
+import org.apache.beam.sdk.values.ValueKind;
+import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList;
+import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableSet;
+import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables;
+import org.joda.time.Duration;
+import org.joda.time.Instant;
+import org.junit.Rule;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
+
+/**
+ * Tests for {@link CommitWindows} (with {@link SplitLateData}), stage 2 of
the CDC sink: event-time
+ * commit windowing, {@code GroupByKey} per {@code (destination, shard)}, the
late-pane dead-letter
+ * split, and the {@code SortValues} sort by the byte {@code (pk, seq, kind)}
key.
+ *
+ * <p>Pure-Beam tests, no Iceberg catalog: {@link TestStream} drives the
watermark for the streaming
+ * cases, {@link Create} provides bounded input for the batch cases, and
{@link PAssert} checks
+ * {@link CommitWindows.Result#getSortedGroups()} and {@link
+ * CommitWindows.Result#getDeadLetterRows()}.
+ */
+@RunWith(JUnit4.class)
+public class CommitWindowsTest {
+
+ @Rule public transient TestPipeline p = TestPipeline.create();
+
+ /** The CDC data schema for these tests: {@code id INT32}, {@code name
STRING}. */
+ private static final Schema DATA_SCHEMA =
+ Schema.builder().addInt32Field("id").addStringField("name").build();
+
+ /** {@link #DATA_SCHEMA} nested under {@code record}, plus the dead-letter
metadata columns. */
+ private static final Schema EXPECTED_DEAD_LETTER_SCHEMA =
+ Schema.builder()
+ .addRowField("record", DATA_SCHEMA)
+ .addStringField("change_type")
+ .addInt64Field("sequence_number")
+ .addStringField("destination")
+ .build();
+
+ /** The stage-1 ({@link AssignCdcKeys#KEYED}) element coder. */
+ private static final KvCoder<DestinationShard, KV<byte[], CdcRecord>>
INPUT_CODER =
+ KvCoder.of(
+ DestinationShard.coder(),
+ KvCoder.of(ByteArrayCoder.of(), CdcRecordCoder.of(DATA_SCHEMA)));
+
+ /** The streaming commit-window size. */
+ private static final Duration WINDOW = Duration.standardSeconds(60);
+
+ /** A test-specific lateness bound, large enough to keep every late test
pane in-window. */
+ private static final Duration TEST_ALLOWED_LATENESS =
Duration.standardDays(7);
+
+ //
---------------------------------------------------------------------------------------------
+ // Helpers
+ //
---------------------------------------------------------------------------------------------
+
+ private static CdcWriteConfig config() {
+ return
CdcWriteConfig.builder().setSinkId("test-sink").setSorterMemoryMB(16).build();
+ }
+
+ /** Boundedness is derived from the input {@link PCollection}, not
configured. */
+ private static CommitWindows batchWindows() {
+ return new CommitWindows(config(), /* triggeringFrequency= */ null,
Duration.ZERO);
+ }
+
+ private static CommitWindows streamingWindows() {
+ return new CommitWindows(config(), WINDOW, TEST_ALLOWED_LATENESS);
+ }
+
+ private static Row data(int id, String name) {
+ return Row.withSchema(DATA_SCHEMA).addValues(id, name).build();
+ }
+
+ /** One primary key shared by every element, so a group's records sort
purely by (seq, kind). */
+ private static final byte[] PK = {42};
+
+ /** Builds one stage-1 output element: {@code KV<DestinationShard,
KV<sortKey, CdcRecord>>}. */
+ private static KV<DestinationShard, KV<byte[], CdcRecord>> element(
+ String dest, int shard, int id, String name, long seq, ValueKind kind) {
+ return KV.of(
+ DestinationShard.of(dest, shard),
+ KV.of(CdcSortKey.encode(PK, seq, kind), CdcRecord.of(data(id, name),
kind, seq)));
+ }
+
+ /** A {@link TimestampedValue} wrapping {@link #element}, for use with
{@link TestStream}. */
+ private static TimestampedValue<KV<DestinationShard, KV<byte[], CdcRecord>>>
at(
+ Instant ts, String dest, int shard, int id, String name, long seq,
ValueKind kind) {
+ return TimestampedValue.of(element(dest, shard, id, name, seq, kind), ts);
+ }
+
+ /** Sequence numbers, in encounter order, of a group's {@link CdcRecord}s. */
+ private static List<Long> seqsOf(KV<DestinationShard, Iterable<KV<byte[],
CdcRecord>>> group) {
+ List<Long> seqs = new ArrayList<>();
+ for (KV<byte[], CdcRecord> kv : group.getValue()) {
+ seqs.add(kv.getValue().getSequenceNumber());
+ }
+ return seqs;
+ }
+
+ /** Change kinds, in encounter order, of a group's {@link CdcRecord}s. */
+ private static List<ValueKind> kindsOf(
+ KV<DestinationShard, Iterable<KV<byte[], CdcRecord>>> group) {
+ List<ValueKind> kinds = new ArrayList<>();
+ for (KV<byte[], CdcRecord> kv : group.getValue()) {
+ kinds.add(kv.getValue().getKind());
+ }
+ return kinds;
+ }
+
+ /** Data-row {@code name} values, in encounter order, of a group's {@link
CdcRecord}s. */
+ private static List<String> namesOf(KV<DestinationShard, Iterable<KV<byte[],
CdcRecord>>> group) {
+ List<String> names = new ArrayList<>();
+ for (KV<byte[], CdcRecord> kv : group.getValue()) {
+ names.add(kv.getValue().getData().getString("name"));
+ }
+ return names;
+ }
+
+ /** Sums the committed values of the named {@link SplitLateData} counter (0
if never fired). */
+ private static long counterTotal(PipelineResult result, String name) {
+ Iterable<MetricResult<Long>> counters =
+ result
+ .metrics()
+ .queryMetrics(
+ MetricsFilter.builder()
+ .addNameFilter(MetricNameFilter.named(SplitLateData.class,
name))
+ .build())
+ .getCounters();
+ long total = 0;
+ for (MetricResult<Long> counter : counters) {
+ total += counter.getCommitted();
+ }
+ return total;
+ }
+
+ //
---------------------------------------------------------------------------------------------
+ // 1. Batch (bounded): one global-window group per (dest, shard), sorted
+ //
---------------------------------------------------------------------------------------------
+
+ @Test
+ public void batchGlobalWindowGroupsAndSortsBySeqAndKind() {
+ // Records deliberately out of order, including an equal-seq (9, 9) pair
where the
+ // UPDATE_AFTER is added BEFORE the UPDATE_BEFORE: the byte sort key must
order by seq,
+ // then before-image (UPDATE_BEFORE) ahead of after-image (UPDATE_AFTER)
at an equal seq.
+ List<KV<DestinationShard, KV<byte[], CdcRecord>>> input =
+ ImmutableList.of(
+ element("db.t", 0, 1, "b7", 7L, ValueKind.UPDATE_AFTER),
+ element("db.t", 0, 1, "a5", 5L, ValueKind.INSERT),
+ element("db.t", 0, 2, "ua9", 9L, ValueKind.UPDATE_AFTER),
+ element("db.t", 0, 2, "ub9", 9L, ValueKind.UPDATE_BEFORE));
+
+ CommitWindows.Result r =
p.apply(Create.of(input).withCoder(INPUT_CODER)).apply(batchWindows());
+
+ PAssert.that(r.getSortedGroups())
+ .satisfies(
+ groups -> {
+ KV<DestinationShard, Iterable<KV<byte[], CdcRecord>>> g =
+ Iterables.getOnlyElement(groups);
+ assertThat(g.getKey(), equalTo(DestinationShard.of("db.t", 0)));
+ assertThat(seqsOf(g), contains(5L, 7L, 9L, 9L));
+ assertThat(
+ kindsOf(g),
+ contains(
+ ValueKind.INSERT,
+ ValueKind.UPDATE_AFTER,
+ ValueKind.UPDATE_BEFORE,
+ ValueKind.UPDATE_AFTER));
+ return null;
+ });
+ PAssert.that(r.getDeadLetterRows()).empty();
+ p.run().waitUntilFinish();
+ }
+
+ //
---------------------------------------------------------------------------------------------
+ // 2. Shard separation: two shards -> two groups
+ //
---------------------------------------------------------------------------------------------
+
+ @Test
+ public void differentShardsProduceSeparateGroups() {
+ List<KV<DestinationShard, KV<byte[], CdcRecord>>> input =
+ ImmutableList.of(
+ element("db.t", 0, 1, "a", 1L, ValueKind.INSERT),
+ element("db.t", 1, 2, "b", 2L, ValueKind.INSERT));
+
+ CommitWindows.Result r =
p.apply(Create.of(input).withCoder(INPUT_CODER)).apply(batchWindows());
+
+ PAssert.that(r.getSortedGroups())
+ .satisfies(
+ groups -> {
+ List<DestinationShard> keys = new ArrayList<>();
+ for (KV<DestinationShard, Iterable<KV<byte[], CdcRecord>>> g :
groups) {
+ keys.add(g.getKey());
+ assertThat(Iterables.size(g.getValue()), equalTo(1));
+ }
+ assertThat(
+ keys,
+ containsInAnyOrder(
+ DestinationShard.of("db.t", 0),
DestinationShard.of("db.t", 1)));
+ return null;
+ });
+ PAssert.that(r.getDeadLetterRows()).empty();
+ p.run().waitUntilFinish();
+ }
+
+ //
---------------------------------------------------------------------------------------------
+ // 4. Streaming: two event-time windows -> two groups, each sorted
+ //
---------------------------------------------------------------------------------------------
+
+ @Test
+ public void streamingWindowsProduceSeparateSortedGroups() {
+ Instant t0 = new Instant(0);
+ TestStream<KV<DestinationShard, KV<byte[], CdcRecord>>> stream =
+ TestStream.create(INPUT_CODER)
+ // Window [0, 60s): two records added out of sequence order.
+ .addElements(
+ at(
+ t0.plus(Duration.millis(1_000)),
+ "db.t",
+ 0,
+ 1,
+ "w1b",
+ 7L,
+ ValueKind.UPDATE_AFTER))
+ .addElements(
+ at(t0.plus(Duration.millis(1_500)), "db.t", 0, 1, "w1a", 5L,
ValueKind.INSERT))
+ .advanceWatermarkTo(t0.plus(Duration.standardSeconds(70))) //
close window [0, 60s)
+ // Window [60s, 120s):
+ .addElements(
+ at(t0.plus(Duration.millis(61_000)), "db.t", 0, 2, "w2", 9L,
ValueKind.INSERT))
+ .advanceWatermarkToInfinity();
+
+ CommitWindows.Result r = p.apply(stream).apply(streamingWindows());
+
+ PAssert.that(r.getSortedGroups())
+ .satisfies(
+ groups -> {
+ Set<List<Long>> seqGroups = new HashSet<>();
+ for (KV<DestinationShard, Iterable<KV<byte[], CdcRecord>>> g :
groups) {
+ assertThat(g.getKey(), equalTo(DestinationShard.of("db.t",
0)));
+ seqGroups.add(seqsOf(g));
+ }
+ assertThat(
+ seqGroups,
+ equalTo(ImmutableSet.of(ImmutableList.of(5L, 7L),
ImmutableList.of(9L))));
+ return null;
+ });
+ PAssert.that(r.getDeadLetterRows()).empty();
+ p.run().waitUntilFinish();
+ }
+
+ //
---------------------------------------------------------------------------------------------
+ // 5. Late non-first pane -> replayable dead letters (+ metric); on-time
group unaffected
+ //
---------------------------------------------------------------------------------------------
+
+ @Test
+ public void lateNonFirstPaneDivertsToReplayableDeadLetters() {
+ Instant t0 = new Instant(0);
+ TestStream<KV<DestinationShard, KV<byte[], CdcRecord>>> stream =
+ TestStream.create(INPUT_CODER)
+ // On-time element; the watermark then closes window [0, 60s) and
fires its
+ // on-time pane.
+ .addElements(
+ at(t0.plus(Duration.millis(1_000)), "db.t", 0, 1, "first", 1L,
ValueKind.INSERT))
+ .advanceWatermarkTo(t0.plus(Duration.standardSeconds(70)))
+ // Two elements timestamped INSIDE the already-fired window arrive
late (within
+ // allowed lateness): every record of a non-first late pane
becomes one dead letter.
+ .addElements(
+ at(
+ t0.plus(Duration.millis(2_000)),
+ "db.t",
+ 0,
+ 2,
+ "late1",
+ 2L,
+ ValueKind.UPDATE_AFTER),
+ at(t0.plus(Duration.millis(3_000)), "db.t", 0, 3, "late2", 3L,
ValueKind.DELETE))
+ .advanceWatermarkToInfinity();
+
+ CommitWindows.Result r = p.apply(stream).apply(streamingWindows());
+
+ assertThat(r.getDeadLetterSchema(), equalTo(EXPECTED_DEAD_LETTER_SCHEMA));
+
+ PAssert.that(r.getSortedGroups())
+ .satisfies(
+ groups -> {
+ KV<DestinationShard, Iterable<KV<byte[], CdcRecord>>> g =
+ Iterables.getOnlyElement(groups);
+ assertThat(g.getKey(), equalTo(DestinationShard.of("db.t", 0)));
+ assertThat(seqsOf(g), contains(1L)); // only the on-time record
+ return null;
+ });
+ // Row.equals compares schemas too, so this pins the exact dead-letter
schema AND values:
+ // nested data row + change type name + sequence number + destination
string.
+ PAssert.that(r.getDeadLetterRows())
+ .containsInAnyOrder(
+ Row.withSchema(EXPECTED_DEAD_LETTER_SCHEMA)
+ .addValues(data(2, "late1"), "UPDATE_AFTER", 2L, "db.t")
+ .build(),
+ Row.withSchema(EXPECTED_DEAD_LETTER_SCHEMA)
+ .addValues(data(3, "late2"), "DELETE", 3L, "db.t")
+ .build());
+
+ PipelineResult result = p.run();
+ result.waitUntilFinish();
+ assertThat(counterTotal(result, "deadLetterRecords"), equalTo(2L));
+ }
+
+ //
---------------------------------------------------------------------------------------------
+ // 6. First-late pane is diverted too
+ //
---------------------------------------------------------------------------------------------
+
+ /**
+ * A window whose FIRST pane is late is dead-lettered like any other late
pane. The pane's timing
+ * is a fact about one {@code (destination, shard, window)}; the committer's
already-committed
+ * skip is per {@code (destination, window)}. A shard that saw no on-time
data therefore has a
+ * first pane that is late even when its destination-window committed long
ago, so {@code
+ * isFirst()} cannot be used to let records through.
+ */
+ @Test
+ public void firstLatePaneIsAlsoDivertedToDeadLetters() {
+ Instant t0 = new Instant(0);
+ TestStream<KV<DestinationShard, KV<byte[], CdcRecord>>> stream =
+ TestStream.create(INPUT_CODER)
+ // The watermark passes the end of window [0, 60s) with NO data
for this key ...
+ .advanceWatermarkTo(t0.plus(Duration.standardSeconds(70)))
+ // ... then the window's ONLY records arrive, late.
+ .addElements(
+ at(t0.plus(Duration.millis(2_000)), "db.t", 0, 1, "jitter",
4L, ValueKind.INSERT))
+ .advanceWatermarkToInfinity();
+
+ CommitWindows.Result r = p.apply(stream).apply(streamingWindows());
+
+ PAssert.that(r.getSortedGroups()).empty();
+ PAssert.that(r.getDeadLetterRows())
+ .containsInAnyOrder(
+ Row.withSchema(EXPECTED_DEAD_LETTER_SCHEMA)
+ .addValues(data(1, "jitter"), "INSERT", 4L, "db.t")
+ .build());
+
+ PipelineResult result = p.run();
+ result.waitUntilFinish();
+ assertThat(counterTotal(result, "deadLetterRecords"), equalTo(1L));
+ }
+
+ //
---------------------------------------------------------------------------------------------
+ // 7. Equal sort keys: both records survive the sort
+ //
---------------------------------------------------------------------------------------------
+
+ @Test
+ public void equalSortKeysBothSurviveSort() {
+ // Identical (seq, kind) -> byte-identical sort keys; the sort must keep
both records.
+ List<KV<DestinationShard, KV<byte[], CdcRecord>>> input =
+ ImmutableList.of(
+ element("db.t", 0, 1, "first", 5L, ValueKind.INSERT),
+ element("db.t", 0, 2, "second", 5L, ValueKind.INSERT));
+
+ CommitWindows.Result r =
p.apply(Create.of(input).withCoder(INPUT_CODER)).apply(batchWindows());
+
+ PAssert.that(r.getSortedGroups())
+ .satisfies(
+ groups -> {
+ KV<DestinationShard, Iterable<KV<byte[], CdcRecord>>> g =
+ Iterables.getOnlyElement(groups);
+ assertThat(seqsOf(g), contains(5L, 5L));
+ assertThat(namesOf(g), containsInAnyOrder("first", "second"));
+ return null;
+ });
+ p.run().waitUntilFinish();
+ }
+}