This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new c65d2550fe [codegen] Support normalized key for timestamp with local 
time zone (#8821)
c65d2550fe is described below

commit c65d2550fefcb6ae1f4e03166c1f00a1e324bdfa
Author: Vova Kolmakov <[email protected]>
AuthorDate: Thu Jul 23 17:32:31 2026 +0700

    [codegen] Support normalized key for timestamp with local time zone (#8821)
---
 .../apache/paimon/codegen/SortCodeGenerator.scala  | 17 ++--
 .../paimon/codegen/NormalizedKeyComputerTest.java  | 95 ++++++++++++++++------
 2 files changed, 80 insertions(+), 32 deletions(-)

diff --git 
a/paimon-codegen/src/main/scala/org/apache/paimon/codegen/SortCodeGenerator.scala
 
b/paimon-codegen/src/main/scala/org/apache/paimon/codegen/SortCodeGenerator.scala
index 6ed773fb3b..16bc719d15 100644
--- 
a/paimon-codegen/src/main/scala/org/apache/paimon/codegen/SortCodeGenerator.scala
+++ 
b/paimon-codegen/src/main/scala/org/apache/paimon/codegen/SortCodeGenerator.scala
@@ -20,7 +20,8 @@ package org.apache.paimon.codegen
 
 import org.apache.paimon.codegen.GenerateUtils.{newName, ROW_DATA, SEGMENT}
 import org.apache.paimon.data.{BinaryRow, Decimal, Timestamp}
-import org.apache.paimon.types.{DataType, DecimalType, RowType, TimestampType}
+import org.apache.paimon.types.{DataType, DecimalType, 
LocalZonedTimestampType, RowType, TimestampType}
+import org.apache.paimon.types.DataTypeChecks.getPrecision
 import org.apache.paimon.types.DataTypeRoot._
 import org.apache.paimon.utils.{SortUtil, TypeUtils}
 
@@ -368,8 +369,8 @@ class SortCodeGenerator(val input: RowType, val sortSpec: 
SortSpec) {
     t match {
       case dt: DecimalType =>
         s"get$prefix($index, ${dt.getPrecision}, ${dt.getScale})"
-      case dt: TimestampType =>
-        s"get$prefix($index, ${dt.getPrecision})"
+      case _: TimestampType | _: LocalZonedTimestampType =>
+        s"get$prefix($index, ${getPrecision(t)})"
       case _ =>
         s"get$prefix($index)"
     }
@@ -392,7 +393,7 @@ class SortCodeGenerator(val input: RowType, val sortSpec: 
SortSpec) {
     case DECIMAL => "Decimal"
     case DATE => "Int"
     case TIME_WITHOUT_TIME_ZONE => "Int"
-    case TIMESTAMP_WITHOUT_TIME_ZONE => "Timestamp"
+    case TIMESTAMP_WITHOUT_TIME_ZONE | TIMESTAMP_WITH_LOCAL_TIME_ZONE => 
"Timestamp"
     case _ => null
   }
 
@@ -410,7 +411,7 @@ class SortCodeGenerator(val input: RowType, val sortSpec: 
SortSpec) {
     t.getTypeRoot match {
       case _ if TypeUtils.isPrimitive(t) => true
       case VARCHAR | CHAR | VARBINARY | BINARY | DATE | TIME_WITHOUT_TIME_ZONE 
=> true
-      case TIMESTAMP_WITHOUT_TIME_ZONE => true
+      case TIMESTAMP_WITHOUT_TIME_ZONE | TIMESTAMP_WITH_LOCAL_TIME_ZONE => true
       case DECIMAL => 
Decimal.isCompact(t.asInstanceOf[DecimalType].getPrecision)
       case _ => false
     }
@@ -425,11 +426,11 @@ class SortCodeGenerator(val input: RowType, val sortSpec: 
SortSpec) {
       case FLOAT => 4
       case DOUBLE => 8
       case BIGINT => 8
-      case TIMESTAMP_WITHOUT_TIME_ZONE
-          if Timestamp.isCompact(t.asInstanceOf[TimestampType].getPrecision) =>
+      case TIMESTAMP_WITHOUT_TIME_ZONE | TIMESTAMP_WITH_LOCAL_TIME_ZONE
+          if Timestamp.isCompact(getPrecision(t)) =>
         8
       // non-compact timestamp: millisecond (8) + nanoOfMillisecond (4)
-      case TIMESTAMP_WITHOUT_TIME_ZONE => 12
+      case TIMESTAMP_WITHOUT_TIME_ZONE | TIMESTAMP_WITH_LOCAL_TIME_ZONE => 12
       case DATE => 4
       case TIME_WITHOUT_TIME_ZONE => 4
       case DECIMAL if 
Decimal.isCompact(t.asInstanceOf[DecimalType].getPrecision) => 8
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/codegen/NormalizedKeyComputerTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/codegen/NormalizedKeyComputerTest.java
index 770b19be57..4ef2207446 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/codegen/NormalizedKeyComputerTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/codegen/NormalizedKeyComputerTest.java
@@ -29,49 +29,73 @@ import org.apache.paimon.memory.MemorySegment;
 import org.apache.paimon.memory.MemorySegmentPool;
 import org.apache.paimon.sort.BinaryInMemorySortBuffer;
 import org.apache.paimon.types.DataType;
+import org.apache.paimon.types.DataTypes;
 import org.apache.paimon.utils.MutableObjectIterator;
 
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.EnumSource;
 
 import java.util.ArrayList;
 import java.util.Collections;
 import java.util.List;
 import java.util.Random;
+import java.util.function.IntFunction;
 
 import static org.apache.paimon.codegen.CodeGenUtils.newNormalizedKeyComputer;
 import static org.apache.paimon.codegen.CodeGenUtils.newRecordComparator;
-import static org.apache.paimon.types.DataTypes.TIMESTAMP;
 import static org.assertj.core.api.Assertions.assertThat;
 
 /**
  * Tests the {@link NormalizedKeyComputer} generated by {@link 
SortCodeGenerator} for timestamps,
- * covering non-compact precisions (sub-millisecond ordering).
+ * covering non-compact precisions (sub-millisecond ordering). Both timestamp 
roots share the same
+ * {@link Timestamp} runtime representation, so every case is run for {@code 
TIMESTAMP} and {@code
+ * TIMESTAMP WITH LOCAL TIME ZONE}.
  */
 class NormalizedKeyComputerTest {
 
     private static final int NON_COMPACT_PRECISION = 9;
 
-    @Test
-    public void testCompactTimestampKeyMetadata() {
-        NormalizedKeyComputer computer = timestampKeyComputer(3);
+    /** The two timestamp roots that share the {@link Timestamp} 
representation and key layout. */
+    private enum TimestampKind {
+        WITHOUT_TIME_ZONE(DataTypes::TIMESTAMP),
+        WITH_LOCAL_TIME_ZONE(DataTypes::TIMESTAMP_WITH_LOCAL_TIME_ZONE);
+
+        private final IntFunction<DataType> factory;
+
+        TimestampKind(IntFunction<DataType> factory) {
+            this.factory = factory;
+        }
+
+        DataType of(int precision) {
+            return factory.apply(precision);
+        }
+    }
+
+    @ParameterizedTest
+    @EnumSource(TimestampKind.class)
+    public void testCompactTimestampKeyMetadata(TimestampKind kind) {
+        NormalizedKeyComputer computer = keyComputer(kind.of(3));
         // 1 null-aware byte + 8 bytes millisecond
         assertThat(computer.getNumKeyBytes()).isEqualTo(9);
         assertThat(computer.isKeyFullyDetermines()).isTrue();
     }
 
-    @Test
-    public void testNonCompactTimestampKeyMetadata() {
+    @ParameterizedTest
+    @EnumSource(TimestampKind.class)
+    public void testNonCompactTimestampKeyMetadata(TimestampKind kind) {
         for (int precision : new int[] {4, 6, 9}) {
-            NormalizedKeyComputer computer = timestampKeyComputer(precision);
+            NormalizedKeyComputer computer = keyComputer(kind.of(precision));
             // 1 null-aware byte + 8 bytes millisecond + 4 bytes 
nanoOfMillisecond
             assertThat(computer.getNumKeyBytes()).as("precision %d", 
precision).isEqualTo(13);
             assertThat(computer.isKeyFullyDetermines()).as("precision %d", 
precision).isTrue();
         }
     }
 
-    @Test
-    public void testSubMillisecondOrdering() {
-        NormalizedKeyComputer computer = 
timestampKeyComputer(NON_COMPACT_PRECISION);
+    @ParameterizedTest
+    @EnumSource(TimestampKind.class)
+    public void testSubMillisecondOrdering(TimestampKind kind) {
+        NormalizedKeyComputer computer = 
keyComputer(kind.of(NON_COMPACT_PRECISION));
         Timestamp lo = Timestamp.fromEpochMillis(1000, 111_111);
         Timestamp hi = Timestamp.fromEpochMillis(1000, 222_222);
 
@@ -80,9 +104,10 @@ class NormalizedKeyComputerTest {
         assertThat(normalizedCompare(computer, lo, lo)).isZero();
     }
 
-    @Test
-    public void testCompareKeyMatchesCompareToForAllPairs() {
-        NormalizedKeyComputer computer = 
timestampKeyComputer(NON_COMPACT_PRECISION);
+    @ParameterizedTest
+    @EnumSource(TimestampKind.class)
+    public void testCompareKeyMatchesCompareToForAllPairs(TimestampKind kind) {
+        NormalizedKeyComputer computer = 
keyComputer(kind.of(NON_COMPACT_PRECISION));
         List<Timestamp> values = timestampSpread();
         for (Timestamp a : values) {
             for (Timestamp b : values) {
@@ -93,8 +118,9 @@ class NormalizedKeyComputerTest {
         }
     }
 
-    @Test
-    public void testEndToEndSortByNonCompactTimestamp() throws Exception {
+    @ParameterizedTest
+    @EnumSource(TimestampKind.class)
+    public void testEndToEndSortByNonCompactTimestamp(TimestampKind kind) 
throws Exception {
         List<Timestamp> values = new ArrayList<>();
         for (Timestamp value : timestampSpread()) {
             // duplicate each value so equal normalized keys are exercised
@@ -106,12 +132,32 @@ class NormalizedKeyComputerTest {
         List<Timestamp> expected = new ArrayList<>(values);
         expected.sort(Timestamp::compareTo);
 
-        
assertThat(sortThroughBuffer(values)).containsExactlyElementsOf(expected);
+        assertThat(sortThroughBuffer(values, kind.of(NON_COMPACT_PRECISION), 
NON_COMPACT_PRECISION))
+                .containsExactlyElementsOf(expected);
+    }
+
+    @Test
+    public void testLocalZonedKeyMatchesTimestampKey() {
+        // both roots must emit byte-for-byte identical normalized keys for 
the same value
+        NormalizedKeyComputer tsComputer =
+                
keyComputer(TimestampKind.WITHOUT_TIME_ZONE.of(NON_COMPACT_PRECISION));
+        NormalizedKeyComputer ltzComputer =
+                
keyComputer(TimestampKind.WITH_LOCAL_TIME_ZONE.of(NON_COMPACT_PRECISION));
+        for (Timestamp value : timestampSpread()) {
+            assertThat(keyBytes(ltzComputer, value))
+                    .as("key bytes for %s", value)
+                    .isEqualTo(keyBytes(tsComputer, value));
+        }
+    }
+
+    private static NormalizedKeyComputer keyComputer(DataType type) {
+        return newNormalizedKeyComputer(Collections.singletonList(type), new 
int[] {0});
     }
 
-    private static NormalizedKeyComputer timestampKeyComputer(int precision) {
-        return newNormalizedKeyComputer(
-                Collections.singletonList(TIMESTAMP(precision)), new int[] 
{0});
+    private static byte[] keyBytes(NormalizedKeyComputer computer, Timestamp 
value) {
+        byte[] bytes = new byte[computer.getNumKeyBytes()];
+        computer.putKey(GenericRow.of(value), MemorySegment.wrap(bytes), 0);
+        return bytes;
     }
 
     private static int normalizedCompare(NormalizedKeyComputer computer, 
Timestamp a, Timestamp b) {
@@ -136,8 +182,9 @@ class NormalizedKeyComputerTest {
     }
 
     @SuppressWarnings({"unchecked", "rawtypes"})
-    private static List<Timestamp> sortThroughBuffer(List<Timestamp> values) 
throws Exception {
-        List<DataType> fieldTypes = 
Collections.singletonList(TIMESTAMP(NON_COMPACT_PRECISION));
+    private static List<Timestamp> sortThroughBuffer(
+            List<Timestamp> values, DataType type, int precision) throws 
Exception {
+        List<DataType> fieldTypes = Collections.singletonList(type);
         BinaryRowSerializer serializer = new BinaryRowSerializer(1);
         MemorySegmentPool pool =
                 new HeapMemorySegmentPool(1024 * 1024, 
MemorySegmentPool.DEFAULT_PAGE_SIZE);
@@ -152,7 +199,7 @@ class NormalizedKeyComputerTest {
         BinaryRowWriter writer = new BinaryRowWriter(row);
         for (Timestamp value : values) {
             writer.reset();
-            writer.writeTimestamp(0, value, NON_COMPACT_PRECISION);
+            writer.writeTimestamp(0, value, precision);
             writer.complete();
             assertThat(buffer.write(row)).isTrue();
         }
@@ -162,7 +209,7 @@ class NormalizedKeyComputerTest {
         BinaryRow reuse = serializer.createInstance();
         BinaryRow next;
         while ((next = iterator.next(reuse)) != null) {
-            sorted.add(next.getTimestamp(0, NON_COMPACT_PRECISION));
+            sorted.add(next.getTimestamp(0, precision));
         }
         buffer.clear();
         return sorted;

Reply via email to