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 c5a2eac1178 [IcebergIO] CDC record and sort key (#39981)
c5a2eac1178 is described below

commit c5a2eac11780c26e31482ff7c6e75d2b7c00ceda
Author: Ahmed Abualsaud <[email protected]>
AuthorDate: Thu Sep 3 23:03:27 2026 +0200

    [IcebergIO] CDC record and sort key (#39981)
---
 sdks/java/io/iceberg/build.gradle                  |   1 +
 .../beam/sdk/io/iceberg/cdc/sink/CdcRecord.java    |  84 ++++++++
 .../sdk/io/iceberg/cdc/sink/CdcRecordCoder.java    |  99 +++++++++
 .../beam/sdk/io/iceberg/cdc/sink/CdcSortKey.java   |  77 +++++++
 .../beam/sdk/io/iceberg/cdc/sink/package-info.java |  23 +++
 .../io/iceberg/cdc/sink/CdcRecordCoderTest.java    | 113 ++++++++++
 .../sdk/io/iceberg/cdc/sink/CdcSortKeyTest.java    | 230 +++++++++++++++++++++
 7 files changed, 627 insertions(+)

diff --git a/sdks/java/io/iceberg/build.gradle 
b/sdks/java/io/iceberg/build.gradle
index e2e8a12d01e..7340d1eeb41 100644
--- a/sdks/java/io/iceberg/build.gradle
+++ b/sdks/java/io/iceberg/build.gradle
@@ -49,6 +49,7 @@ dependencies {
     implementation library.java.avro
     implementation library.java.slf4j_api
     implementation library.java.joda_time
+    implementation library.java.guava
     implementation "org.apache.parquet:parquet-column:$parquet_version"
     implementation "org.apache.parquet:parquet-hadoop:$parquet_version"
     implementation "org.apache.parquet:parquet-common:$parquet_version"
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcRecord.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcRecord.java
new file mode 100644
index 00000000000..1aa82914fde
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcRecord.java
@@ -0,0 +1,84 @@
+/*
+ * 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.common.base.MoreObjects;
+import java.util.Objects;
+import org.apache.beam.sdk.values.Row;
+import org.apache.beam.sdk.values.ValueKind;
+import org.checkerframework.checker.nullness.qual.Nullable;
+
+/**
+ * One change record carried through the CDC sink's shuffle.
+ *
+ * <p>{@link ValueKind} is reified because it's not preserved across a {@code 
GroupByKey}.
+ */
+final class CdcRecord {
+
+  private final Row data;
+  private final ValueKind kind;
+  private final long sequenceNumber;
+
+  private CdcRecord(Row data, ValueKind kind, long sequenceNumber) {
+    this.data = data;
+    this.kind = kind;
+    this.sequenceNumber = sequenceNumber;
+  }
+
+  public static CdcRecord of(Row data, ValueKind kind, long sequenceNumber) {
+    return new CdcRecord(data, kind, sequenceNumber);
+  }
+
+  public Row getData() {
+    return data;
+  }
+
+  public ValueKind getKind() {
+    return kind;
+  }
+
+  public long getSequenceNumber() {
+    return sequenceNumber;
+  }
+
+  @Override
+  public boolean equals(@Nullable Object o) {
+    if (this == o) {
+      return true;
+    }
+    if (!(o instanceof CdcRecord)) {
+      return false;
+    }
+    CdcRecord that = (CdcRecord) o;
+    return sequenceNumber == that.sequenceNumber && kind == that.kind && 
data.equals(that.data);
+  }
+
+  @Override
+  public int hashCode() {
+    return Objects.hash(data, kind, sequenceNumber);
+  }
+
+  @Override
+  public String toString() {
+    return MoreObjects.toStringHelper(CdcRecord.class)
+        .add("data", data)
+        .add("kind", kind)
+        .add("sequenceNumber", sequenceNumber)
+        .toString();
+  }
+}
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcRecordCoder.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcRecordCoder.java
new file mode 100644
index 00000000000..d477fe10050
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcRecordCoder.java
@@ -0,0 +1,99 @@
+/*
+ * 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 java.io.IOException;
+import java.io.InputStream;
+import java.io.OutputStream;
+import org.apache.beam.sdk.coders.Coder;
+import org.apache.beam.sdk.coders.CustomCoder;
+import org.apache.beam.sdk.coders.RowCoder;
+import org.apache.beam.sdk.coders.ValueKindCoder;
+import org.apache.beam.sdk.coders.VarLongCoder;
+import org.apache.beam.sdk.schemas.Schema;
+import org.apache.beam.sdk.values.Row;
+import org.apache.beam.sdk.values.ValueKind;
+import org.checkerframework.checker.nullness.qual.Nullable;
+
+/**
+ * {@link CdcRecord} carries a {@link Row} field whose schema is known only at 
pipeline-construction
+ * time. We need a custom coder because {@code AutoValueSchema} only infers 
schemas at the class
+ * level and cannot infer a dynamic {@link Row} field, so {@code 
@DefaultSchema} alone cannot
+ * produce a working coder for {@link CdcRecord}.
+ */
+final class CdcRecordCoder extends CustomCoder<CdcRecord> {
+
+  private final RowCoder dataCoder;
+  private final ValueKindCoder kindCoder = ValueKindCoder.of();
+  private final VarLongCoder seqCoder = VarLongCoder.of();
+
+  private CdcRecordCoder(Schema dataSchema) {
+    this.dataCoder = RowCoder.of(dataSchema);
+  }
+
+  public static CdcRecordCoder of(Schema dataSchema) {
+    return new CdcRecordCoder(dataSchema);
+  }
+
+  public Schema getDataSchema() {
+    return dataCoder.getSchema();
+  }
+
+  @Override
+  public void encode(CdcRecord value, OutputStream outStream) throws 
IOException {
+    dataCoder.encode(value.getData(), outStream);
+    kindCoder.encode(value.getKind(), outStream);
+    seqCoder.encode(value.getSequenceNumber(), outStream);
+  }
+
+  @Override
+  public CdcRecord decode(InputStream inStream) throws IOException {
+    Row data = dataCoder.decode(inStream);
+    ValueKind kind = kindCoder.decode(inStream);
+    long seq = seqCoder.decode(inStream);
+    return CdcRecord.of(data, kind, seq);
+  }
+
+  @Override
+  public void verifyDeterministic() throws NonDeterministicException {
+    Coder.verifyDeterministic(this, "Data coder must be deterministic", 
dataCoder);
+  }
+
+  @Override
+  public boolean consistentWithEquals() {
+    // decode() rebuilds the row with this coder's schema object, which may 
differ from the
+    // original row's; false is always safe.
+    return false;
+  }
+
+  @Override
+  public boolean equals(@Nullable Object o) {
+    if (this == o) {
+      return true;
+    }
+    if (o == null || getClass() != o.getClass()) {
+      return false;
+    }
+    return getDataSchema().equals(((CdcRecordCoder) o).getDataSchema());
+  }
+
+  @Override
+  public int hashCode() {
+    return getDataSchema().hashCode();
+  }
+}
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcSortKey.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcSortKey.java
new file mode 100644
index 00000000000..21a297b3824
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcSortKey.java
@@ -0,0 +1,77 @@
+/*
+ * 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 java.nio.ByteBuffer;
+import java.util.Arrays;
+import org.apache.beam.sdk.values.ValueKind;
+
+/**
+ * Builds the byte-comparable secondary sort key used by {@code SortValues} 
(extensions/sorter) to
+ * order each (destination, shard, window) group: one primary key's records 
come out contiguous,
+ * ordered by sequence number then {@link #kindRank(ValueKind)} within the key.
+ *
+ * <p>The key is {@code [pkLen:4][pkBytes][seq ^ 
Long.MIN_VALUE:8][kindRank:1]}, big-endian.
+ */
+final class CdcSortKey {
+
+  private CdcSortKey() {}
+
+  /** Ranks change kinds so before-images sort before after-images at an equal 
{@code seq}. */
+  public static byte kindRank(ValueKind kind) {
+    switch (kind) {
+      case UPDATE_BEFORE:
+        return 0;
+      case DELETE:
+        return 1;
+      case UPDATE_AFTER:
+        return 2;
+      case INSERT:
+        return 3;
+      default:
+        throw new IllegalArgumentException("Unknown ValueKind: " + kind);
+    }
+  }
+
+  /**
+   * Encodes the deterministic, byte-comparable sort key {@code 
[pkLen:4][pkBytes][seq ^
+   * Long.MIN_VALUE:8][kindRank:1]} for one CDC record.
+   *
+   * <p>SortValues compares unsigned lexicographic byte order. The length 
prefix is needed to
+   * accurately compare two primary keys of varying byte-lengths. Flipping the 
sequence number's
+   * sign bit makes unsigned byte order match signed numeric order. kindRank 
breaks equal-seq ties.
+   */
+  public static byte[] encode(byte[] pkBytes, long seq, ValueKind kind) {
+    return ByteBuffer.allocate(4 + pkBytes.length + 9)
+        .putInt(pkBytes.length)
+        .put(pkBytes)
+        .putLong(seq ^ Long.MIN_VALUE)
+        .put(kindRank(kind))
+        .array();
+  }
+
+  /**
+   * Whether two encoded sort keys carry the same primary key, compared on the 
raw {@code
+   * [pkLen:4][pkBytes]} prefix.
+   */
+  public static boolean samePk(byte[] a, byte[] b) {
+    int aPkEnd = 4 + ByteBuffer.wrap(a).getInt(0);
+    int bPkEnd = 4 + ByteBuffer.wrap(b).getInt(0);
+    return Arrays.equals(a, 0, aPkEnd, b, 0, bPkEnd);
+  }
+}
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/package-info.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/package-info.java
new file mode 100644
index 00000000000..2bcd8b5f622
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/package-info.java
@@ -0,0 +1,23 @@
+/*
+ * 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.
+ */
+
+/**
+ * A change data capture (CDC) write sink for Apache Iceberg: applies inserts, 
updates, and deletes
+ * to a table, identifying rows by equality columns.
+ */
+package org.apache.beam.sdk.io.iceberg.cdc.sink;
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcRecordCoderTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcRecordCoderTest.java
new file mode 100644
index 00000000000..a53b98a10aa
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcRecordCoderTest.java
@@ -0,0 +1,113 @@
+/*
+ * 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.junit.Assert.assertFalse;
+import static org.junit.Assert.assertThrows;
+
+import java.io.ByteArrayInputStream;
+import java.io.ByteArrayOutputStream;
+import org.apache.beam.sdk.coders.CoderException;
+import org.apache.beam.sdk.coders.RowCoder;
+import org.apache.beam.sdk.coders.VarIntCoder;
+import org.apache.beam.sdk.coders.VarLongCoder;
+import org.apache.beam.sdk.schemas.Schema;
+import org.apache.beam.sdk.testing.CoderProperties;
+import org.apache.beam.sdk.values.Row;
+import org.apache.beam.sdk.values.ValueKind;
+import org.hamcrest.Matchers;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
+
+/** Tests for {@link CdcRecordCoder}. */
+@RunWith(JUnit4.class)
+public class CdcRecordCoderTest {
+
+  private static final Schema DATA_SCHEMA =
+      
Schema.builder().addInt32Field("id").addStringField("name").addStringField("data").build();
+
+  @Test
+  public void roundTripsAllValueKinds() throws Exception {
+    for (ValueKind kind : ValueKind.values()) {
+      CdcRecord record =
+          CdcRecord.of(Row.withSchema(DATA_SCHEMA).addValues(1, "a", 
"x").build(), kind, 12L);
+      CoderProperties.coderDecodeEncodeEqual(CdcRecordCoder.of(DATA_SCHEMA), 
record);
+    }
+  }
+
+  /** The coder's own properties: deterministic, serializable, schema-keyed 
equality. */
+  @Test
+  public void coderIsDeterministicSerializableAndSchemaKeyed() throws 
Exception {
+    CdcRecordCoder.of(DATA_SCHEMA).verifyDeterministic();
+    CoderProperties.coderSerializable(CdcRecordCoder.of(DATA_SCHEMA));
+    // decode() rebuilds the row with the coder's schema object, which may 
differ from the
+    // original row's, so the coder must not claim consistency with equals.
+    assertFalse(CdcRecordCoder.of(DATA_SCHEMA).consistentWithEquals());
+    assertThat(CdcRecordCoder.of(DATA_SCHEMA).getDataSchema(), 
Matchers.equalTo(DATA_SCHEMA));
+
+    Schema otherSchema = Schema.builder().addInt32Field("other").build();
+    assertThat(CdcRecordCoder.of(DATA_SCHEMA), 
Matchers.equalTo(CdcRecordCoder.of(DATA_SCHEMA)));
+    assertThat(
+        CdcRecordCoder.of(DATA_SCHEMA).hashCode(),
+        Matchers.equalTo(CdcRecordCoder.of(DATA_SCHEMA).hashCode()));
+    assertThat(
+        CdcRecordCoder.of(DATA_SCHEMA),
+        Matchers.not(Matchers.equalTo(CdcRecordCoder.of(otherSchema))));
+  }
+
+  @Test
+  public void encodePinnedWireMapping() throws Exception {
+    // The stream layout is [row bytes][kind VarInt][seq VarLong]. Codes 0-3 
each fit in a single
+    // VarInt byte equal to the code itself, so the byte immediately following 
the row bytes must
+    // be the pinned code for that ValueKind.
+    Row row = Row.withSchema(DATA_SCHEMA).addValues(1, "a", "x").build();
+    ByteArrayOutputStream rowOnly = new ByteArrayOutputStream();
+    RowCoder.of(DATA_SCHEMA).encode(row, rowOnly);
+    int rowLen = rowOnly.toByteArray().length;
+
+    assertPinnedKindCode(row, ValueKind.INSERT, rowLen, 0);
+    assertPinnedKindCode(row, ValueKind.UPDATE_BEFORE, rowLen, 1);
+    assertPinnedKindCode(row, ValueKind.UPDATE_AFTER, rowLen, 2);
+    assertPinnedKindCode(row, ValueKind.DELETE, rowLen, 3);
+  }
+
+  private static void assertPinnedKindCode(Row row, ValueKind kind, int 
rowLen, int expectedCode)
+      throws Exception {
+    ByteArrayOutputStream out = new ByteArrayOutputStream();
+    CdcRecordCoder.of(DATA_SCHEMA).encode(CdcRecord.of(row, kind, 1L), out);
+    byte[] bytes = out.toByteArray();
+    assertThat(bytes[rowLen] & 0xFF, Matchers.equalTo(expectedCode));
+  }
+
+  @Test
+  public void decodeRejectsUnknownKindCode() throws Exception {
+    // Hand-encode a stream with a valid data row and seq, but a ValueKind 
code (4) at the pinned
+    // mapping's boundary: one past the highest valid code (3).
+    ByteArrayOutputStream out = new ByteArrayOutputStream();
+    RowCoder.of(DATA_SCHEMA)
+        .encode(Row.withSchema(DATA_SCHEMA).addValues(1, "a", "x").build(), 
out);
+    VarIntCoder.of().encode(4, out);
+    VarLongCoder.of().encode(1L, out);
+
+    CdcRecordCoder coder = CdcRecordCoder.of(DATA_SCHEMA);
+    assertThrows(
+        CoderException.class, () -> coder.decode(new 
ByteArrayInputStream(out.toByteArray())));
+  }
+}
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcSortKeyTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcSortKeyTest.java
new file mode 100644
index 00000000000..6e13f385826
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcSortKeyTest.java
@@ -0,0 +1,230 @@
+/*
+ * 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.junit.Assert.assertArrayEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertTrue;
+
+import java.nio.ByteBuffer;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Set;
+import org.apache.beam.sdk.values.ValueKind;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.primitives.UnsignedBytes;
+import org.hamcrest.Matchers;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
+
+/** Tests for {@link CdcSortKey}. */
+@RunWith(JUnit4.class)
+public class CdcSortKeyTest {
+
+  // Ascending sequence numbers, and ValueKinds in ascending kindRank order.
+  private static final long[] SEQS_ASCENDING = {-5L, -1L, 0L, 1L, 5L};
+  private static final ValueKind[] KINDS_BY_RANK_ASCENDING = {
+    ValueKind.UPDATE_BEFORE, ValueKind.DELETE, ValueKind.UPDATE_AFTER, 
ValueKind.INSERT
+  };
+
+  // Primary keys of assorted lengths and leading bytes.
+  private static final byte[][] PKS = {{0x02}, {0x01, 0x7F}, {(byte) 0xF0, 
0x00, 0x10}};
+
+  @Test
+  public void kindRankOrdersBeforeImagesFirst() {
+    assertThat(CdcSortKey.kindRank(ValueKind.UPDATE_BEFORE), 
Matchers.equalTo((byte) 0));
+    assertThat(CdcSortKey.kindRank(ValueKind.DELETE), Matchers.equalTo((byte) 
1));
+    assertThat(CdcSortKey.kindRank(ValueKind.UPDATE_AFTER), 
Matchers.equalTo((byte) 2));
+    assertThat(CdcSortKey.kindRank(ValueKind.INSERT), Matchers.equalTo((byte) 
3));
+  }
+
+  /**
+   * The sorted keys must group each primary key's entries into one unbroken 
run, ordered by (seq,
+   * kind) within the run. Which primary key's run comes first is NOT part of 
the contract and is
+   * deliberately unasserted.
+   */
+  @Test
+  public void sortGroupsEachPkContiguouslyOrderedBySeqThenKindWithin() {
+    // Per pk, its keys in expected within-key order: seq ascending, kindRank 
breaking ties.
+    List<List<byte[]>> expectedByPk = new ArrayList<>();
+    for (byte[] pk : PKS) {
+      List<byte[]> expected = new ArrayList<>();
+      for (long seq : SEQS_ASCENDING) {
+        for (ValueKind kind : KINDS_BY_RANK_ASCENDING) {
+          expected.add(CdcSortKey.encode(pk, seq, kind));
+        }
+      }
+      expectedByPk.add(expected);
+    }
+
+    // Feed the sorter the same keys interleaved across pks (the shape a shard 
group arrives in);
+    // the sort must both regroup and reorder.
+    List<byte[]> sorted = new ArrayList<>();
+    for (int i = 0; i < expectedByPk.get(0).size(); i++) {
+      for (List<byte[]> keys : expectedByPk) {
+        sorted.add(keys.get(i));
+      }
+    }
+    sorted.sort(UnsignedBytes.lexicographicalComparator());
+
+    // Contiguity: each pk appears in exactly one run.
+    Set<Integer> seenRuns = new HashSet<>();
+    int previous = -1;
+    for (byte[] key : sorted) {
+      int pkIndex = pkIndexOf(key);
+      if (pkIndex != previous) {
+        assertTrue("entries of pk " + pkIndex + " are split across runs", 
seenRuns.add(pkIndex));
+        previous = pkIndex;
+      }
+    }
+
+    // Within each run, the (seq, kind) order pinned above.
+    for (int pkIndex = 0; pkIndex < PKS.length; pkIndex++) {
+      List<byte[]> run = new ArrayList<>();
+      for (byte[] key : sorted) {
+        if (pkIndexOf(key) == pkIndex) {
+          run.add(key);
+        }
+      }
+      List<byte[]> expected = expectedByPk.get(pkIndex);
+      assertThat(run, Matchers.hasSize(expected.size()));
+      for (int i = 0; i < expected.size(); i++) {
+        assertArrayEquals(expected.get(i), run.get(i));
+      }
+    }
+  }
+
+  /**
+   * Two primary keys where one's bytes are a strict prefix of the other's, 
with sequence numbers
+   * chosen so the unprefixed layout {@code [pkBytes][seq ^ 
Long.MIN_VALUE:8][kindRank:1]} WOULD
+   * interleave them; the length prefix must keep the short key's entries 
adjacent.
+   */
+  @Test
+  public void lengthPrefixKeepsAPrefixPkContiguous() {
+    byte[] shortPk = {0x0A};
+    byte[] longPk = {0x0A, (byte) 0x80};
+    byte[] shortLow = CdcSortKey.encode(shortPk, 0L, ValueKind.INSERT);
+    byte[] shortHigh = CdcSortKey.encode(shortPk, Long.MAX_VALUE, 
ValueKind.INSERT);
+    byte[] longMid = CdcSortKey.encode(longPk, 0L, ValueKind.INSERT);
+
+    // Fixture self-check: with the length prefix stripped, the long pk's 
entry lands BETWEEN the
+    // short pk's two entries (its second byte 0x80 ties the flipped seq 0 and 
loses to the
+    // flipped Long.MAX_VALUE).
+    List<byte[]> naive =
+        new ArrayList<>(Arrays.asList(stripLengthPrefix(shortHigh), 
stripLengthPrefix(longMid)));
+    naive.add(stripLengthPrefix(shortLow));
+    naive.sort(UnsignedBytes.lexicographicalComparator());
+    assertArrayEquals(stripLengthPrefix(shortLow), naive.get(0));
+    assertArrayEquals(stripLengthPrefix(longMid), naive.get(1));
+    assertArrayEquals(stripLengthPrefix(shortHigh), naive.get(2));
+
+    // The real keys: the short pk's entries stay adjacent, low seq first.
+    List<byte[]> sorted = new ArrayList<>(Arrays.asList(shortHigh, longMid, 
shortLow));
+    sorted.sort(UnsignedBytes.lexicographicalComparator());
+    int low = indexOfKey(sorted, shortLow);
+    int high = indexOfKey(sorted, shortHigh);
+    assertThat(high, Matchers.equalTo(low + 1));
+  }
+
+  @Test
+  public void encodeOrdersAcrossLongExtremes() {
+    byte[] pk = {0x01};
+    byte[] min = CdcSortKey.encode(pk, Long.MIN_VALUE, ValueKind.INSERT);
+    byte[] negOne = CdcSortKey.encode(pk, -1L, ValueKind.INSERT);
+    byte[] zero = CdcSortKey.encode(pk, 0L, ValueKind.INSERT);
+    byte[] one = CdcSortKey.encode(pk, 1L, ValueKind.INSERT);
+    byte[] max = CdcSortKey.encode(pk, Long.MAX_VALUE, ValueKind.INSERT);
+
+    List<byte[]> sorted = new ArrayList<>(Arrays.asList(max, zero, min, one, 
negOne));
+    sorted.sort(UnsignedBytes.lexicographicalComparator());
+
+    assertArrayEquals(min, sorted.get(0));
+    assertArrayEquals(negOne, sorted.get(1));
+    assertArrayEquals(zero, sorted.get(2));
+    assertArrayEquals(one, sorted.get(3));
+    assertArrayEquals(max, sorted.get(4));
+  }
+
+  /**
+   * The encoding is a frozen wire format: an in-place pipeline update replays 
in-flight groups
+   * through it, so changing any of these bytes silently re-orders live data. 
Update only with a
+   * migration story.
+   */
+  @Test
+  public void encodePinnedByteLayout() {
+    // pk {0x01, 0x02} -> pkLen 2; seq 5 -> flippedSeq 0x8000000000000005; 
INSERT -> kindRank 3.
+    assertArrayEquals(
+        new byte[] {0, 0, 0, 2, 1, 2, (byte) 0x80, 0, 0, 0, 0, 0, 0, 5, 3},
+        CdcSortKey.encode(new byte[] {1, 2}, 5L, ValueKind.INSERT));
+    // pk {0xAB} -> pkLen 1; seq 0 -> flippedSeq 0x8000000000000000; 
UPDATE_BEFORE -> kindRank 0.
+    assertArrayEquals(
+        new byte[] {0, 0, 0, 1, (byte) 0xAB, (byte) 0x80, 0, 0, 0, 0, 0, 0, 0, 
0},
+        CdcSortKey.encode(new byte[] {(byte) 0xAB}, 0L, 
ValueKind.UPDATE_BEFORE));
+  }
+
+  /**
+   * {@code samePk} is the writer's block-boundary test: equal for any two 
keys of one pk whatever
+   * their seq/kind bytes, unequal across pks, including the prefix-pk pair 
whose seq bytes tie.
+   */
+  @Test
+  public void samePkComparesOnlyThePkPrefix() {
+    byte[] pk = {0x0A};
+    assertTrue(
+        CdcSortKey.samePk(
+            CdcSortKey.encode(pk, Long.MIN_VALUE, ValueKind.UPDATE_BEFORE),
+            CdcSortKey.encode(pk, Long.MAX_VALUE, ValueKind.INSERT)));
+    assertFalse(
+        CdcSortKey.samePk(
+            CdcSortKey.encode(new byte[] {0x01}, 5L, ValueKind.INSERT),
+            CdcSortKey.encode(new byte[] {0x02}, 5L, ValueKind.INSERT)));
+    // A strict-prefix pk with seq bytes continuing the longer pk's bytes must 
still differ.
+    assertFalse(
+        CdcSortKey.samePk(
+            CdcSortKey.encode(pk, 0L, ValueKind.INSERT),
+            CdcSortKey.encode(new byte[] {0x0A, (byte) 0x80}, 0L, 
ValueKind.INSERT)));
+  }
+
+  /** The index in {@link #PKS} of the pk carried in {@code key}'s 
length-prefixed prefix. */
+  private static int pkIndexOf(byte[] key) {
+    int pkLen = ByteBuffer.wrap(key).getInt();
+    byte[] pk = Arrays.copyOfRange(key, 4, 4 + pkLen);
+    for (int i = 0; i < PKS.length; i++) {
+      if (Arrays.equals(PKS[i], pk)) {
+        return i;
+      }
+    }
+    throw new AssertionError("unknown pk " + Arrays.toString(pk));
+  }
+
+  /** Drops the 4-byte length prefix, leaving the naive {@code 
[pkBytes][seq][kind]} layout. */
+  private static byte[] stripLengthPrefix(byte[] key) {
+    return Arrays.copyOfRange(key, 4, key.length);
+  }
+
+  private static int indexOfKey(List<byte[]> keys, byte[] key) {
+    for (int i = 0; i < keys.size(); i++) {
+      if (Arrays.equals(keys.get(i), key)) {
+        return i;
+      }
+    }
+    throw new AssertionError("key not found: " + Arrays.toString(key));
+  }
+}

Reply via email to