From fe0d00c8b17bbd037108d9aaca2e0b33652530d8 Mon Sep 17 00:00:00 2001 From: aibrahiim Date: Mon, 3 Aug 2026 17:30:14 +0300 Subject: [PATCH 1/4] add Timestamp.MICROS for iceberg timestamptz --- .../trigger_files/beam_PostCommit_SQL.json | 2 +- .github/trigger_files/beam_PreCommit_SQL.json | 2 +- .../iceberg/BeamSqlCliIcebergTest.java | 7 +++++-- .../provider/iceberg/IcebergReadWriteIT.java | 9 +++++--- .../extensions/sql/impl/rel/BeamCalcRel.java | 21 +++++++++++++++++++ .../sql/impl/utils/CalciteUtils.java | 6 +++++- 6 files changed, 39 insertions(+), 8 deletions(-) diff --git a/.github/trigger_files/beam_PostCommit_SQL.json b/.github/trigger_files/beam_PostCommit_SQL.json index 5df3841d2363..5ac8a7f3f6ee 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": 4 } diff --git a/.github/trigger_files/beam_PreCommit_SQL.json b/.github/trigger_files/beam_PreCommit_SQL.json index 07d1fb889961..5abe02fc09c7 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": 1 } 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 9ac96652d340..f5ec529a5c00 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 @@ -40,7 +40,6 @@ import org.apache.beam.sdk.values.Row; import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.runtime.CalciteContextException; import org.checkerframework.checker.nullness.qual.Nullable; -import org.joda.time.DateTime; import org.junit.Before; import org.junit.ClassRule; import org.junit.Rule; @@ -229,7 +228,11 @@ public void testCrossCatalogTableWriteAndRead() throws IOException { PAssert.that(output) .containsInAnyOrder( Row.withSchema(expectedSchema) - .addValues(2147483647, true, DateTime.parse("2025-07-31T20:17:40.123Z"), "varchar") + .addValues( + 2147483647, + true, + java.time.Instant.parse("2025-07-31T20:17:40.123Z"), + "varchar") .build()); p3.run().waitUntilFinish(); assertEquals("catalog_1", catalogManager.currentCatalog().name()); 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 417db09a2210..9220d1e49798 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.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; @@ -55,6 +54,8 @@ import org.apache.beam.sdk.io.iceberg.IcebergUtils; import org.apache.beam.sdk.options.PipelineOptionsFactory; import org.apache.beam.sdk.schemas.Schema; +import org.apache.beam.sdk.schemas.Schema.FieldType; +import org.apache.beam.sdk.schemas.logicaltypes.Timestamp; import org.apache.beam.sdk.testing.PAssert; import org.apache.beam.sdk.testing.TestPipeline; import org.apache.beam.sdk.values.PCollection; @@ -200,8 +201,10 @@ public void runSqlWriteAndRead(boolean withPartitionFields) 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())); + Schema fromIceberg = IcebergUtils.icebergSchemaToBeamSchema(icebergTable.schema()); + assertEquals( + FieldType.logicalType(Timestamp.MICROS).withNullable(true), + fromIceberg.getField("c_timestamp").getType()); // 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 fad96abb29a5..25655ae087d9 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 @@ -439,6 +439,12 @@ static Object toBeamObject(Object value, FieldType fieldType, boolean verifyValu 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 +597,15 @@ private static Expression getBeamField( 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 +699,12 @@ private static Expression toCalciteValue( 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( + Expressions.convert_(value, java.time.Instant.class), "toEpochMilli")); } 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 d55c227e7b45..f8f37ca9fa5e 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.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 static boolean isDateTimeType(FieldType fieldType) { 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; } @@ -222,6 +224,8 @@ public static SqlTypeName toSqlTypeName(FieldType type) { 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" From 5b71eb5761e842ee03d5bfb0b57b2a42d9a9f9bb Mon Sep 17 00:00:00 2001 From: aibrahiim Date: Mon, 3 Aug 2026 21:21:09 +0300 Subject: [PATCH 2/4] fix schema assert --- .github/trigger_files/beam_PostCommit_SQL.json | 2 +- .github/trigger_files/beam_PreCommit_SQL.json | 2 +- .../provider/iceberg/BeamSqlCliIcebergTest.java | 13 ++++--------- 3 files changed, 6 insertions(+), 11 deletions(-) diff --git a/.github/trigger_files/beam_PostCommit_SQL.json b/.github/trigger_files/beam_PostCommit_SQL.json index 5ac8a7f3f6ee..ae43f2754d6e 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": 4 + "modification": 5 } diff --git a/.github/trigger_files/beam_PreCommit_SQL.json b/.github/trigger_files/beam_PreCommit_SQL.json index 5abe02fc09c7..3a009261f4f9 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": 1 + "modification": 2 } 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 f5ec529a5c00..567c3bbc7fae 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; @@ -40,6 +39,7 @@ import org.apache.beam.sdk.values.Row; import org.apache.beam.vendor.calcite.v1_40_0.org.apache.calcite.runtime.CalciteContextException; import org.checkerframework.checker.nullness.qual.Nullable; +import org.joda.time.DateTime; import org.junit.Before; import org.junit.ClassRule; import org.junit.Rule; @@ -222,17 +222,12 @@ public void testCrossCatalogTableWriteAndRead() throws IOException { PCollection 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) - .addValues( - 2147483647, - true, - java.time.Instant.parse("2025-07-31T20:17:40.123Z"), - "varchar") + .addValues(2147483647, true, DateTime.parse("2025-07-31T20:17:40.123Z"), "varchar") .build()); p3.run().waitUntilFinish(); assertEquals("catalog_1", catalogManager.currentCatalog().name()); From 1def02993feac4d2b0d162e38d9284c7cec3f106 Mon Sep 17 00:00:00 2001 From: aibrahiim Date: Tue, 4 Aug 2026 12:21:57 +0300 Subject: [PATCH 3/4] address comments --- .../trigger_files/beam_PostCommit_SQL.json | 2 +- .github/trigger_files/beam_PreCommit_SQL.json | 2 +- .../provider/iceberg/IcebergReadWriteIT.java | 6 ---- .../extensions/sql/impl/rel/BeamCalcRel.java | 21 +++++++++++- .../sql/impl/utils/CalciteUtils.java | 2 ++ .../extensions/sql/BeamComplexTypeTest.java | 32 +++++++++++++++++++ .../sql/impl/rel/BeamCalcRelTest.java | 20 ++++++++++++ 7 files changed, 76 insertions(+), 9 deletions(-) diff --git a/.github/trigger_files/beam_PostCommit_SQL.json b/.github/trigger_files/beam_PostCommit_SQL.json index ae43f2754d6e..1b6aa099172e 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": 5 + "modification": 6 } diff --git a/.github/trigger_files/beam_PreCommit_SQL.json b/.github/trigger_files/beam_PreCommit_SQL.json index 3a009261f4f9..ab4daeae2349 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": 2 + "modification": 3 } 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 9220d1e49798..3a791c6fe88c 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 @@ -54,8 +54,6 @@ import org.apache.beam.sdk.io.iceberg.IcebergUtils; import org.apache.beam.sdk.options.PipelineOptionsFactory; import org.apache.beam.sdk.schemas.Schema; -import org.apache.beam.sdk.schemas.Schema.FieldType; -import org.apache.beam.sdk.schemas.logicaltypes.Timestamp; import org.apache.beam.sdk.testing.PAssert; import org.apache.beam.sdk.testing.TestPipeline; import org.apache.beam.sdk.values.PCollection; @@ -201,10 +199,6 @@ public void runSqlWriteAndRead(boolean withPartitionFields) assertEquals("my_catalog." + tableIdentifier, icebergTable.name()); assertTrue(icebergTable.location().startsWith(warehouse)); assertEquals(expectedSpec, icebergTable.spec()); - Schema fromIceberg = IcebergUtils.icebergSchemaToBeamSchema(icebergTable.schema()); - assertEquals( - FieldType.logicalType(Timestamp.MICROS).withNullable(true), - fromIceberg.getField("c_timestamp").getType()); // 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 25655ae087d9..b9525bd07dd0 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 rows = new TupleTag() {}; private static final TupleTag errors = new TupleTag() {}; + /** + * 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); } @@ -704,7 +721,9 @@ private static Expression toCalciteValue( return nullOr( value, Expressions.call( - Expressions.convert_(value, java.time.Instant.class), "toEpochMilli")); + 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 f8f37ca9fa5e..2627f1c0f868 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 @@ -116,6 +116,8 @@ public static boolean isStringType(FieldType fieldType) { 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); 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 5ef081b92c3f..39599b7f148c 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.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 void testUnknownLogicalType() { assertEquals(inputRow.getSchema(), outputRow.getSchema()); pipeline.run().waitUntilFinish(Duration.standardMinutes(1)); } + + @Test + public void testSqlTimestampMicrosLogicalType() { + // 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 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 8aeb77fc0490..019832ffe964 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 void testNoFieldAccess() throws IllegalAccessException { 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")); + } } From 45c845062f11316c2c07448cd68438b6e2ff8767 Mon Sep 17 00:00:00 2001 From: aibrahiim Date: Wed, 5 Aug 2026 10:38:05 +0300 Subject: [PATCH 4/4] rename test --- .../org/apache/beam/sdk/extensions/sql/BeamComplexTypeTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 39599b7f148c..062651161f59 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 @@ -800,7 +800,7 @@ public void testUnknownLogicalType() { } @Test - public void testSqlTimestampMicrosLogicalType() { + public void testSqlTimestampLogicalType() { // Calcite TIMESTAMP is millis-based; SQL projection of Timestamp.MICROS uses // FieldType.DATETIME. Schema inputSchema =