This is an automated email from the ASF dual-hosted git repository.
snuyanzin pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git
The following commit(s) were added to refs/heads/master by this push:
new 37cf4088c1f [FLINK-40222][table] Fix precision for `INTERVAL` literals
37cf4088c1f is described below
commit 37cf4088c1f4c309cbfdfec8cf0d0ec25b1e97a2
Author: Timo Theusner <[email protected]>
AuthorDate: Wed Jul 29 07:04:19 2026 +0200
[FLINK-40222][table] Fix precision for `INTERVAL` literals
---
.../table/expressions/ValueLiteralExpression.java | 62 ++++++++--
.../table/types/utils/ValueDataTypeConverter.java | 12 +-
.../flink/table/expressions/ExpressionTest.java | 126 +++++++++++++++++++++
.../table/types/ValueDataTypeConverterTest.java | 4 +
.../table/planner/calcite/FlinkTypeFactory.java | 3 +-
.../table/api/QueryOperationTestPrograms.java | 6 +-
.../planner/calcite/FlinkTypeFactoryTest.java | 48 ++++++++
.../LiteralExpressionsSerializationITCase.java | 49 ++++++--
8 files changed, 280 insertions(+), 30 deletions(-)
diff --git
a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/expressions/ValueLiteralExpression.java
b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/expressions/ValueLiteralExpression.java
index 36adeb8a2d6..72a2534b64f 100644
---
a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/expressions/ValueLiteralExpression.java
+++
b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/expressions/ValueLiteralExpression.java
@@ -23,11 +23,13 @@ import org.apache.flink.table.api.TableException;
import org.apache.flink.table.api.ValidationException;
import org.apache.flink.table.types.DataType;
import org.apache.flink.table.types.inference.CallContext;
+import org.apache.flink.table.types.logical.DayTimeIntervalType;
import org.apache.flink.table.types.logical.DecimalType;
import org.apache.flink.table.types.logical.LocalZonedTimestampType;
import org.apache.flink.table.types.logical.LogicalType;
import org.apache.flink.table.types.logical.LogicalTypeFamily;
import org.apache.flink.table.types.logical.LogicalTypeRoot;
+import org.apache.flink.table.types.logical.YearMonthIntervalType;
import org.apache.flink.table.types.utils.ValueDataTypeConverter;
import org.apache.flink.table.utils.DateTimeUtils;
import org.apache.flink.table.utils.EncodingUtils;
@@ -288,18 +290,11 @@ public final class ValueLiteralExpression implements
ResolvedExpression {
getValueAs(Instant.class).get(),
((LocalZonedTimestampType)
dataType.getLogicalType()).getPrecision());
case INTERVAL_YEAR_MONTH:
- final Period period =
getValueAs(Period.class).get().normalized();
- return String.format(
- "INTERVAL '%d-%d' YEAR TO MONTH", period.getYears(),
period.getMonths());
+ return formatYearMonthIntervalLiteral(
+ (YearMonthIntervalType) logicalType,
getValueAs(Period.class).get());
case INTERVAL_DAY_TIME:
- final Duration duration = getValueAs(Duration.class).get();
- return String.format(
- "INTERVAL '%d %02d:%02d:%02d.%d' DAY TO SECOND(3)",
- duration.toDays(),
- duration.toHours() % 24,
- duration.toMinutes() % 60,
- duration.getSeconds() % 60,
- duration.getNano() / 1_000_000);
+ return formatDayTimeIntervalLiteral(
+ (DayTimeIntervalType) logicalType,
getValueAs(Duration.class).get());
case DESCRIPTOR:
final ColumnList columnList =
getValueAs(ColumnList.class).get();
if (!columnList.getDataTypes().isEmpty()) {
@@ -411,6 +406,51 @@ public final class ValueLiteralExpression implements
ResolvedExpression {
}
}
+ /**
+ * Formats a year-month interval value as a SQL literal, including
explicit precision on the
+ * {@code YEAR} field.
+ */
+ private static String formatYearMonthIntervalLiteral(
+ YearMonthIntervalType type, Period period) {
+ final long totalMonths = period.toTotalMonths();
+ final long years = totalMonths / 12;
+ final int months = (int) (totalMonths % 12);
+ final int yearPrecision = type.getYearPrecision();
+ return String.format("INTERVAL '%d-%d' YEAR(%d) TO MONTH", years,
months, yearPrecision);
+ }
+
+ /**
+ * Formats a day-time interval value as a SQL literal, including explicit
precision on the
+ * {@code DAY} and {@code SECOND} field.
+ */
+ private static String formatDayTimeIntervalLiteral(
+ DayTimeIntervalType type, Duration duration) {
+ final long days = duration.toDays();
+ final int hours = duration.toHoursPart();
+ final int minutes = duration.toMinutesPart();
+ final int seconds = duration.toSecondsPart();
+ final int dayPrecision = type.getDayPrecision();
+ final int fractionalPrecision = type.getFractionalPrecision();
+ final String fraction = formatFractionalSeconds(duration.getNano(),
fractionalPrecision);
+
+ return String.format(
+ "INTERVAL '%d %02d:%02d:%02d%s' DAY(%d) TO SECOND(%d)",
+ days, hours, minutes, seconds, fraction, dayPrecision,
fractionalPrecision);
+ }
+
+ /**
+ * Formats the fractional-seconds suffix (including the leading {@code .})
from a raw nanosecond
+ * value to {@code fractionalPrecision} digits. Returns an empty string
when {@code
+ * fractionalPrecision} is 0.
+ */
+ private static String formatFractionalSeconds(int nanos, int
fractionalPrecision) {
+ if (fractionalPrecision == 0) {
+ return "";
+ }
+ final String nanosString = String.format("%09d", nanos);
+ return "." + nanosString.substring(0, fractionalPrecision);
+ }
+
/** Supports (nested) arrays and makes string values more explicit. */
private static String stringifyValue(Object value) {
if (value instanceof String[]) {
diff --git
a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/utils/ValueDataTypeConverter.java
b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/utils/ValueDataTypeConverter.java
index d27b87efae4..347f562d6d3 100644
---
a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/utils/ValueDataTypeConverter.java
+++
b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/utils/ValueDataTypeConverter.java
@@ -83,8 +83,7 @@ public final class ValueDataTypeConverter {
convertedDataType =
convertToLocalZonedTimestampType(((java.time.Instant)
value).getNano());
} else if (value instanceof java.time.Period) {
- convertedDataType =
- convertToYearMonthIntervalType(((java.time.Period)
value).getYears());
+ convertedDataType =
convertToYearMonthIntervalType((java.time.Period) value);
} else if (value instanceof java.time.Duration) {
final java.time.Duration duration = (java.time.Duration) value;
convertedDataType =
convertToDayTimeIntervalType(duration.toDays(), duration.getNano());
@@ -156,7 +155,8 @@ public final class ValueDataTypeConverter {
return
DataTypes.TIMESTAMP_WITH_LOCAL_TIME_ZONE(fractionalSecondPrecision(nanos));
}
- private static DataType convertToYearMonthIntervalType(int years) {
+ private static DataType convertToYearMonthIntervalType(java.time.Period
period) {
+ final long years = period.toTotalMonths() / 12;
return DataTypes.INTERVAL(DataTypes.YEAR(yearPrecision(years)),
DataTypes.MONTH());
}
@@ -217,12 +217,12 @@ public final class ValueDataTypeConverter {
return String.format("%09d", nanos).replaceAll("0+$", "").length();
}
- private static int yearPrecision(int years) {
- return String.valueOf(years).length();
+ private static int yearPrecision(long years) {
+ return String.valueOf(Math.abs(years)).length();
}
private static int dayPrecision(long days) {
- return String.valueOf(days).length();
+ return String.valueOf(Math.abs(days)).length();
}
private ValueDataTypeConverter() {
diff --git
a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/expressions/ExpressionTest.java
b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/expressions/ExpressionTest.java
index f39d36fba58..fca9d0a8f66 100644
---
a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/expressions/ExpressionTest.java
+++
b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/expressions/ExpressionTest.java
@@ -22,6 +22,7 @@ import org.apache.flink.table.api.DataTypes;
import org.apache.flink.table.api.ValidationException;
import org.apache.flink.table.functions.ScalarFunction;
import org.apache.flink.table.functions.ScalarFunctionDefinition;
+import org.apache.flink.table.types.DataType;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
@@ -33,6 +34,7 @@ import java.nio.charset.StandardCharsets;
import java.sql.Date;
import java.sql.Time;
import java.sql.Timestamp;
+import java.time.Duration;
import java.time.Instant;
import java.time.LocalDate;
import java.time.LocalDateTime;
@@ -293,6 +295,16 @@ class ExpressionTest {
.isEqualTo(expected);
}
+ @ParameterizedTest(name = "{0}")
+ @MethodSource("intervalLiteralTestCases")
+ void testIntervalAsSerializableString(
+ String caseName, Object value, DataType dataType, String expected)
{
+ assertThat(
+ new ValueLiteralExpression(value, dataType)
+
.asSerializableString(DefaultSqlFactory.INSTANCE))
+ .isEqualTo(expected);
+ }
+
//
--------------------------------------------------------------------------------------------
private static Expression createExpressionTree(Integer nestedValue) {
@@ -365,4 +377,118 @@ class ExpressionTest {
9,
"TO_TIMESTAMP_LTZ(-9223372036000000000, 9)"));
}
+
+ private static Stream<Arguments> intervalLiteralTestCases() {
+ return Stream.of(
+ Arguments.of(
+ "YEAR_TO_MONTH maximum",
+ Period.of(9999, 11, 0),
+ DataTypes.INTERVAL(DataTypes.YEAR(4),
DataTypes.MONTH()).notNull(),
+ "INTERVAL '9999-11' YEAR(4) TO MONTH"),
+ Arguments.of(
+ "YEAR_TO_MONTH default precision",
+ Period.ofMonths(470),
+ DataTypes.INTERVAL(DataTypes.YEAR(),
DataTypes.MONTH()).notNull(),
+ "INTERVAL '39-2' YEAR(2) TO MONTH"),
+ Arguments.of(
+ "YEAR with months not dropped",
+ Period.of(5, 3, 0),
+ DataTypes.INTERVAL(DataTypes.YEAR()).notNull(),
+ "INTERVAL '5-3' YEAR(2) TO MONTH"),
+ Arguments.of(
+ "YEAR only whole years",
+ Period.ofYears(120),
+ DataTypes.INTERVAL(DataTypes.YEAR(3)).notNull(),
+ "INTERVAL '120-0' YEAR(3) TO MONTH"),
+ Arguments.of(
+ "MONTH only",
+ Period.ofMonths(50),
+ DataTypes.INTERVAL(DataTypes.MONTH()).notNull(),
+ "INTERVAL '4-2' YEAR(2) TO MONTH"),
+ Arguments.of(
+ "DAY(3) three-digit value",
+ Duration.ofDays(100),
+ DataTypes.INTERVAL(DataTypes.DAY(3)).notNull(),
+ "INTERVAL '100 00:00:00.000000' DAY(3) TO SECOND(6)"),
+ Arguments.of(
+ "DAY_TO_HOUR",
+ Duration.ofDays(5).plusHours(7),
+ DataTypes.INTERVAL(DataTypes.DAY(2),
DataTypes.HOUR()).notNull(),
+ "INTERVAL '5 07:00:00.000000' DAY(2) TO SECOND(6)"),
+ Arguments.of(
+ "DAY_TO_MINUTE",
+ Duration.ofDays(5).plusHours(7).plusMinutes(8),
+ DataTypes.INTERVAL(DataTypes.DAY(2),
DataTypes.MINUTE()).notNull(),
+ "INTERVAL '5 07:08:00.000000' DAY(2) TO SECOND(6)"),
+ Arguments.of(
+ "DAY_TO_SECOND fractional padding",
+ Duration.ofDays(1).plusMillis(5),
+ DataTypes.INTERVAL(DataTypes.DAY(2),
DataTypes.SECOND(3)).notNull(),
+ "INTERVAL '1 00:00:00.005' DAY(2) TO SECOND(3)"),
+ Arguments.of(
+ "DAY_TO_SECOND nanosecond precision (string only;
sub-ms not round-trippable)",
+
Duration.ofDays(99).plusSeconds(34).plusNanos(999_000_001),
+ DataTypes.INTERVAL(DataTypes.DAY(2),
DataTypes.SECOND(9)).notNull(),
+ "INTERVAL '99 00:00:34.999000001' DAY(2) TO
SECOND(9)"),
+ Arguments.of(
+ "DAY_TO_SECOND zero fractional precision",
+ Duration.ofDays(1).plusSeconds(30),
+ DataTypes.INTERVAL(DataTypes.DAY(2),
DataTypes.SECOND(0)).notNull(),
+ "INTERVAL '1 00:00:30' DAY(2) TO SECOND(0)"),
+ Arguments.of(
+ "HOUR only preserves whole value",
+ Duration.ofHours(30),
+ DataTypes.INTERVAL(DataTypes.HOUR()).notNull(),
+ "INTERVAL '1 06:00:00.000000' DAY(2) TO SECOND(6)"),
+ Arguments.of(
+ "HOUR with minutes not dropped",
+ Duration.ofHours(5).plusMinutes(30),
+ DataTypes.INTERVAL(DataTypes.HOUR()).notNull(),
+ "INTERVAL '0 05:30:00.000000' DAY(2) TO SECOND(6)"),
+ Arguments.of(
+ "HOUR_TO_MINUTE",
+ Duration.ofHours(30).plusMinutes(15),
+ DataTypes.INTERVAL(DataTypes.HOUR(),
DataTypes.MINUTE()).notNull(),
+ "INTERVAL '1 06:15:00.000000' DAY(2) TO SECOND(6)"),
+ Arguments.of(
+ "HOUR_TO_SECOND",
+
Duration.ofHours(30).plusMinutes(15).plusSeconds(20).plusNanos(500_000_000),
+ DataTypes.INTERVAL(DataTypes.HOUR(),
DataTypes.SECOND(6)).notNull(),
+ "INTERVAL '1 06:15:20.500000' DAY(2) TO SECOND(6)"),
+ Arguments.of(
+ "MINUTE only",
+ Duration.ofMinutes(45),
+ DataTypes.INTERVAL(DataTypes.MINUTE()).notNull(),
+ "INTERVAL '0 00:45:00.000000' DAY(2) TO SECOND(6)"),
+ Arguments.of(
+ "MINUTE >= 100 (0 days, fits DAY(2))",
+ Duration.ofMinutes(100),
+ DataTypes.INTERVAL(DataTypes.MINUTE()).notNull(),
+ "INTERVAL '0 01:40:00.000000' DAY(2) TO SECOND(6)"),
+ Arguments.of(
+ "MINUTE_TO_SECOND",
+ Duration.ofMinutes(45).plusSeconds(45).plusMillis(120),
+ DataTypes.INTERVAL(DataTypes.MINUTE(),
DataTypes.SECOND(2)).notNull(),
+ "INTERVAL '0 00:45:45.12' DAY(2) TO SECOND(2)"),
+ Arguments.of(
+ "SECOND standalone",
+ Duration.ofSeconds(45).plusMillis(750),
+ DataTypes.INTERVAL(DataTypes.SECOND(4)).notNull(),
+ "INTERVAL '0 00:00:45.7500' DAY(2) TO SECOND(4)"),
+ Arguments.of(
+ "SECOND >= 100 (0 days, fits DAY(2))",
+ Duration.ofSeconds(150),
+ DataTypes.INTERVAL(DataTypes.SECOND(3)).notNull(),
+ "INTERVAL '0 00:02:30.000' DAY(2) TO SECOND(3)"),
+ Arguments.of(
+ "Large day value under its natural DAY(3) type (=
10000 hours)",
+ Duration.ofDays(416).plusHours(16),
+ DataTypes.INTERVAL(DataTypes.DAY(3),
DataTypes.HOUR()).notNull(),
+ "INTERVAL '416 16:00:00.000000' DAY(3) TO SECOND(6)"),
+ Arguments.of(
+ "DAY maximum precision",
+ Duration.ofDays(999_999),
+ DataTypes.INTERVAL(DataTypes.DAY(6)).notNull(),
+ "INTERVAL '999999 00:00:00.000000' DAY(6) TO
SECOND(6)"));
+ }
}
diff --git
a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/ValueDataTypeConverterTest.java
b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/ValueDataTypeConverterTest.java
index 14c4a7466f5..9f767c44afd 100644
---
a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/ValueDataTypeConverterTest.java
+++
b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/ValueDataTypeConverterTest.java
@@ -86,6 +86,10 @@ class ValueDataTypeConverterTest {
Period.ofYears(1000),
DataTypes.INTERVAL(DataTypes.YEAR(4),
DataTypes.MONTH())
.bridgedTo(Period.class)),
+ of(
+ Period.ofMonths(470),
+ DataTypes.INTERVAL(DataTypes.YEAR(2),
DataTypes.MONTH())
+ .bridgedTo(Period.class)),
of(
Duration.ofMillis(1100),
DataTypes.INTERVAL(DataTypes.DAY(1),
DataTypes.SECOND(1))
diff --git
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/calcite/FlinkTypeFactory.java
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/calcite/FlinkTypeFactory.java
index bb99230775c..f5bd64adeef 100644
---
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/calcite/FlinkTypeFactory.java
+++
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/calcite/FlinkTypeFactory.java
@@ -42,6 +42,7 @@ import org.apache.flink.table.types.logical.BitmapType;
import org.apache.flink.table.types.logical.BooleanType;
import org.apache.flink.table.types.logical.CharType;
import org.apache.flink.table.types.logical.DateType;
+import org.apache.flink.table.types.logical.DayTimeIntervalType;
import org.apache.flink.table.types.logical.DecimalType;
import org.apache.flink.table.types.logical.DescriptorType;
import org.apache.flink.table.types.logical.DoubleType;
@@ -798,7 +799,7 @@ public class FlinkTypeFactory extends JavaTypeFactoryImpl
implements ExtendedRel
case INTERVAL_MINUTE:
case INTERVAL_MINUTE_SECOND:
case INTERVAL_SECOND:
- if (relDataType.getPrecision() > 3) {
+ if (relDataType.getPrecision() >
DayTimeIntervalType.MAX_DAY_PRECISION) {
throw new TableException(
"DAY_INTERVAL_TYPES precision is not supported: "
+ relDataType.getPrecision());
diff --git
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/api/QueryOperationTestPrograms.java
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/api/QueryOperationTestPrograms.java
index 6436d5ef854..16ea3502ab8 100644
---
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/api/QueryOperationTestPrograms.java
+++
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/api/QueryOperationTestPrograms.java
@@ -281,7 +281,7 @@ public class QueryOperationTestPrograms {
+ " TUMBLE((\n"
+ " SELECT `$$T_SOURCE`.`a`,
`$$T_SOURCE`.`b`, "
+ "`$$T_SOURCE`.`ts` FROM
`default_catalog`.`default_database`.`s` $$T_SOURCE\n"
- + " ), DESCRIPTOR(`ts`), INTERVAL
'0 00:00:05.0' DAY TO SECOND(3))\n"
+ + " ), DESCRIPTOR(`ts`), INTERVAL
'0 00:00:05.000' DAY(2) TO SECOND(3))\n"
+ " ) $$T_WIN_AGG GROUP BY
window_start, window_end, `$$T_WIN_AGG`.`b`\n"
+ ") $$T_PROJECT")
.build();
@@ -961,8 +961,8 @@ public class QueryOperationTestPrograms {
.runSql(
"SELECT `$$T_PROJECT`.`k`,
(LAST_VALUE(`$$T_PROJECT`.`v`) "
+ "OVER(PARTITION BY `$$T_PROJECT`.`k` "
- + "ORDER BY `$$T_PROJECT`.`ts` RANGE
BETWEEN INTERVAL '0 "
- + "00:00:02.0' DAY TO SECOND(3) PRECEDING
AND CURRENT ROW)) AS `_c1`, `$$T_PROJECT`.`ts` FROM (\n"
+ + "ORDER BY `$$T_PROJECT`.`ts` RANGE
BETWEEN INTERVAL "
+ + "'0 00:00:02.000' DAY(2) TO SECOND(3)
PRECEDING AND CURRENT ROW)) AS `_c1`, `$$T_PROJECT`.`ts` FROM (\n"
+ " SELECT `$$T_SOURCE`.`k`,
`$$T_SOURCE`.`v`, "
+ "`$$T_SOURCE`.`ts` FROM
`default_catalog`.`default_database`.`data` $$T_SOURCE\n"
+ ") $$T_PROJECT")
diff --git
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/calcite/FlinkTypeFactoryTest.java
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/calcite/FlinkTypeFactoryTest.java
index 7715f5bacdd..6e0e7a8c3d3 100644
---
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/calcite/FlinkTypeFactoryTest.java
+++
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/calcite/FlinkTypeFactoryTest.java
@@ -22,12 +22,15 @@ import
org.apache.flink.api.common.serialization.SerializerConfigImpl;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.common.typeinfo.Types;
import org.apache.flink.api.java.typeutils.runtime.kryo.KryoSerializer;
+import org.apache.flink.table.api.DataTypes;
+import org.apache.flink.table.api.TableException;
import org.apache.flink.table.legacy.types.logical.TypeInformationRawType;
import org.apache.flink.table.types.logical.ArrayType;
import org.apache.flink.table.types.logical.BigIntType;
import org.apache.flink.table.types.logical.BooleanType;
import org.apache.flink.table.types.logical.CharType;
import org.apache.flink.table.types.logical.DateType;
+import org.apache.flink.table.types.logical.DayTimeIntervalType;
import org.apache.flink.table.types.logical.DecimalType;
import org.apache.flink.table.types.logical.DoubleType;
import org.apache.flink.table.types.logical.FloatType;
@@ -47,7 +50,10 @@ import org.apache.flink.table.types.logical.VarBinaryType;
import org.apache.flink.table.types.logical.VarCharType;
import org.apache.flink.table.types.logical.utils.LogicalTypeMerging;
+import org.apache.calcite.avatica.util.TimeUnit;
import org.apache.calcite.rel.type.RelDataType;
+import org.apache.calcite.sql.SqlIntervalQualifier;
+import org.apache.calcite.sql.parser.SqlParserPos;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.parallel.Execution;
import org.junit.jupiter.api.parallel.ExecutionMode;
@@ -62,6 +68,7 @@ import java.util.stream.Collectors;
import java.util.stream.Stream;
import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
/** Tests for {@link FlinkTypeFactory}. */
@Execution(ExecutionMode.CONCURRENT)
@@ -137,6 +144,47 @@ class FlinkTypeFactoryTest {
.isEqualTo(logicalType.copy(true));
}
+ @Test
+ void testDayTimeIntervalLeadingPrecisionUpToMaxIsSupported() {
+ FlinkTypeFactory typeFactory =
+ new FlinkTypeFactory(
+ Thread.currentThread().getContextClassLoader(),
FlinkTypeSystem.INSTANCE);
+
+ RelDataType intervalType =
+ typeFactory.createSqlIntervalType(
+ new SqlIntervalQualifier(
+ TimeUnit.DAY,
+ DayTimeIntervalType.MAX_DAY_PRECISION,
+ TimeUnit.SECOND,
+ RelDataType.PRECISION_NOT_SPECIFIED,
+ SqlParserPos.ZERO));
+
+ assertThat(FlinkTypeFactory.toLogicalType(intervalType))
+
.isEqualTo(DataTypes.INTERVAL(DataTypes.SECOND(3)).notNull().getLogicalType());
+ }
+
+ @Test
+ void testDayTimeIntervalLeadingPrecisionAboveMaxIsRejected() {
+ FlinkTypeFactory typeFactory =
+ new FlinkTypeFactory(
+ Thread.currentThread().getContextClassLoader(),
FlinkTypeSystem.INSTANCE);
+
+ RelDataType intervalType =
+ typeFactory.createSqlIntervalType(
+ new SqlIntervalQualifier(
+ TimeUnit.DAY,
+ DayTimeIntervalType.MAX_DAY_PRECISION + 1,
+ TimeUnit.SECOND,
+ RelDataType.PRECISION_NOT_SPECIFIED,
+ SqlParserPos.ZERO));
+
+ assertThatThrownBy(() -> FlinkTypeFactory.toLogicalType(intervalType))
+ .isInstanceOf(TableException.class)
+ .hasMessageContaining(
+ "DAY_INTERVAL_TYPES precision is not supported: "
+ + (DayTimeIntervalType.MAX_DAY_PRECISION + 1));
+ }
+
@Test
void testDecimalInferType() {
assertThat(LogicalTypeMerging.findSumAggType(new DecimalType(10, 5)))
diff --git
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/expressions/LiteralExpressionsSerializationITCase.java
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/expressions/LiteralExpressionsSerializationITCase.java
index c24b84b1b00..7737003f65f 100644
---
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/expressions/LiteralExpressionsSerializationITCase.java
+++
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/expressions/LiteralExpressionsSerializationITCase.java
@@ -63,8 +63,13 @@ public class LiteralExpressionsSerializationITCase {
final LocalDateTime localDateTimeWithoutSeconds =
LocalDateTime.of(localDate, localTimeWithoutSeconds);
final Instant instant = Instant.ofEpochMilli(1234567);
- final Duration duration =
Duration.ofDays(99).plusSeconds(34).plusMillis(999);
- final Period period = Period.ofMonths(470);
+ final Duration defaultDuration =
Duration.ofDays(99).plusSeconds(34).plusMillis(999);
+ final Period defaultPeriod = Period.ofMonths(470);
+ final Period yearMonthPeriod = Period.ofMonths(9999 * 12 + 11);
+ final Duration dayHourDuration = Duration.ofDays(100).plusHours(5);
+ final Duration secondPrecisionDuration =
Duration.ofSeconds(45).plusMillis(750);
+ final Duration maxDayDuration = Duration.ofDays(999_999);
+ final Period derivedLargePeriod = Period.ofMonths(1200);
final Table t =
env.fromValues(1)
.select(
@@ -90,13 +95,27 @@ public class LiteralExpressionsSerializationITCase {
lit(instant,
DataTypes.TIMESTAMP_LTZ(6).notNull()),
lit(instant,
DataTypes.TIMESTAMP_LTZ(9).notNull()),
lit(
- duration,
+ defaultDuration,
DataTypes.INTERVAL(DataTypes.DAY(),
DataTypes.SECOND(9))
.notNull()),
lit(
- period,
+ defaultPeriod,
DataTypes.INTERVAL(DataTypes.YEAR(),
DataTypes.MONTH())
- .notNull()));
+ .notNull()),
+ lit(
+ yearMonthPeriod,
+ DataTypes.INTERVAL(DataTypes.YEAR(4),
DataTypes.MONTH())
+ .notNull()),
+ lit(
+ dayHourDuration,
+ DataTypes.INTERVAL(DataTypes.DAY(3),
DataTypes.HOUR())
+ .notNull()),
+ lit(
+ secondPrecisionDuration,
+
DataTypes.INTERVAL(DataTypes.SECOND(4)).notNull()),
+ lit(maxDayDuration,
DataTypes.INTERVAL(DataTypes.DAY(6)).notNull()),
+ lit(defaultPeriod),
+ lit(derivedLargePeriod));
final ProjectQueryOperation operation = (ProjectQueryOperation)
t.getQueryOperation();
final String exprStr =
operation.getProjectList().stream()
@@ -129,8 +148,14 @@ public class LiteralExpressionsSerializationITCase {
+ "TO_TIMESTAMP_LTZ(1234567, 3),\n"
+ "TO_TIMESTAMP_LTZ(1234567000, 6),\n"
+ "TO_TIMESTAMP_LTZ(1234567000000, 9),\n"
- + "INTERVAL '99 00:00:34.999' DAY TO
SECOND(3),\n"
- + "INTERVAL '39-2' YEAR TO MONTH");
+ + "INTERVAL '99 00:00:34.999000000' DAY(2) TO
SECOND(9),\n"
+ + "INTERVAL '39-2' YEAR(2) TO MONTH,\n"
+ + "INTERVAL '9999-11' YEAR(4) TO MONTH,\n"
+ + "INTERVAL '100 05:00:00.000000' DAY(3) TO
SECOND(6),\n"
+ + "INTERVAL '0 00:00:45.7500' DAY(2) TO
SECOND(4),\n"
+ + "INTERVAL '999999 00:00:00.000000' DAY(6) TO
SECOND(6),\n"
+ + "INTERVAL '39-2' YEAR(2) TO MONTH,\n"
+ + "INTERVAL '100-0' YEAR(3) TO MONTH");
final TableResult tableResult = env.sqlQuery(String.format("SELECT
%s", exprStr)).execute();
final List<Row> results =
CollectionUtil.iteratorToList(tableResult.collect());
@@ -158,7 +183,13 @@ public class LiteralExpressionsSerializationITCase {
instant,
instant,
instant,
- duration,
- period));
+ defaultDuration,
+ defaultPeriod,
+ yearMonthPeriod,
+ dayHourDuration,
+ secondPrecisionDuration,
+ maxDayDuration,
+ defaultPeriod,
+ derivedLargePeriod));
}
}