MartijnVisser commented on code in PR #29063:
URL: https://github.com/apache/flink/pull/29063#discussion_r4145042958
##########
flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/calls/JsonValueCallGen.scala:
##########
@@ -82,26 +83,48 @@ class JsonValueCallGen extends CallGenerator {
)
val rawResultTerm = CodeGenUtils.newName(ctx, "rawResult")
- val call = s"""
- |Object $rawResultTerm =
- |
${qualifyMethod(BuiltInMethods.JSON_VALUE)}(${terms
- .mkString(", ")});
+ val baseCall = s"""
+ |Object $rawResultTerm =
+ |
${qualifyMethod(BuiltInMethods.JSON_VALUE)}(${terms
+ .mkString(", ")});
""".stripMargin
- val convertedResult = returnType.getTypeRoot match {
+ returnType.getTypeRoot match {
case LogicalTypeRoot.VARCHAR =>
-
s"$BINARY_STRING.fromString(java.lang.String.valueOf($rawResultTerm))"
- case LogicalTypeRoot.BOOLEAN => s"(java.lang.Boolean)
$rawResultTerm"
- case LogicalTypeRoot.INTEGER => s"(java.lang.Integer)
$rawResultTerm"
- case LogicalTypeRoot.DOUBLE => s"((java.math.BigDecimal)
$rawResultTerm).doubleValue()"
- case _ =>
+ val result =
+ s"($rawResultTerm == null) ? null : " +
+
s"($BINARY_STRING.fromString(java.lang.String.valueOf($rawResultTerm)))"
+ (baseCall, result)
+ case typeRoot if
!SqlJsonUtils.isSupportedJsonReturningType(typeRoot) =>
throw new CodeGenException(
- s"Unsupported type '$returnType' "
- + "for RETURNING in JSON_VALUE().")
- }
+ s"Unsupported type '$returnType' for RETURNING in
JSON_VALUE().")
+ case _ =>
+ val jsonUtils =
+ "org.apache.flink.table.runtime.functions.SqlJsonUtils"
+ val typeRootEnum =
+
s"org.apache.flink.table.types.logical.LogicalTypeRoot.${returnType.getTypeRoot.name()}"
+ val (precisionStr, scaleStr) = returnType.getTypeRoot match {
+ case LogicalTypeRoot.DECIMAL =>
+ val dt = returnType.asInstanceOf[DecimalType]
+ (dt.getPrecision.toString, dt.getScale.toString)
+ case _ => ("0", "0")
+ }
+ val errorBehaviorEnum =
+
qualifyEnum(errorBehavior._1.asInstanceOf[JsonValueOnEmptyOrError])
+ val defaultExpr =
+ if (errorBehavior._2 != null) errorBehavior._2 else "null"
- val result = s"($rawResultTerm == null) ? null : ($convertedResult)"
- (call, result)
+ val boxedType = CodeGenUtils.boxedTypeTermForType(returnType)
Review Comment:
The boxed result also breaks conditions: `IF(JSON_VALUE(f0, '$.missing'
RETURNING BOOLEAN), 1, 2)` throws an NPE on master and on this branch. Unboxing
the result in this generator fixes that too.
##########
flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/calls/ScalarOperatorGens.scala:
##########
@@ -624,7 +624,14 @@ object ScalarOperatorGens {
isNumeric(left.resultType) && isNumeric(right.resultType)
|| isTime(left.resultType) &&
isTime(right.resultType)
- ) { (leftTerm, rightTerm) => s"$leftTerm $operator $rightTerm" }
+ ) {
+ (leftTerm, rightTerm) =>
+ {
+ val l = s"((${primitiveTypeTermForType(left.resultType)})
$leftTerm)"
Review Comment:
This touches every comparison and cast in every query. With the unboxing in
`JsonValueCallGen` instead, all 890 cases pass without changes to this file.
##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/functions/SqlJsonUtils.java:
##########
@@ -815,6 +830,274 @@ public String toString() {
}
}
+ // --- Type conversion for JSON_VALUE / JSON_QUERY RETURNING ---
+
+ private static final java.util.EnumSet<LogicalTypeRoot>
SUPPORTED_JSON_RETURNING_TYPES =
Review Comment:
The docs still list only VARCHAR, BOOLEAN, INTEGER and DOUBLE for
`JSON_VALUE`, and only `ARRAY<STRING>` for `JSON_QUERY` (`sql_functions.yml`
and the zh file, `BaseExpressions`, PyFlink's `json_value`). Please update them
in this PR.
##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/functions/SqlJsonUtils.java:
##########
@@ -815,6 +830,274 @@ 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,
+ boolean elementNullable,
+ 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);
+ } else if (!elementNullable) {
+ throw new JsonConversionException(
+ "Null element at index " + i + " in NOT NULL
array");
+ }
+ }
+ 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 {
+ String trimmed = ((String) raw).trim();
+ return convertToType(new BigDecimal(trimmed), typeRoot,
precision, scale);
+ } catch (NumberFormatException e) {
+ throw new JsonConversionException(
+ "Cannot parse string '" + raw + "' as " + typeRoot, e);
+ }
+ }
+ try {
+ switch (typeRoot) {
+ case BOOLEAN:
+ if (raw instanceof Number) {
+ double d = ((Number) raw).doubleValue();
+ if (d == 0.0) {
+ return false;
+ }
+ if (d == 1.0) {
+ return true;
+ }
+ throw new JsonConversionException("Cannot convert " +
raw + " to 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(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);
+ }
+ }
+
+ private static @Nullable BigInteger toBigIntegerTruncated(Number n) {
+ if (n instanceof BigDecimal) {
+ return ((BigDecimal) n).toBigInteger();
+ }
+ if (n instanceof BigInteger) {
+ return (BigInteger) n;
+ }
+ return null;
+ }
+
+ static byte toCheckedByte(Number n) {
+ BigInteger bi = toBigIntegerTruncated(n);
+ if (bi != null) {
+ if (bi.compareTo(BigInteger.valueOf(Byte.MAX_VALUE)) > 0
+ || bi.compareTo(BigInteger.valueOf(Byte.MIN_VALUE)) < 0) {
+ throw new JsonConversionException("Value " + n + " is out of
range for TINYINT");
+ }
+ return bi.byteValue();
+ }
+ 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) {
+ BigInteger bi = toBigIntegerTruncated(n);
+ if (bi != null) {
+ if (bi.compareTo(BigInteger.valueOf(Short.MAX_VALUE)) > 0
+ || bi.compareTo(BigInteger.valueOf(Short.MIN_VALUE)) < 0) {
+ throw new JsonConversionException("Value " + n + " is out of
range for SMALLINT");
+ }
+ return bi.shortValue();
+ }
+ 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) {
+ BigInteger bi = toBigIntegerTruncated(n);
+ if (bi != null) {
+ if (bi.compareTo(BigInteger.valueOf(Integer.MAX_VALUE)) > 0
+ || bi.compareTo(BigInteger.valueOf(Integer.MIN_VALUE)) <
0) {
+ throw new JsonConversionException("Value " + n + " is out of
range for INTEGER");
+ }
+ return bi.intValue();
+ }
+ long v = n.longValue();
+ if (v < Integer.MIN_VALUE || v > Integer.MAX_VALUE) {
+ throw new JsonConversionException("Value " + n + " is out of range
for INTEGER");
+ }
+ return (int) v;
+ }
+
+ static boolean parseStringAsBoolean(String s) {
Review Comment:
Flink's `CAST(... AS BOOLEAN)` also accepts `y` and `n`, and so does
PostgreSQL. Here `"y"` gives NULL while casting the same value gives TRUE.
Could this reuse `BinaryStringDataUtil.toBoolean`?
##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/functions/SqlJsonUtils.java:
##########
@@ -815,6 +830,274 @@ 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,
+ boolean elementNullable,
+ 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);
+ } else if (!elementNullable) {
+ throw new JsonConversionException(
+ "Null element at index " + i + " in NOT NULL
array");
+ }
+ }
+ 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) {
Review Comment:
A JSON string with `RETURNING INTEGER`, or a number with `BOOLEAN`, used to
fail the job and now converts or goes to `ON ERROR`. Can you suggest a release
note for the Jira?
##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/functions/SqlJsonUtils.java:
##########
@@ -815,6 +830,274 @@ 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,
+ boolean elementNullable,
+ 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);
+ } else if (!elementNullable) {
+ throw new JsonConversionException(
+ "Null element at index " + i + " in NOT NULL
array");
+ }
+ }
+ 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 {
+ String trimmed = ((String) raw).trim();
+ return convertToType(new BigDecimal(trimmed), typeRoot,
precision, scale);
+ } catch (NumberFormatException e) {
+ throw new JsonConversionException(
+ "Cannot parse string '" + raw + "' as " + typeRoot, e);
+ }
+ }
+ try {
+ switch (typeRoot) {
+ case BOOLEAN:
+ if (raw instanceof Number) {
+ double d = ((Number) raw).doubleValue();
Review Comment:
`doubleValue()` rounds, so `1.00000000000000000001` returns TRUE and
`1e-400` returns FALSE. Comparing as `BigDecimal` keeps this exact.
--
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]