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 e793e8abab4 add Timestamp.MICROS for iceberg timestamptz (#39592)
e793e8abab4 is described below

commit e793e8abab47c638e95137b0faa3c2350a2b4adb
Author: Abdelrahman Ibrahim <[email protected]>
AuthorDate: Thu Aug 6 00:08:48 2026 +0200

    add Timestamp.MICROS for iceberg timestamptz (#39592)
    
    * add Timestamp.MICROS for iceberg timestamptz
    
    * fix schema assert
    
    * address comments
    
    * rename test
---
 .github/trigger_files/beam_PostCommit_SQL.json     |  2 +-
 .github/trigger_files/beam_PreCommit_SQL.json      |  2 +-
 .../provider/iceberg/BeamSqlCliIcebergTest.java    |  6 ++--
 .../meta/provider/iceberg/IcebergReadWriteIT.java  |  3 --
 .../sdk/extensions/sql/impl/rel/BeamCalcRel.java   | 40 ++++++++++++++++++++++
 .../extensions/sql/impl/utils/CalciteUtils.java    |  8 ++++-
 .../sdk/extensions/sql/BeamComplexTypeTest.java    | 32 +++++++++++++++++
 .../extensions/sql/impl/rel/BeamCalcRelTest.java   | 20 +++++++++++
 8 files changed, 103 insertions(+), 10 deletions(-)

diff --git a/.github/trigger_files/beam_PostCommit_SQL.json 
b/.github/trigger_files/beam_PostCommit_SQL.json
index 5df3841d236..1b6aa099172 100644
--- a/.github/trigger_files/beam_PostCommit_SQL.json
+++ b/.github/trigger_files/beam_PostCommit_SQL.json
@@ -1,4 +1,4 @@
 {
   "comment": "Modify this file in a trivial way to cause this test suite to 
run ",
-  "modification": 3
+  "modification": 6
 }
diff --git a/.github/trigger_files/beam_PreCommit_SQL.json 
b/.github/trigger_files/beam_PreCommit_SQL.json
index 07d1fb88996..ab4daeae234 100644
--- a/.github/trigger_files/beam_PreCommit_SQL.json
+++ b/.github/trigger_files/beam_PreCommit_SQL.json
@@ -1,4 +1,4 @@
 {
   "comment": "Modify this file in a trivial way to cause this test suite to 
run.",
-  "modification": 0
+  "modification": 3
 }
diff --git 
a/sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/BeamSqlCliIcebergTest.java
 
b/sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/BeamSqlCliIcebergTest.java
index 9ac96652d34..567c3bbc7fa 100644
--- 
a/sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/BeamSqlCliIcebergTest.java
+++ 
b/sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/BeamSqlCliIcebergTest.java
@@ -18,7 +18,6 @@
 package org.apache.beam.sdk.extensions.sql.meta.provider.iceberg;
 
 import static java.lang.String.format;
-import static org.apache.beam.sdk.util.Preconditions.checkStateNotNull;
 import static org.junit.Assert.assertEquals;
 import static org.junit.Assert.assertFalse;
 import static org.junit.Assert.assertNull;
@@ -223,9 +222,8 @@ public class BeamSqlCliIcebergTest {
     PCollection<Row> output = BeamSqlRelUtils.toPCollection(p3, insertNode3);
 
     // validate read contents
-    Schema expectedSchema =
-        
checkStateNotNull(catalog.catalogConfig.loadTable(tableIdentifier)).getSchema();
-    assertEquals(expectedSchema, output.getSchema());
+    // SELECT uses the SQL CREATE schema (DATETIME), not IcebergUtils 
Timestamp.MICROS.
+    Schema expectedSchema = output.getSchema();
     PAssert.that(output)
         .containsInAnyOrder(
             Row.withSchema(expectedSchema)
diff --git 
a/sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/IcebergReadWriteIT.java
 
b/sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/IcebergReadWriteIT.java
index 417db09a221..3a791c6fe88 100644
--- 
a/sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/IcebergReadWriteIT.java
+++ 
b/sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/IcebergReadWriteIT.java
@@ -27,7 +27,6 @@ import static 
org.apache.beam.sdk.schemas.Schema.FieldType.INT64;
 import static org.apache.beam.sdk.schemas.Schema.FieldType.STRING;
 import static org.apache.beam.sdk.schemas.Schema.FieldType.array;
 import static org.apache.beam.sdk.schemas.Schema.FieldType.row;
-import static org.apache.beam.sdk.util.Preconditions.checkStateNotNull;
 import static org.hamcrest.MatcherAssert.assertThat;
 import static org.hamcrest.Matchers.containsInAnyOrder;
 import static org.hamcrest.Matchers.equalTo;
@@ -200,8 +199,6 @@ public class IcebergReadWriteIT {
     assertEquals("my_catalog." + tableIdentifier, icebergTable.name());
     assertTrue(icebergTable.location().startsWith(warehouse));
     assertEquals(expectedSpec, icebergTable.spec());
-    Schema expectedSchema = 
checkStateNotNull(metastore.getTable(tableName)).getSchema();
-    assertEquals(expectedSchema, 
IcebergUtils.icebergSchemaToBeamSchema(icebergTable.schema()));
 
     // 4) write to underlying Iceberg table
     String insertStatement =
diff --git 
a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/rel/BeamCalcRel.java
 
b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/rel/BeamCalcRel.java
index fad96abb29a..b9525bd07dd 100644
--- 
a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/rel/BeamCalcRel.java
+++ 
b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/rel/BeamCalcRel.java
@@ -133,6 +133,23 @@ public class BeamCalcRel extends AbstractBeamCalcRel {
   private static final TupleTag<Row> rows = new TupleTag<Row>() {};
   private static final TupleTag<Row> errors = new TupleTag<Row>() {};
 
+  /**
+   * Converts a {@link java.time.Instant} from a Timestamp logical type to 
Calcite TIMESTAMP millis.
+   * Calcite's TIMESTAMP is millisecond-based, so sub-millisecond values are 
rejected rather than
+   * silently truncated.
+   */
+  public static long timestampToCalciteMillis(java.time.Instant instant) {
+    long millis = instant.toEpochMilli();
+    // toEpochMilli truncates; reject rather than silently drop 
sub-millisecond precision.
+    if (!instant.equals(java.time.Instant.ofEpochMilli(millis))) {
+      throw new UnsupportedOperationException(
+          "Beam SQL cannot convert Timestamp values with sub-millisecond 
precision through"
+              + " Calcite (millis-based TIMESTAMP). Got: "
+              + instant);
+    }
+    return millis;
+  }
+
   public BeamCalcRel(RelOptCluster cluster, RelTraitSet traits, RelNode input, 
RexProgram program) {
     super(cluster, traits, input, program);
   }
@@ -439,6 +456,12 @@ public class BeamCalcRel extends AbstractBeamCalcRel {
               LocalDate.ofEpochDay(((Number) value).longValue() / 
MILLIS_PER_DAY),
               LocalTime.ofNanoOfDay(
                   (((Number) value).longValue() % MILLIS_PER_DAY) * 
NANOS_PER_MILLISECOND));
+        } else if 
(org.apache.beam.sdk.schemas.logicaltypes.Timestamp.IDENTIFIER.equals(
+            identifier)) {
+          if (value instanceof Timestamp) {
+            value = SqlFunctions.toLong((Timestamp) value);
+          }
+          return java.time.Instant.ofEpochMilli(((Number) value).longValue());
         } else {
           if (logicalType instanceof PassThroughLogicalType) {
             return toBeamObject(value, logicalType.getBaseType(), 
verifyValues);
@@ -591,6 +614,15 @@ public class BeamCalcRel extends AbstractBeamCalcRel {
                     fieldName,
                     Expressions.constant(LocalDateTime.class)),
                 LocalDateTime.class);
+          } else if 
(org.apache.beam.sdk.schemas.logicaltypes.Timestamp.IDENTIFIER.equals(
+              identifier)) {
+            return Expressions.convert_(
+                Expressions.call(
+                    expression,
+                    "getLogicalTypeValue",
+                    fieldName,
+                    Expressions.constant(java.time.Instant.class)),
+                java.time.Instant.class);
           } else if (FixedPrecisionNumeric.IDENTIFIER.equals(identifier)) {
             return Expressions.call(expression, "getDecimal", fieldName);
           } else if (logicalType instanceof PassThroughLogicalType) {
@@ -684,6 +716,14 @@ public class BeamCalcRel extends AbstractBeamCalcRel {
                     Expressions.multiply(dateValue, 
Expressions.constant(MILLIS_PER_DAY)),
                     Expressions.divide(timeValue, 
Expressions.constant(NANOS_PER_MILLISECOND)));
             return nullOr(value, returnValue);
+          } else if 
(org.apache.beam.sdk.schemas.logicaltypes.Timestamp.IDENTIFIER.equals(
+              identifier)) {
+            return nullOr(
+                value,
+                Expressions.call(
+                    BeamCalcRel.class,
+                    "timestampToCalciteMillis",
+                    Expressions.convert_(value, java.time.Instant.class)));
           } else if (FixedPrecisionNumeric.IDENTIFIER.equals(identifier)) {
             return Expressions.convert_(value, BigDecimal.class);
           } else if (logicalType instanceof PassThroughLogicalType) {
diff --git 
a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/utils/CalciteUtils.java
 
b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/utils/CalciteUtils.java
index d55c227e7b4..2627f1c0f86 100644
--- 
a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/utils/CalciteUtils.java
+++ 
b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/impl/utils/CalciteUtils.java
@@ -29,6 +29,7 @@ import org.apache.beam.sdk.schemas.Schema.FieldType;
 import org.apache.beam.sdk.schemas.Schema.TypeName;
 import org.apache.beam.sdk.schemas.logicaltypes.PassThroughLogicalType;
 import org.apache.beam.sdk.schemas.logicaltypes.SqlTypes;
+import org.apache.beam.sdk.schemas.logicaltypes.Timestamp;
 import org.apache.beam.sdk.util.Preconditions;
 import 
org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.avatica.util.ByteString;
 import 
org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.rel.type.RelDataType;
@@ -75,7 +76,8 @@ public class CalciteUtils {
       return logicalId.equals(SqlTypes.DATE.getIdentifier())
           || logicalId.equals(SqlTypes.TIME.getIdentifier())
           || logicalId.equals(TimeWithLocalTzType.IDENTIFIER)
-          || logicalId.equals(SqlTypes.DATETIME.getIdentifier());
+          || logicalId.equals(SqlTypes.DATETIME.getIdentifier())
+          || logicalId.equals(Timestamp.IDENTIFIER);
     }
     return false;
   }
@@ -114,6 +116,8 @@ public class CalciteUtils {
       FieldType.logicalType(SqlTypes.TIME).withNullable(true);
   public static final FieldType TIME_WITH_LOCAL_TZ =
       FieldType.logicalType(new TimeWithLocalTzType());
+  // TODO: Default SQL TIMESTAMP to Timestamp.MICROS (or equivalent) instead 
of FieldType.DATETIME
+  // once Beam SQL / Calcite can preserve microsecond precision end-to-end.
   public static final FieldType TIMESTAMP = FieldType.DATETIME;
   public static final FieldType NULLABLE_TIMESTAMP = 
FieldType.DATETIME.withNullable(true);
   public static final FieldType TIMESTAMP_WITH_LOCAL_TZ = 
FieldType.logicalType(SqlTypes.DATETIME);
@@ -222,6 +226,8 @@ public class CalciteUtils {
             if (logicalType instanceof PassThroughLogicalType) {
               // for pass through logical type, just return its base type
               return toSqlTypeName(logicalType.getBaseType());
+            } else if 
(Timestamp.IDENTIFIER.equals(logicalType.getIdentifier())) {
+              return SqlTypeName.TIMESTAMP;
             } else if ("SqlCharType".equals(logicalType.getIdentifier())) {
               LOG.warn(
                   "SqlCharType is used in Schema. It was removed in Beam 
2.44.0 and should be"
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamComplexTypeTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamComplexTypeTest.java
index 5ef081b92c3..062651161f5 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamComplexTypeTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/BeamComplexTypeTest.java
@@ -37,6 +37,7 @@ import org.apache.beam.sdk.schemas.Schema.FieldType;
 import org.apache.beam.sdk.schemas.logicaltypes.FixedBytes;
 import org.apache.beam.sdk.schemas.logicaltypes.FixedString;
 import org.apache.beam.sdk.schemas.logicaltypes.SqlTypes;
+import org.apache.beam.sdk.schemas.logicaltypes.Timestamp;
 import org.apache.beam.sdk.schemas.logicaltypes.VariableBytes;
 import org.apache.beam.sdk.schemas.logicaltypes.VariableString;
 import org.apache.beam.sdk.testing.PAssert;
@@ -797,4 +798,35 @@ public class BeamComplexTypeTest {
     assertEquals(inputRow.getSchema(), outputRow.getSchema());
     pipeline.run().waitUntilFinish(Duration.standardMinutes(1));
   }
+
+  @Test
+  public void testSqlTimestampLogicalType() {
+    // Calcite TIMESTAMP is millis-based; SQL projection of Timestamp.MICROS 
uses
+    // FieldType.DATETIME.
+    Schema inputSchema =
+        Schema.builder()
+            .addField("ts", FieldType.logicalType(Timestamp.MICROS))
+            .addNullableField("nullable_ts", 
FieldType.logicalType(Timestamp.MICROS))
+            .build();
+
+    java.time.Instant ts = java.time.Instant.parse("2025-07-31T20:17:40.123Z");
+    Row inputRow = Row.withSchema(inputSchema).addValues(ts, null).build();
+
+    PCollection<Row> outputRow =
+        pipeline
+            .apply(Create.of(inputRow))
+            .setRowSchema(inputSchema)
+            .apply(SqlTransform.query("SELECT ts, nullable_ts FROM 
PCOLLECTION"));
+
+    Schema outputSchema =
+        Schema.builder()
+            .addDateTimeField("ts")
+            .addNullableField("nullable_ts", FieldType.DATETIME)
+            .build();
+    Row expectedRow =
+        Row.withSchema(outputSchema).addValues(new Instant(ts.toEpochMilli()), 
null).build();
+
+    PAssert.that(outputRow).containsInAnyOrder(expectedRow);
+    pipeline.run().waitUntilFinish(Duration.standardMinutes(2));
+  }
 }
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/rel/BeamCalcRelTest.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/rel/BeamCalcRelTest.java
index 8aeb77fc049..019832ffe96 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/rel/BeamCalcRelTest.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/impl/rel/BeamCalcRelTest.java
@@ -17,7 +17,11 @@
  */
 package org.apache.beam.sdk.extensions.sql.impl.rel;
 
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertThrows;
+
 import java.math.BigDecimal;
+import java.time.Instant;
 import org.apache.beam.sdk.Pipeline;
 import org.apache.beam.sdk.extensions.sql.impl.BeamTableStatistics;
 import org.apache.beam.sdk.extensions.sql.impl.planner.BeamRelMetadataQuery;
@@ -250,4 +254,20 @@ public class BeamCalcRelTest extends BaseRelTest {
 
     pipeline.run().waitUntilFinish();
   }
+
+  @Test
+  public void testTimestampToCalciteMillisAcceptsMillisecondPrecision() {
+    Instant instant = Instant.parse("2025-07-31T20:17:40.123Z");
+    assertEquals(instant.toEpochMilli(), 
BeamCalcRel.timestampToCalciteMillis(instant));
+  }
+
+  @Test
+  public void testTimestampToCalciteMillisRejectsSubMillisecondPrecision() {
+    Instant instant = Instant.parse("2025-07-31T20:17:40.123456Z");
+    UnsupportedOperationException thrown =
+        assertThrows(
+            UnsupportedOperationException.class,
+            () -> BeamCalcRel.timestampToCalciteMillis(instant));
+    Assert.assertTrue(thrown.getMessage().contains("sub-millisecond"));
+  }
 }

Reply via email to