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"));
+ }
}