raminqaf commented on code in PR #29063:
URL: https://github.com/apache/flink/pull/29063#discussion_r3932968384
##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/JsonFunctionsITCase.java:
##########
@@ -480,6 +484,41 @@ private static TestSetSpec jsonValueSpec() {
"JSON_VALUE(f0, '$.longBalance' RETURNING DOUBLE)",
123456789.987654321,
DOUBLE())
+ .testResult(
+ $("f0").jsonValue("$.age", TINYINT()),
+ "JSON_VALUE(f0, '$.age' RETURNING TINYINT)",
+ (byte) 42,
+ TINYINT())
+ .testResult(
+ $("f0").jsonValue("$.age", SMALLINT()),
+ "JSON_VALUE(f0, '$.age' RETURNING SMALLINT)",
+ (short) 42,
+ SMALLINT())
+ .testResult(
+ $("f0").jsonValue("$.age", BIGINT()),
+ "JSON_VALUE(f0, '$.age' RETURNING BIGINT)",
+ 42L,
+ BIGINT())
+ .testResult(
+ $("f0").jsonValue("$.bigCount", BIGINT()),
+ "JSON_VALUE(f0, '$.bigCount' RETURNING BIGINT)",
+ 9999999999L,
+ BIGINT())
+ .testResult(
+ $("f0").jsonValue("$.balance", FLOAT()),
+ "JSON_VALUE(f0, '$.balance' RETURNING FLOAT)",
+ 13.37f,
+ FLOAT())
+ .testResult(
+ $("f0").jsonValue("$.balance", DECIMAL(10, 2)),
+ "JSON_VALUE(f0, '$.balance' RETURNING DECIMAL(10, 2))",
+ new java.math.BigDecimal("13.37"),
+ DECIMAL(10, 2))
+ .testResult(
+ $("f0").jsonValue("$.longBalance", DECIMAL(30, 10)),
Review Comment:
Let's add a test for the case `DECIMAL(9, 2)`
##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/JsonFunctionsITCase.java:
##########
@@ -480,6 +484,41 @@ private static TestSetSpec jsonValueSpec() {
"JSON_VALUE(f0, '$.longBalance' RETURNING DOUBLE)",
123456789.987654321,
DOUBLE())
+ .testResult(
+ $("f0").jsonValue("$.age", TINYINT()),
+ "JSON_VALUE(f0, '$.age' RETURNING TINYINT)",
+ (byte) 42,
+ TINYINT())
+ .testResult(
+ $("f0").jsonValue("$.age", SMALLINT()),
+ "JSON_VALUE(f0, '$.age' RETURNING SMALLINT)",
+ (short) 42,
+ SMALLINT())
+ .testResult(
+ $("f0").jsonValue("$.age", BIGINT()),
+ "JSON_VALUE(f0, '$.age' RETURNING BIGINT)",
+ 42L,
+ BIGINT())
+ .testResult(
+ $("f0").jsonValue("$.bigCount", BIGINT()),
+ "JSON_VALUE(f0, '$.bigCount' RETURNING BIGINT)",
+ 9999999999L,
+ BIGINT())
+ .testResult(
+ $("f0").jsonValue("$.balance", FLOAT()),
+ "JSON_VALUE(f0, '$.balance' RETURNING FLOAT)",
+ 13.37f,
+ FLOAT())
+ .testResult(
+ $("f0").jsonValue("$.balance", DECIMAL(10, 2)),
+ "JSON_VALUE(f0, '$.balance' RETURNING DECIMAL(10, 2))",
+ new java.math.BigDecimal("13.37"),
Review Comment:
```suggestion
new BigDecimal("13.37"),
```
##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/JsonFunctionsITCase.java:
##########
@@ -886,7 +1083,214 @@ private static List<TestSetSpec> jsonQuerySpec() {
.testTableApiRuntimeError(
$("f0").jsonQuery("strict $.err10",
WITHOUT_ARRAY, NULL, ERROR),
TableRuntimeException.class,
- "No results for path"));
+ "No results for path"),
+
+ // Typed RETURNING ARRAY<T> support
+ TestSetSpec.forFunction(BuiltInFunctionDefinitions.JSON_QUERY)
+ .onFieldsWithData(
+ "{\"ints\": [1, 2, 3], \"doubles\": [1.5,
2.5], \"bools\": [true, false], \"withNull\": [1, null, 3], \"bigints\": [1,
9999999999]}")
+ .andDataTypes(STRING())
+ .testResult(
+ $("f0").jsonQuery("$.ints", ARRAY(INT())),
+ "JSON_QUERY(f0, '$.ints' RETURNING
ARRAY<INT>)",
+ new Integer[] {1, 2, 3},
+ ARRAY(INT()))
+ .testResult(
+ $("f0").jsonQuery("$.ints", ARRAY(TINYINT())),
+ "JSON_QUERY(f0, '$.ints' RETURNING
ARRAY<TINYINT>)",
+ new Byte[] {1, 2, 3},
+ ARRAY(TINYINT()))
+ .testResult(
+ $("f0").jsonQuery("$.ints", ARRAY(SMALLINT())),
+ "JSON_QUERY(f0, '$.ints' RETURNING
ARRAY<SMALLINT>)",
+ new Short[] {1, 2, 3},
+ ARRAY(SMALLINT()))
+ .testResult(
+ $("f0").jsonQuery("$.ints", ARRAY(BIGINT())),
+ "JSON_QUERY(f0, '$.ints' RETURNING
ARRAY<BIGINT>)",
+ new Long[] {1L, 2L, 3L},
+ ARRAY(BIGINT()))
+ .testResult(
+ $("f0").jsonQuery("$.bigints",
ARRAY(BIGINT())),
+ "JSON_QUERY(f0, '$.bigints' RETURNING
ARRAY<BIGINT>)",
+ new Long[] {1L, 9999999999L},
+ ARRAY(BIGINT()))
+ .testResult(
+ $("f0").jsonQuery("$.doubles",
ARRAY(DOUBLE())),
+ "JSON_QUERY(f0, '$.doubles' RETURNING
ARRAY<DOUBLE>)",
+ new Double[] {1.5, 2.5},
+ ARRAY(DOUBLE()))
Review Comment:
Add a test with ARRAY<INT>
##########
flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/functions/SqlJsonUtilsConversionTest.java:
##########
@@ -0,0 +1,396 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.flink.table.runtime.functions;
+
+import org.apache.flink.table.api.JsonValueOnEmptyOrError;
+import org.apache.flink.table.api.TableRuntimeException;
+import org.apache.flink.table.data.DecimalData;
+import org.apache.flink.table.data.GenericArrayData;
+import org.apache.flink.table.types.logical.LogicalTypeRoot;
+
+import org.junit.jupiter.api.Nested;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.Arguments;
+import org.junit.jupiter.params.provider.MethodSource;
+
+import java.math.BigDecimal;
+import java.math.BigInteger;
+import java.util.stream.Stream;
+
+import static org.apache.flink.table.api.JsonQueryOnEmptyOrError.EMPTY_ARRAY;
+import static org.apache.flink.table.api.JsonQueryOnEmptyOrError.ERROR;
+import static org.apache.flink.table.api.JsonQueryOnEmptyOrError.NULL;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.BIGINT;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.BOOLEAN;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.DECIMAL;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.DOUBLE;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.FLOAT;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.INTEGER;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.SMALLINT;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.TINYINT;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+class SqlJsonUtilsConversionTest {
+
+ private static Object scalar(Object raw, LogicalTypeRoot type) {
+ return SqlJsonUtils.convertJsonScalar(raw, type, 0, 0,
JsonValueOnEmptyOrError.NULL, null);
+ }
+
+ private static Object scalar(
+ Object raw,
+ LogicalTypeRoot type,
+ JsonValueOnEmptyOrError onError,
+ Object defaultValue) {
+ return SqlJsonUtils.convertJsonScalar(raw, type, 0, 0, onError,
defaultValue);
+ }
+
+ private static Object scalarDecimal(Object raw, int precision, int scale) {
+ return SqlJsonUtils.convertJsonScalar(
+ raw, DECIMAL, precision, scale, JsonValueOnEmptyOrError.NULL,
null);
+ }
+
+ private static GenericArrayData array(Object[] raw, LogicalTypeRoot
elementType) {
+ return SqlJsonUtils.convertJsonArray(raw, elementType, 0, 0, NULL);
+ }
+
+ private static GenericArrayData arrayDecimal(
+ Object[] raw,
+ int precision,
+ int scale,
+ org.apache.flink.table.api.JsonQueryOnEmptyOrError onError) {
Review Comment:
```suggestion
JsonQueryOnEmptyOrError onError) {
```
##########
flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/functions/SqlJsonUtilsConversionTest.java:
##########
@@ -0,0 +1,396 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.flink.table.runtime.functions;
+
+import org.apache.flink.table.api.JsonValueOnEmptyOrError;
+import org.apache.flink.table.api.TableRuntimeException;
+import org.apache.flink.table.data.DecimalData;
+import org.apache.flink.table.data.GenericArrayData;
+import org.apache.flink.table.types.logical.LogicalTypeRoot;
+
+import org.junit.jupiter.api.Nested;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.Arguments;
+import org.junit.jupiter.params.provider.MethodSource;
+
+import java.math.BigDecimal;
+import java.math.BigInteger;
+import java.util.stream.Stream;
+
+import static org.apache.flink.table.api.JsonQueryOnEmptyOrError.EMPTY_ARRAY;
+import static org.apache.flink.table.api.JsonQueryOnEmptyOrError.ERROR;
+import static org.apache.flink.table.api.JsonQueryOnEmptyOrError.NULL;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.BIGINT;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.BOOLEAN;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.DECIMAL;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.DOUBLE;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.FLOAT;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.INTEGER;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.SMALLINT;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.TINYINT;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+class SqlJsonUtilsConversionTest {
+
+ private static Object scalar(Object raw, LogicalTypeRoot type) {
+ return SqlJsonUtils.convertJsonScalar(raw, type, 0, 0,
JsonValueOnEmptyOrError.NULL, null);
+ }
+
+ private static Object scalar(
+ Object raw,
+ LogicalTypeRoot type,
+ JsonValueOnEmptyOrError onError,
+ Object defaultValue) {
+ return SqlJsonUtils.convertJsonScalar(raw, type, 0, 0, onError,
defaultValue);
+ }
+
+ private static Object scalarDecimal(Object raw, int precision, int scale) {
+ return SqlJsonUtils.convertJsonScalar(
+ raw, DECIMAL, precision, scale, JsonValueOnEmptyOrError.NULL,
null);
+ }
+
+ private static GenericArrayData array(Object[] raw, LogicalTypeRoot
elementType) {
+ return SqlJsonUtils.convertJsonArray(raw, elementType, 0, 0, NULL);
+ }
+
+ private static GenericArrayData arrayDecimal(
+ Object[] raw,
+ int precision,
+ int scale,
+ org.apache.flink.table.api.JsonQueryOnEmptyOrError onError) {
+ return SqlJsonUtils.convertJsonArray(raw, DECIMAL, precision, scale,
onError);
+ }
+
+ static Stream<Arguments> validConversions() {
+ return Stream.of(
+ Arguments.of(42, TINYINT, (byte) 42),
+ Arguments.of(127L, TINYINT, Byte.MAX_VALUE),
+ Arguments.of(-128L, TINYINT, Byte.MIN_VALUE),
+ Arguments.of(1000, SMALLINT, (short) 1000),
+ Arguments.of((long) Short.MAX_VALUE, SMALLINT,
Short.MAX_VALUE),
+ Arguments.of(42L, INTEGER, 42),
+ Arguments.of((long) Integer.MAX_VALUE, INTEGER,
Integer.MAX_VALUE),
+ Arguments.of(42, BIGINT, 42L),
+ Arguments.of(9_999_999_999L, BIGINT, 9_999_999_999L),
+ Arguments.of(13.37, FLOAT, 13.37f),
+ Arguments.of(13.37, DOUBLE, 13.37));
+ }
+
+ static Stream<Arguments> overflowCases() {
+ return Stream.of(
+ Arguments.of(128, TINYINT),
+ Arguments.of(-129, TINYINT),
+ Arguments.of(200L, TINYINT),
+ Arguments.of(-200L, TINYINT),
+ Arguments.of((long) Short.MAX_VALUE + 1, SMALLINT),
+ Arguments.of((long) Short.MIN_VALUE - 1, SMALLINT),
+ Arguments.of(40_000L, SMALLINT),
+ Arguments.of(-40_000L, SMALLINT),
+ Arguments.of(9_999_999_999L, INTEGER),
+ Arguments.of(-9_999_999_999L, INTEGER));
+ }
+
+ @Nested
+ class RangeCheckedConversions {
+
+ @ParameterizedTest(name = "{0} -> {1}")
+ @MethodSource(
+
"org.apache.flink.table.runtime.functions.SqlJsonUtilsConversionTest#validConversions")
Review Comment:
If the method is used only in this nested class move it in and
```suggestion
"validConversions")
```
##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/JsonFunctionsITCase.java:
##########
@@ -528,6 +567,159 @@ private static TestSetSpec jsonValueSpec() {
"JSON_VALUE(f0, 'strict $.invalid' RETURNING INTEGER
NULL ON EMPTY DEFAULT 42 ON ERROR)",
42,
INT())
+ .testResult(
+ $("f0").jsonValue(
+ "lax $.invalid",
+ BIGINT(),
+ JsonValueOnEmptyOrError.DEFAULT,
+ 99L,
+ JsonValueOnEmptyOrError.ERROR,
+ null),
+ "JSON_VALUE(f0, 'lax $.invalid' RETURNING BIGINT
DEFAULT 99 ON EMPTY ERROR ON ERROR)",
+ 99L,
+ BIGINT())
+
+ // JSON null at a valid path triggers ON EMPTY (default: NULL)
+ .testResult(
+ $("f0").jsonValue("$.nullField", INT()),
+ "JSON_VALUE(f0, '$.nullField' RETURNING INTEGER)",
+ null,
+ INT())
+
+ // Type mismatch: string value cast to numeric triggers ON
ERROR
+ .testResult(
+ $("f0").jsonValue(
+ "$.type",
+ INT(),
+ JsonValueOnEmptyOrError.NULL,
+ null,
+ JsonValueOnEmptyOrError.DEFAULT,
+ 42),
+ "JSON_VALUE(f0, '$.type' RETURNING INTEGER DEFAULT 42
ON ERROR)",
+ 42,
+ INT())
+ .testResult(
+ $("f0").jsonValue(
+ "$.type",
+ INT(),
+ JsonValueOnEmptyOrError.NULL,
+ null,
+ JsonValueOnEmptyOrError.NULL,
+ null),
+ "JSON_VALUE(f0, '$.type' RETURNING INTEGER NULL ON
ERROR)",
+ null,
+ INT())
+ .testResult(
+ $("f0").jsonValue(
+ "$.type",
+ BIGINT(),
+ JsonValueOnEmptyOrError.NULL,
+ null,
+ JsonValueOnEmptyOrError.DEFAULT,
+ 0L),
+ "JSON_VALUE(f0, '$.type' RETURNING BIGINT DEFAULT 0 ON
ERROR)",
+ 0L,
+ BIGINT())
+ .testSqlRuntimeError(
+ "JSON_VALUE(f0, '$.type' RETURNING INTEGER ERROR ON
ERROR)",
+ TableRuntimeException.class,
+ "Cannot cast")
+
+ // Numeric overflow triggers ON ERROR (not silent wrapping)
+ .testResult(
+ $("f0").jsonValue("$.bigCount", INT()),
+ "JSON_VALUE(f0, '$.bigCount' RETURNING INTEGER)",
+ null,
+ INT())
+ // Fractional truncation (13.37 -> 13): MySQL truncates,
PostgreSQL errors.
+ // We match MySQL (CAST(JSON_UNQUOTE(JSON_EXTRACT(...)) AS
type)).
Review Comment:
We should mention that we truncate toward zero. 49.99 -> 49 and 49.01 -> 49
##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/functions/SqlJsonUtils.java:
##########
@@ -815,6 +828,233 @@ public String toString() {
}
}
+ // --- Type conversion for JSON_VALUE / JSON_QUERY RETURNING ---
+
+ private static final java.util.EnumSet<LogicalTypeRoot>
SUPPORTED_JSON_RETURNING_TYPES =
+ java.util.EnumSet.of(
+ LogicalTypeRoot.BOOLEAN,
+ LogicalTypeRoot.TINYINT,
+ LogicalTypeRoot.SMALLINT,
+ LogicalTypeRoot.INTEGER,
+ LogicalTypeRoot.BIGINT,
+ LogicalTypeRoot.FLOAT,
+ LogicalTypeRoot.DOUBLE,
+ LogicalTypeRoot.DECIMAL);
+
+ public static boolean isSupportedJsonReturningType(LogicalTypeRoot
typeRoot) {
+ return SUPPORTED_JSON_RETURNING_TYPES.contains(typeRoot);
+ }
+
+ public static Object convertJsonScalar(
+ Object raw,
+ LogicalTypeRoot typeRoot,
+ int precision,
+ int scale,
+ JsonValueOnEmptyOrError errorBehavior,
+ Object defaultValue) {
+ if (raw == null) {
+ return null;
+ }
+ try {
+ return convertToType(raw, typeRoot, precision, scale);
+ } catch (JsonConversionException e) {
+ switch (errorBehavior) {
+ case NULL:
+ return null;
+ case DEFAULT:
+ return convertDefault(defaultValue, typeRoot, precision,
scale);
+ case ERROR:
+ throw new TableRuntimeException(
+ "Cannot cast " + raw.getClass().getName() + " to "
+ typeRoot, e);
+ default:
+ throw new TableRuntimeException(
+ "Unreachable: unknown error behavior " +
errorBehavior);
+ }
+ }
+ }
+
+ public static GenericArrayData convertJsonArray(
+ Object rawResult,
+ LogicalTypeRoot elementTypeRoot,
+ int precision,
+ int scale,
+ JsonQueryOnEmptyOrError errorBehavior) {
+ if (rawResult == null) {
+ return null;
+ }
+ try {
+ Object[] rawArr = (Object[]) rawResult;
+ Object[] converted = new Object[rawArr.length];
+ for (int i = 0; i < rawArr.length; i++) {
+ if (rawArr[i] != null) {
+ converted[i] = convertToType(rawArr[i], elementTypeRoot,
precision, scale);
+ }
+ }
+ return new GenericArrayData(converted);
+ } catch (JsonConversionException e) {
+ switch (errorBehavior) {
+ case NULL:
+ return null;
+ case EMPTY_ARRAY:
+ return new GenericArrayData(new Object[0]);
+ case ERROR:
+ throw new TableRuntimeException("Array element type
mismatch in JSON_QUERY", e);
+ default:
+ return null;
+ }
+ }
+ }
+
+ private static Object convertToType(
+ Object raw, LogicalTypeRoot typeRoot, int precision, int scale) {
+ if (raw instanceof StringData) {
+ return convertToType(raw.toString(), typeRoot, precision, scale);
+ }
+ if (raw instanceof String) {
+ if (typeRoot == LogicalTypeRoot.BOOLEAN) {
+ return parseStringAsBoolean((String) raw);
+ }
+ try {
+ return convertToType(new BigDecimal((String) raw), typeRoot,
precision, scale);
+ } catch (NumberFormatException e) {
+ throw new JsonConversionException(
+ "Cannot parse string '" + raw + "' as " + typeRoot, e);
+ }
+ }
+ try {
+ switch (typeRoot) {
+ case BOOLEAN:
+ return (Boolean) raw;
+ case TINYINT:
+ return toCheckedByte((Number) raw);
+ case SMALLINT:
+ return toCheckedShort((Number) raw);
+ case INTEGER:
+ return toCheckedInt((Number) raw);
+ case BIGINT:
+ return toCheckedLong((Number) raw);
+ case FLOAT:
+ return toCheckedFloat((Number) raw);
+ case DOUBLE:
+ return toCheckedDouble((Number) raw);
+ case DECIMAL:
+ return toCheckedDecimal(((Number) raw).toString(),
precision, scale);
+ default:
+ throw new JsonConversionException(
+ "Unsupported type for JSON conversion: " +
typeRoot);
+ }
+ } catch (ClassCastException e) {
+ throw new JsonConversionException(
+ "Cannot convert " + raw.getClass().getName() + " to " +
typeRoot, e);
+ }
+ }
+
+ private static Object convertDefault(
+ Object defaultValue, LogicalTypeRoot typeRoot, int precision, int
scale) {
+ if (defaultValue == null) {
+ return null;
+ }
+ try {
+ return convertToType(defaultValue, typeRoot, precision, scale);
+ } catch (JsonConversionException e) {
+ throw new TableRuntimeException(
+ "Default value " + defaultValue + " cannot be represented
as " + typeRoot, e);
+ }
+ }
+
+ static byte toCheckedByte(Number n) {
+ long v = n.longValue();
+ if (v < Byte.MIN_VALUE || v > Byte.MAX_VALUE) {
+ throw new JsonConversionException("Value " + n + " is out of range
for TINYINT");
+ }
+ return (byte) v;
+ }
+
+ static short toCheckedShort(Number n) {
+ long v = n.longValue();
+ if (v < Short.MIN_VALUE || v > Short.MAX_VALUE) {
+ throw new JsonConversionException("Value " + n + " is out of range
for SMALLINT");
+ }
+ return (short) v;
+ }
+
+ static int toCheckedInt(Number n) {
+ long v = n.longValue();
Review Comment:
This will be error prone... if the input is >2^64:
```sql
JSON_VALUE('{"x": 18446744073709551616}', 'lax $.x' RETURNING INTEGER)
-- expected: overflow → NULL (or ON ERROR); actual: 0
JSON_VALUE('{"x": 18446744073709551616}', 'lax $.x' RETURNING TINYINT)
-- actual: 0 (18446744073709551616 = 2^64)
```
##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/functions/SqlJsonUtils.java:
##########
@@ -815,6 +828,233 @@ public String toString() {
}
}
+ // --- Type conversion for JSON_VALUE / JSON_QUERY RETURNING ---
+
+ private static final java.util.EnumSet<LogicalTypeRoot>
SUPPORTED_JSON_RETURNING_TYPES =
+ java.util.EnumSet.of(
+ LogicalTypeRoot.BOOLEAN,
+ LogicalTypeRoot.TINYINT,
+ LogicalTypeRoot.SMALLINT,
+ LogicalTypeRoot.INTEGER,
+ LogicalTypeRoot.BIGINT,
+ LogicalTypeRoot.FLOAT,
+ LogicalTypeRoot.DOUBLE,
+ LogicalTypeRoot.DECIMAL);
+
+ public static boolean isSupportedJsonReturningType(LogicalTypeRoot
typeRoot) {
+ return SUPPORTED_JSON_RETURNING_TYPES.contains(typeRoot);
+ }
+
+ public static Object convertJsonScalar(
+ Object raw,
+ LogicalTypeRoot typeRoot,
+ int precision,
+ int scale,
+ JsonValueOnEmptyOrError errorBehavior,
+ Object defaultValue) {
+ if (raw == null) {
+ return null;
+ }
+ try {
+ return convertToType(raw, typeRoot, precision, scale);
+ } catch (JsonConversionException e) {
+ switch (errorBehavior) {
+ case NULL:
+ return null;
+ case DEFAULT:
+ return convertDefault(defaultValue, typeRoot, precision,
scale);
+ case ERROR:
+ throw new TableRuntimeException(
+ "Cannot cast " + raw.getClass().getName() + " to "
+ typeRoot, e);
+ default:
+ throw new TableRuntimeException(
+ "Unreachable: unknown error behavior " +
errorBehavior);
+ }
+ }
+ }
+
+ public static GenericArrayData convertJsonArray(
+ Object rawResult,
+ LogicalTypeRoot elementTypeRoot,
+ int precision,
+ int scale,
+ JsonQueryOnEmptyOrError errorBehavior) {
+ if (rawResult == null) {
+ return null;
+ }
+ try {
+ Object[] rawArr = (Object[]) rawResult;
+ Object[] converted = new Object[rawArr.length];
+ for (int i = 0; i < rawArr.length; i++) {
+ if (rawArr[i] != null) {
+ converted[i] = convertToType(rawArr[i], elementTypeRoot,
precision, scale);
+ }
+ }
+ return new GenericArrayData(converted);
+ } catch (JsonConversionException e) {
+ switch (errorBehavior) {
+ case NULL:
+ return null;
+ case EMPTY_ARRAY:
+ return new GenericArrayData(new Object[0]);
+ case ERROR:
+ throw new TableRuntimeException("Array element type
mismatch in JSON_QUERY", e);
+ default:
+ return null;
+ }
+ }
+ }
+
+ private static Object convertToType(
+ Object raw, LogicalTypeRoot typeRoot, int precision, int scale) {
+ if (raw instanceof StringData) {
+ return convertToType(raw.toString(), typeRoot, precision, scale);
+ }
+ if (raw instanceof String) {
+ if (typeRoot == LogicalTypeRoot.BOOLEAN) {
+ return parseStringAsBoolean((String) raw);
+ }
+ try {
+ return convertToType(new BigDecimal((String) raw), typeRoot,
precision, scale);
+ } catch (NumberFormatException e) {
+ throw new JsonConversionException(
+ "Cannot parse string '" + raw + "' as " + typeRoot, e);
+ }
+ }
+ try {
+ switch (typeRoot) {
+ case BOOLEAN:
+ return (Boolean) raw;
Review Comment:
```suggestion
return raw;
```
##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/functions/SqlJsonUtils.java:
##########
@@ -815,6 +828,233 @@ public String toString() {
}
}
+ // --- Type conversion for JSON_VALUE / JSON_QUERY RETURNING ---
+
+ private static final java.util.EnumSet<LogicalTypeRoot>
SUPPORTED_JSON_RETURNING_TYPES =
+ java.util.EnumSet.of(
+ LogicalTypeRoot.BOOLEAN,
+ LogicalTypeRoot.TINYINT,
+ LogicalTypeRoot.SMALLINT,
+ LogicalTypeRoot.INTEGER,
+ LogicalTypeRoot.BIGINT,
+ LogicalTypeRoot.FLOAT,
+ LogicalTypeRoot.DOUBLE,
+ LogicalTypeRoot.DECIMAL);
+
+ public static boolean isSupportedJsonReturningType(LogicalTypeRoot
typeRoot) {
+ return SUPPORTED_JSON_RETURNING_TYPES.contains(typeRoot);
+ }
+
+ public static Object convertJsonScalar(
+ Object raw,
+ LogicalTypeRoot typeRoot,
+ int precision,
+ int scale,
+ JsonValueOnEmptyOrError errorBehavior,
+ Object defaultValue) {
+ if (raw == null) {
+ return null;
+ }
+ try {
+ return convertToType(raw, typeRoot, precision, scale);
+ } catch (JsonConversionException e) {
+ switch (errorBehavior) {
+ case NULL:
+ return null;
+ case DEFAULT:
+ return convertDefault(defaultValue, typeRoot, precision,
scale);
+ case ERROR:
+ throw new TableRuntimeException(
+ "Cannot cast " + raw.getClass().getName() + " to "
+ typeRoot, e);
+ default:
+ throw new TableRuntimeException(
+ "Unreachable: unknown error behavior " +
errorBehavior);
+ }
+ }
+ }
+
+ public static GenericArrayData convertJsonArray(
+ Object rawResult,
+ LogicalTypeRoot elementTypeRoot,
+ int precision,
+ int scale,
+ JsonQueryOnEmptyOrError errorBehavior) {
+ if (rawResult == null) {
+ return null;
+ }
+ try {
+ Object[] rawArr = (Object[]) rawResult;
+ Object[] converted = new Object[rawArr.length];
+ for (int i = 0; i < rawArr.length; i++) {
+ if (rawArr[i] != null) {
+ converted[i] = convertToType(rawArr[i], elementTypeRoot,
precision, scale);
+ }
+ }
+ return new GenericArrayData(converted);
+ } catch (JsonConversionException e) {
+ switch (errorBehavior) {
+ case NULL:
+ return null;
+ case EMPTY_ARRAY:
+ return new GenericArrayData(new Object[0]);
+ case ERROR:
+ throw new TableRuntimeException("Array element type
mismatch in JSON_QUERY", e);
+ default:
+ return null;
+ }
+ }
+ }
+
+ private static Object convertToType(
+ Object raw, LogicalTypeRoot typeRoot, int precision, int scale) {
+ if (raw instanceof StringData) {
+ return convertToType(raw.toString(), typeRoot, precision, scale);
+ }
+ if (raw instanceof String) {
+ if (typeRoot == LogicalTypeRoot.BOOLEAN) {
+ return parseStringAsBoolean((String) raw);
+ }
+ try {
+ return convertToType(new BigDecimal((String) raw), typeRoot,
precision, scale);
+ } catch (NumberFormatException e) {
+ throw new JsonConversionException(
+ "Cannot parse string '" + raw + "' as " + typeRoot, e);
+ }
+ }
+ try {
+ switch (typeRoot) {
+ case BOOLEAN:
+ return (Boolean) raw;
+ case TINYINT:
+ return toCheckedByte((Number) raw);
+ case SMALLINT:
+ return toCheckedShort((Number) raw);
+ case INTEGER:
+ return toCheckedInt((Number) raw);
+ case BIGINT:
+ return toCheckedLong((Number) raw);
+ case FLOAT:
+ return toCheckedFloat((Number) raw);
+ case DOUBLE:
+ return toCheckedDouble((Number) raw);
+ case DECIMAL:
+ return toCheckedDecimal(((Number) raw).toString(),
precision, scale);
Review Comment:
```suggestion
return toCheckedDecimal(raw.toString(), precision,
scale);
```
##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/functions/SqlJsonUtils.java:
##########
@@ -815,6 +828,233 @@ public String toString() {
}
}
+ // --- Type conversion for JSON_VALUE / JSON_QUERY RETURNING ---
+
+ private static final java.util.EnumSet<LogicalTypeRoot>
SUPPORTED_JSON_RETURNING_TYPES =
+ java.util.EnumSet.of(
+ LogicalTypeRoot.BOOLEAN,
+ LogicalTypeRoot.TINYINT,
+ LogicalTypeRoot.SMALLINT,
+ LogicalTypeRoot.INTEGER,
+ LogicalTypeRoot.BIGINT,
+ LogicalTypeRoot.FLOAT,
+ LogicalTypeRoot.DOUBLE,
+ LogicalTypeRoot.DECIMAL);
+
+ public static boolean isSupportedJsonReturningType(LogicalTypeRoot
typeRoot) {
+ return SUPPORTED_JSON_RETURNING_TYPES.contains(typeRoot);
+ }
+
+ public static Object convertJsonScalar(
+ Object raw,
+ LogicalTypeRoot typeRoot,
+ int precision,
+ int scale,
+ JsonValueOnEmptyOrError errorBehavior,
+ Object defaultValue) {
+ if (raw == null) {
+ return null;
+ }
+ try {
+ return convertToType(raw, typeRoot, precision, scale);
+ } catch (JsonConversionException e) {
+ switch (errorBehavior) {
+ case NULL:
+ return null;
+ case DEFAULT:
+ return convertDefault(defaultValue, typeRoot, precision,
scale);
+ case ERROR:
+ throw new TableRuntimeException(
+ "Cannot cast " + raw.getClass().getName() + " to "
+ typeRoot, e);
+ default:
+ throw new TableRuntimeException(
+ "Unreachable: unknown error behavior " +
errorBehavior);
+ }
+ }
+ }
+
+ public static GenericArrayData convertJsonArray(
+ Object rawResult,
+ LogicalTypeRoot elementTypeRoot,
+ int precision,
+ int scale,
+ JsonQueryOnEmptyOrError errorBehavior) {
+ if (rawResult == null) {
+ return null;
+ }
+ try {
+ Object[] rawArr = (Object[]) rawResult;
+ Object[] converted = new Object[rawArr.length];
+ for (int i = 0; i < rawArr.length; i++) {
+ if (rawArr[i] != null) {
+ converted[i] = convertToType(rawArr[i], elementTypeRoot,
precision, scale);
+ }
+ }
+ return new GenericArrayData(converted);
+ } catch (JsonConversionException e) {
+ switch (errorBehavior) {
+ case NULL:
+ return null;
+ case EMPTY_ARRAY:
+ return new GenericArrayData(new Object[0]);
+ case ERROR:
+ throw new TableRuntimeException("Array element type
mismatch in JSON_QUERY", e);
+ default:
+ return null;
+ }
+ }
+ }
+
+ private static Object convertToType(
+ Object raw, LogicalTypeRoot typeRoot, int precision, int scale) {
+ if (raw instanceof StringData) {
+ return convertToType(raw.toString(), typeRoot, precision, scale);
+ }
+ if (raw instanceof String) {
+ if (typeRoot == LogicalTypeRoot.BOOLEAN) {
+ return parseStringAsBoolean((String) raw);
Review Comment:
Do we also consider integer `1` as `true` and `0` as `false`?
##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/JsonFunctionsITCase.java:
##########
@@ -886,7 +1083,214 @@ private static List<TestSetSpec> jsonQuerySpec() {
.testTableApiRuntimeError(
$("f0").jsonQuery("strict $.err10",
WITHOUT_ARRAY, NULL, ERROR),
TableRuntimeException.class,
- "No results for path"));
+ "No results for path"),
+
+ // Typed RETURNING ARRAY<T> support
+ TestSetSpec.forFunction(BuiltInFunctionDefinitions.JSON_QUERY)
+ .onFieldsWithData(
+ "{\"ints\": [1, 2, 3], \"doubles\": [1.5,
2.5], \"bools\": [true, false], \"withNull\": [1, null, 3], \"bigints\": [1,
9999999999]}")
Review Comment:
how about `"zeroAndOne": [0, 1]`, `"yesAndNo": ['Y', 'n']`"?
##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/JsonFunctionsITCase.java:
##########
@@ -886,7 +1083,214 @@ private static List<TestSetSpec> jsonQuerySpec() {
.testTableApiRuntimeError(
$("f0").jsonQuery("strict $.err10",
WITHOUT_ARRAY, NULL, ERROR),
TableRuntimeException.class,
- "No results for path"));
+ "No results for path"),
+
+ // Typed RETURNING ARRAY<T> support
+ TestSetSpec.forFunction(BuiltInFunctionDefinitions.JSON_QUERY)
+ .onFieldsWithData(
+ "{\"ints\": [1, 2, 3], \"doubles\": [1.5,
2.5], \"bools\": [true, false], \"withNull\": [1, null, 3], \"bigints\": [1,
9999999999]}")
Review Comment:
JSON_QUERY does not support MAP?
##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/JsonFunctionsITCase.java:
##########
@@ -886,7 +1083,214 @@ private static List<TestSetSpec> jsonQuerySpec() {
.testTableApiRuntimeError(
$("f0").jsonQuery("strict $.err10",
WITHOUT_ARRAY, NULL, ERROR),
TableRuntimeException.class,
- "No results for path"));
+ "No results for path"),
+
+ // Typed RETURNING ARRAY<T> support
+ TestSetSpec.forFunction(BuiltInFunctionDefinitions.JSON_QUERY)
+ .onFieldsWithData(
+ "{\"ints\": [1, 2, 3], \"doubles\": [1.5,
2.5], \"bools\": [true, false], \"withNull\": [1, null, 3], \"bigints\": [1,
9999999999]}")
+ .andDataTypes(STRING())
+ .testResult(
+ $("f0").jsonQuery("$.ints", ARRAY(INT())),
+ "JSON_QUERY(f0, '$.ints' RETURNING
ARRAY<INT>)",
+ new Integer[] {1, 2, 3},
+ ARRAY(INT()))
+ .testResult(
+ $("f0").jsonQuery("$.ints", ARRAY(TINYINT())),
+ "JSON_QUERY(f0, '$.ints' RETURNING
ARRAY<TINYINT>)",
+ new Byte[] {1, 2, 3},
+ ARRAY(TINYINT()))
+ .testResult(
+ $("f0").jsonQuery("$.ints", ARRAY(SMALLINT())),
+ "JSON_QUERY(f0, '$.ints' RETURNING
ARRAY<SMALLINT>)",
+ new Short[] {1, 2, 3},
+ ARRAY(SMALLINT()))
+ .testResult(
+ $("f0").jsonQuery("$.ints", ARRAY(BIGINT())),
+ "JSON_QUERY(f0, '$.ints' RETURNING
ARRAY<BIGINT>)",
+ new Long[] {1L, 2L, 3L},
+ ARRAY(BIGINT()))
+ .testResult(
+ $("f0").jsonQuery("$.bigints",
ARRAY(BIGINT())),
+ "JSON_QUERY(f0, '$.bigints' RETURNING
ARRAY<BIGINT>)",
+ new Long[] {1L, 9999999999L},
+ ARRAY(BIGINT()))
+ .testResult(
+ $("f0").jsonQuery("$.doubles",
ARRAY(DOUBLE())),
+ "JSON_QUERY(f0, '$.doubles' RETURNING
ARRAY<DOUBLE>)",
+ new Double[] {1.5, 2.5},
+ ARRAY(DOUBLE()))
+ .testResult(
+ $("f0").jsonQuery("$.doubles", ARRAY(FLOAT())),
+ "JSON_QUERY(f0, '$.doubles' RETURNING
ARRAY<FLOAT>)",
+ new Float[] {1.5f, 2.5f},
+ ARRAY(FLOAT()))
+ .testResult(
+ $("f0").jsonQuery("$.bools", ARRAY(BOOLEAN())),
+ "JSON_QUERY(f0, '$.bools' RETURNING
ARRAY<BOOLEAN>)",
+ new Boolean[] {true, false},
+ ARRAY(BOOLEAN()))
+ .testResult(
+ $("f0").jsonQuery("$.withNull", ARRAY(INT())),
+ "JSON_QUERY(f0, '$.withNull' RETURNING
ARRAY<INT>)",
+ new Integer[] {1, null, 3},
+ ARRAY(INT()))
Review Comment:
Add a test with `ARRAY<INT NOT NULL>`
##########
flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/functions/SqlJsonUtilsConversionTest.java:
##########
@@ -0,0 +1,396 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.flink.table.runtime.functions;
+
+import org.apache.flink.table.api.JsonValueOnEmptyOrError;
+import org.apache.flink.table.api.TableRuntimeException;
+import org.apache.flink.table.data.DecimalData;
+import org.apache.flink.table.data.GenericArrayData;
+import org.apache.flink.table.types.logical.LogicalTypeRoot;
+
+import org.junit.jupiter.api.Nested;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.Arguments;
+import org.junit.jupiter.params.provider.MethodSource;
+
+import java.math.BigDecimal;
+import java.math.BigInteger;
+import java.util.stream.Stream;
+
+import static org.apache.flink.table.api.JsonQueryOnEmptyOrError.EMPTY_ARRAY;
+import static org.apache.flink.table.api.JsonQueryOnEmptyOrError.ERROR;
+import static org.apache.flink.table.api.JsonQueryOnEmptyOrError.NULL;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.BIGINT;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.BOOLEAN;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.DECIMAL;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.DOUBLE;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.FLOAT;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.INTEGER;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.SMALLINT;
+import static org.apache.flink.table.types.logical.LogicalTypeRoot.TINYINT;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+class SqlJsonUtilsConversionTest {
+
+ private static Object scalar(Object raw, LogicalTypeRoot type) {
+ return SqlJsonUtils.convertJsonScalar(raw, type, 0, 0,
JsonValueOnEmptyOrError.NULL, null);
+ }
+
+ private static Object scalar(
+ Object raw,
+ LogicalTypeRoot type,
+ JsonValueOnEmptyOrError onError,
+ Object defaultValue) {
+ return SqlJsonUtils.convertJsonScalar(raw, type, 0, 0, onError,
defaultValue);
+ }
+
+ private static Object scalarDecimal(Object raw, int precision, int scale) {
+ return SqlJsonUtils.convertJsonScalar(
+ raw, DECIMAL, precision, scale, JsonValueOnEmptyOrError.NULL,
null);
+ }
+
+ private static GenericArrayData array(Object[] raw, LogicalTypeRoot
elementType) {
+ return SqlJsonUtils.convertJsonArray(raw, elementType, 0, 0, NULL);
+ }
+
+ private static GenericArrayData arrayDecimal(
+ Object[] raw,
+ int precision,
+ int scale,
+ org.apache.flink.table.api.JsonQueryOnEmptyOrError onError) {
+ return SqlJsonUtils.convertJsonArray(raw, DECIMAL, precision, scale,
onError);
+ }
+
+ static Stream<Arguments> validConversions() {
+ return Stream.of(
+ Arguments.of(42, TINYINT, (byte) 42),
+ Arguments.of(127L, TINYINT, Byte.MAX_VALUE),
+ Arguments.of(-128L, TINYINT, Byte.MIN_VALUE),
+ Arguments.of(1000, SMALLINT, (short) 1000),
+ Arguments.of((long) Short.MAX_VALUE, SMALLINT,
Short.MAX_VALUE),
+ Arguments.of(42L, INTEGER, 42),
+ Arguments.of((long) Integer.MAX_VALUE, INTEGER,
Integer.MAX_VALUE),
+ Arguments.of(42, BIGINT, 42L),
+ Arguments.of(9_999_999_999L, BIGINT, 9_999_999_999L),
+ Arguments.of(13.37, FLOAT, 13.37f),
+ Arguments.of(13.37, DOUBLE, 13.37));
+ }
+
+ static Stream<Arguments> overflowCases() {
+ return Stream.of(
+ Arguments.of(128, TINYINT),
+ Arguments.of(-129, TINYINT),
+ Arguments.of(200L, TINYINT),
+ Arguments.of(-200L, TINYINT),
+ Arguments.of((long) Short.MAX_VALUE + 1, SMALLINT),
+ Arguments.of((long) Short.MIN_VALUE - 1, SMALLINT),
+ Arguments.of(40_000L, SMALLINT),
+ Arguments.of(-40_000L, SMALLINT),
+ Arguments.of(9_999_999_999L, INTEGER),
+ Arguments.of(-9_999_999_999L, INTEGER));
+ }
+
+ @Nested
+ class RangeCheckedConversions {
+
+ @ParameterizedTest(name = "{0} -> {1}")
+ @MethodSource(
+
"org.apache.flink.table.runtime.functions.SqlJsonUtilsConversionTest#validConversions")
+ void convertsWithinRange(Number input, LogicalTypeRoot type, Object
expected) {
+ assertThat(scalar(input, type)).isEqualTo(expected);
+ }
+
+ @ParameterizedTest(name = "{0} -> {1}")
+ @MethodSource(
+
"org.apache.flink.table.runtime.functions.SqlJsonUtilsConversionTest#overflowCases")
Review Comment:
Can't we just use a `CsvSource`? It is easer to follow and read
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]