From 4fdb155f62b2e38d6af7ac0150d61d0cb84a5c5d Mon Sep 17 00:00:00 2001 From: Zihan Dai Date: Wed, 19 Aug 2026 00:08:28 +1000 Subject: [PATCH] Portable Time type (Java + model changes) Adds TIME = 10 to the LogicalTypes enum with URN beam:logical_type:time:v1, registers Time in SchemaTranslation's standard logical types, and derives Time.IDENTIFIER from the enum rather than hardcoding the string. Mirrors the portable Date type added in #38077. As there, the Python side follows in a separate PR so the URN reaches a snapshot container before the cross-language tests need it. --- .../beam/model/pipeline/v1/schema.proto | 7 ++++ .../beam/sdk/schemas/SchemaTranslation.java | 2 + .../beam/sdk/schemas/logicaltypes/Time.java | 20 +++++----- .../sdk/schemas/SchemaTranslationTest.java | 40 +++++++++++++++++++ 4 files changed, 59 insertions(+), 10 deletions(-) diff --git a/model/pipeline/src/main/proto/org/apache/beam/model/pipeline/v1/schema.proto b/model/pipeline/src/main/proto/org/apache/beam/model/pipeline/v1/schema.proto index c114c35f2d90..97d68440a797 100644 --- a/model/pipeline/src/main/proto/org/apache/beam/model/pipeline/v1/schema.proto +++ b/model/pipeline/src/main/proto/org/apache/beam/model/pipeline/v1/schema.proto @@ -247,6 +247,13 @@ message LogicalTypes { // e.g. -1.5s at precision 6 is {seconds: -2, subseconds: 500000}. TIMESTAMP = 9 [(org.apache.beam.model.pipeline.v1.beam_urn) = "beam:logical_type:timestamp:v1"]; + + // A URN for Time type + // - Representation type: INT64 + // - A time without a timezone, represented by the number of + // nanoseconds since midnight. + TIME = 10 [(org.apache.beam.model.pipeline.v1.beam_urn) = + "beam:logical_type:time:v1"]; } } diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/SchemaTranslation.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/SchemaTranslation.java index 8122704e6436..4589ae918fcb 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/SchemaTranslation.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/SchemaTranslation.java @@ -50,6 +50,7 @@ import org.apache.beam.sdk.schemas.logicaltypes.MicrosInstant; import org.apache.beam.sdk.schemas.logicaltypes.PythonCallable; import org.apache.beam.sdk.schemas.logicaltypes.SchemaLogicalType; +import org.apache.beam.sdk.schemas.logicaltypes.Time; import org.apache.beam.sdk.schemas.logicaltypes.Timestamp; import org.apache.beam.sdk.schemas.logicaltypes.UnknownLogicalType; import org.apache.beam.sdk.schemas.logicaltypes.VariableBytes; @@ -116,6 +117,7 @@ private static String getLogicalTypeUrn(String identifier) { .put(FixedString.IDENTIFIER, FixedString.class) .put(VariableString.IDENTIFIER, VariableString.class) .put(Date.IDENTIFIER, Date.class) + .put(Time.IDENTIFIER, Time.class) .put(Timestamp.IDENTIFIER, Timestamp.class) .build(); diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/logicaltypes/Time.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/logicaltypes/Time.java index 04f307063e77..809d52601660 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/logicaltypes/Time.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/logicaltypes/Time.java @@ -18,7 +18,10 @@ package org.apache.beam.sdk.schemas.logicaltypes; import java.time.LocalTime; +import org.apache.beam.model.pipeline.v1.RunnerApi; +import org.apache.beam.model.pipeline.v1.SchemaApi; import org.apache.beam.sdk.schemas.Schema; +import org.checkerframework.checker.nullness.qual.Nullable; /** * A time without a time-zone. @@ -30,23 +33,20 @@ * of time in nanoseconds. */ public class Time implements Schema.LogicalType { - public static final String IDENTIFIER = "beam:logical_type:time:v1"; + public static final String IDENTIFIER = + SchemaApi.LogicalTypes.Enum.TIME + .getValueDescriptor() + .getOptions() + .getExtension(RunnerApi.beamUrn); @Override public String getIdentifier() { return IDENTIFIER; } - // unused @Override - public Schema.FieldType getArgumentType() { - return Schema.FieldType.STRING; - } - - // unused - @Override - public String getArgument() { - return ""; + public Schema.@Nullable FieldType getArgumentType() { + return null; } @Override diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/SchemaTranslationTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/SchemaTranslationTest.java index 925b1e02b786..6ecf0aad04fe 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/SchemaTranslationTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/SchemaTranslationTest.java @@ -27,6 +27,7 @@ import java.nio.charset.StandardCharsets; import java.time.LocalDate; import java.time.LocalDateTime; +import java.time.LocalTime; import java.util.ArrayList; import java.util.HashMap; import java.util.List; @@ -41,6 +42,7 @@ import org.apache.beam.model.pipeline.v1.SchemaApi.LogicalType; import org.apache.beam.sdk.schemas.Schema.Field; import org.apache.beam.sdk.schemas.Schema.FieldType; +import org.apache.beam.sdk.schemas.logicaltypes.Date; import org.apache.beam.sdk.schemas.logicaltypes.DateTime; import org.apache.beam.sdk.schemas.logicaltypes.FixedBytes; import org.apache.beam.sdk.schemas.logicaltypes.FixedPrecisionNumeric; @@ -51,6 +53,7 @@ import org.apache.beam.sdk.schemas.logicaltypes.PythonCallable; import org.apache.beam.sdk.schemas.logicaltypes.SchemaLogicalType; import org.apache.beam.sdk.schemas.logicaltypes.SqlTypes; +import org.apache.beam.sdk.schemas.logicaltypes.Time; import org.apache.beam.sdk.schemas.logicaltypes.Timestamp; import org.apache.beam.sdk.schemas.logicaltypes.UnknownLogicalType; import org.apache.beam.sdk.schemas.logicaltypes.VariableBytes; @@ -145,6 +148,7 @@ public static Iterable data() { .add(Schema.of(Field.of("fixed_bytes", FieldType.logicalType(FixedBytes.of(24))))) .add(Schema.of(Field.of("micros_instant", FieldType.logicalType(new MicrosInstant())))) .add(Schema.of(Field.of("date", FieldType.logicalType(SqlTypes.DATE)))) + .add(Schema.of(Field.of("time", FieldType.logicalType(SqlTypes.TIME)))) .add(Schema.of(Field.of("python_callable", FieldType.logicalType(new PythonCallable())))) .add( Schema.of( @@ -392,6 +396,7 @@ public static Iterable data() { .add(simpleRow(FieldType.logicalType(new PortableNullArgLogicalType()), "str")) .add(simpleRow(FieldType.logicalType(new DateTime()), LocalDateTime.of(2000, 1, 3, 3, 1))) .add(simpleRow(FieldType.logicalType(SqlTypes.DATE), LocalDate.of(2000, 1, 3))) + .add(simpleRow(FieldType.logicalType(SqlTypes.TIME), LocalTime.of(3, 1, 2, 3000))) .add(simpleNullRow(FieldType.STRING)) .add(simpleNullRow(FieldType.INT32)) .add(simpleNullRow(FieldType.map(FieldType.STRING, FieldType.INT32))) @@ -444,6 +449,41 @@ public void typeInfoNotSet() { } } + /** + * A portable logical type has to be recoverable from its URN alone, because that is all a schema + * coming from another SDK carries. {@link LogicalTypesTest#testLogicalTypeFromToProtoCorrectly} + * cannot check this: it branches on {@code STANDARD_LOGICAL_TYPES} itself, so it passes whether + * or not the type is registered there. + */ + @RunWith(JUnit4.class) + public static class PortableLogicalTypeFromUrnTest { + + @Test + public void timeIsRecoveredFromItsUrnAlone() { + FieldType fieldType = FieldType.logicalType(SqlTypes.TIME); + + // serializeLogicalType = false, so the proto carries the URN and no Java payload + SchemaApi.FieldType proto = SchemaTranslation.fieldTypeToProto(fieldType, false, false); + assertThat(proto.getLogicalType().getUrn(), equalTo("beam:logical_type:time:v1")); + assertThat(proto.getLogicalType().getPayload().size(), equalTo(0)); + + Schema.FieldType translated = SchemaTranslation.fieldTypeFromProto(proto); + assertThat(translated.getLogicalType().getClass(), equalTo(Time.class)); + assertThat(translated.getLogicalType().getBaseType(), equalTo(FieldType.INT64)); + } + + @Test + public void dateIsRecoveredFromItsUrnAlone() { + FieldType fieldType = FieldType.logicalType(SqlTypes.DATE); + + SchemaApi.FieldType proto = SchemaTranslation.fieldTypeToProto(fieldType, false, false); + assertThat(proto.getLogicalType().getUrn(), equalTo("beam:logical_type:date:v1")); + + Schema.FieldType translated = SchemaTranslation.fieldTypeFromProto(proto); + assertThat(translated.getLogicalType().getClass(), equalTo(Date.class)); + } + } + /** Test schema translation of logical types. */ @RunWith(Parameterized.class) public static class LogicalTypesTest {