diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/JsonAggregationFunctionsITCase.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/JsonAggregationFunctionsITCase.java index 34b00813d759d5..e0baa5012bcddc 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/JsonAggregationFunctionsITCase.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/JsonAggregationFunctionsITCase.java @@ -157,6 +157,18 @@ public Stream getTestCaseSpecs() { Arrays.asList( Row.of(1, "{\"A\":1,\"B\":3}", 3), Row.of(2, "{\"A\":2,\"C\":5}", 5))), + TestSpec.forFunction(BuiltInFunctionDefinitions.JSON_OBJECTAGG_NULL_ON_NULL) + .withDescription("Trailing Garbage After A Valid JSON Value") + .withSource( + ROW(STRING(), STRING()), + Collections.singletonList(Row.ofKind(INSERT, "A", "{\"a\":1} x"))) + .testResult( + source -> "SELECT JSON_OBJECTAGG(f0 VALUE f1) FROM " + source, + TableApiAggSpec.select( + jsonObjectAgg(JsonOnNull.NULL, $("f0"), $("f1"))), + ROW(VARCHAR(2000).notNull()), + ROW(STRING().notNull()), + Collections.singletonList(Row.of("{\"A\":\"{\\\"a\\\":1} x\"}"))), // JSON_ARRAYAGG TestSpec.forFunction(BuiltInFunctionDefinitions.JSON_ARRAYAGG_ABSENT_ON_NULL) @@ -240,6 +252,17 @@ public Stream getTestCaseSpecs() { ROW(INT(), STRING(), STRING().notNull()), Arrays.asList( Row.of(1, "A", "[\"A\"]"), - Row.of(2, "D", "[\"C\",\"D\"]")))); + Row.of(2, "D", "[\"C\",\"D\"]"))), + TestSpec.forFunction(BuiltInFunctionDefinitions.JSON_ARRAYAGG_ABSENT_ON_NULL) + .withDescription("Trailing Garbage After A Valid JSON Value") + .withSource( + ROW(STRING()), + Collections.singletonList(Row.ofKind(INSERT, "{\"a\":1} x"))) + .testResult( + source -> "SELECT JSON_ARRAYAGG(f0) FROM " + source, + TableApiAggSpec.select(jsonArrayAgg(JsonOnNull.ABSENT, $("f0"))), + ROW(VARCHAR(2000).notNull()), + ROW(STRING().notNull()), + Collections.singletonList(Row.of("[\"{\\\"a\\\":1} x\"]")))); } } diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/JsonFunctionsITCase.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/JsonFunctionsITCase.java index 89a0448e95eef5..57bdfb77a73b62 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/JsonFunctionsITCase.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/JsonFunctionsITCase.java @@ -100,8 +100,8 @@ Stream getTestSetSpecs() { private static TestSetSpec jsonExistsSpec() { final String jsonValue = getJsonFromResource("/json/json-exists.json"); return TestSetSpec.forFunction(BuiltInFunctionDefinitions.JSON_EXISTS) - .onFieldsWithData(jsonValue) - .andDataTypes(STRING()) + .onFieldsWithData(jsonValue, "{\"a\":1} x") + .andDataTypes(STRING(), STRING()) // NULL .testResult( @@ -174,14 +174,17 @@ private static TestSetSpec jsonExistsSpec() { .testTableApiRuntimeError( $("f0").jsonExists("strict $.invalid", JsonExistsOnError.ERROR), TableRuntimeException.class, - "No results for path: $['invalid']"); + "No results for path: $['invalid']") + + // Trailing garbage after a valid JSON value + .testResult($("f1").jsonExists("$.a"), "JSON_EXISTS(f1, '$.a')", false, BOOLEAN()); } private static TestSetSpec jsonValueSpec() { final String jsonValue = getJsonFromResource("/json/json-value.json"); return TestSetSpec.forFunction(BuiltInFunctionDefinitions.JSON_VALUE) - .onFieldsWithData(jsonValue) - .andDataTypes(STRING()) + .onFieldsWithData(jsonValue, "{\"a\":1} x") + .andDataTypes(STRING(), STRING()) // NULL and invalid types .testResult( @@ -294,6 +297,9 @@ private static TestSetSpec jsonValueSpec() { "right", STRING()) + // Trailing garbage after a valid JSON value + .testResult($("f1").jsonValue("$.a"), "JSON_VALUE(f1, '$.a')", null, STRING()) + // Multiple JSON_VALUE calls on the same input should reuse parsed JSON .testSqlResult( "JSON_VALUE(f0, '$.type'), JSON_VALUE(f0, '$.age')", @@ -368,7 +374,12 @@ private static List isJsonSpec() { $("f0").isJson(JsonType.OBJECT), "f0 IS JSON OBJECT", true, - BOOLEAN().notNull())); + BOOLEAN().notNull()), + // Trailing garbage after a valid JSON value + TestSetSpec.forFunction(BuiltInFunctionDefinitions.IS_JSON) + .onFieldsWithData("{\"a\":1} x") + .andDataTypes(STRING()) + .testResult($("f0").isJson(), "f0 IS JSON", false, BOOLEAN().notNull())); } private static List jsonQuerySpec() { @@ -593,7 +604,14 @@ private static List jsonQuerySpec() { .testTableApiRuntimeError( $("f0").jsonQuery("strict $.err10", WITHOUT_ARRAY, NULL, ERROR), TableRuntimeException.class, - "No results for path")); + "No results for path"), + // Trailing garbage after a valid JSON value. The path must select an object, + // because a scalar would return NULL even if the trailing garbage was accepted. + TestSetSpec.forFunction(BuiltInFunctionDefinitions.JSON_QUERY) + .onFieldsWithData("{\"a\":{\"b\":1}} x") + .andDataTypes(STRING()) + .testResult( + $("f0").jsonQuery("$.a"), "JSON_QUERY(f0, '$.a')", null, STRING())); } private static List jsonStringSpec() { @@ -719,7 +737,17 @@ private static List jsonStringSpec() { jsonString($("f0")), "JSON_STRING(f0)", "{\"field\\ttab\":\"val4\",\"field\\nline\":\"val3\",\"field\\rreturn\":\"val5\",\"field\\\"quote\":\"val1\",\"field\\\\slash\":\"val2\"}", - STRING().notNull())); + STRING().notNull()), + // Trailing garbage after a valid JSON value. JSON_STRING serializes its argument + // instead of parsing it, so the garbage is kept as part of the string. + TestSetSpec.forFunction(BuiltInFunctionDefinitions.JSON_STRING) + .onFieldsWithData("{\"a\":1} x") + .andDataTypes(STRING()) + .testResult( + jsonString($("f0")), + "JSON_STRING(f0)", + "\"{\\\"a\\\":1} x\"", + STRING())); } private static List jsonSpec() { @@ -1106,7 +1134,19 @@ private static List jsonObjectSpec() { jsonObject(JsonOnNull.NULL, "testRow", $("f0")), "JSON_OBJECT(KEY 'testRow' VALUE f0 NULL ON NULL)", "{\"testRow\":{\"field\\ttab\":\"val4\",\"field\\nline\":\"val3\",\"field\\rreturn\":\"val5\",\"field\\\"quote\":\"val1\",\"field\\\\slash\":\"val2\"}}", - STRING().notNull())); + STRING().notNull()), + // Trailing garbage after a valid JSON value + TestSetSpec.forFunction(BuiltInFunctionDefinitions.JSON_OBJECT) + .onFieldsWithData("{\"a\":1} x") + .andDataTypes(STRING()) + .testSqlRuntimeError( + "JSON_OBJECT(KEY 'K' VALUE JSON(f0))", + TableRuntimeException.class, + "Invalid JSON string in JSON(value) function") + .testTableApiRuntimeError( + jsonObject(JsonOnNull.NULL, "K", json($("f0"))), + TableRuntimeException.class, + "Invalid JSON string in JSON(value) function")); } private static List jsonQuoteSpec() { @@ -1178,7 +1218,17 @@ private static List jsonQuoteSpec() { "\"\\u2260 will be escaped\"", STRING().notNull()) .testResult( - $("f8").jsonQuote(), "JSON_QUOTE(f8)", null, STRING().nullable())); + $("f8").jsonQuote(), "JSON_QUOTE(f8)", null, STRING().nullable()), + // Trailing garbage after a valid JSON value. JSON_QUOTE only escapes its input, so + // the garbage is preserved. + TestSetSpec.forFunction(BuiltInFunctionDefinitions.JSON_QUOTE) + .onFieldsWithData("{\"a\":1} x") + .andDataTypes(STRING()) + .testResult( + $("f0").jsonQuote(), + "JSON_QUOTE(f0)", + "\"{\\\"a\\\":1} x\"", + STRING())); } private static List jsonUnquoteSpecWithValidInput() { @@ -1194,13 +1244,13 @@ private static List jsonUnquoteSpecWithValidInput() { TestSetSpec.forFunction(BuiltInFunctionDefinitions.JSON_UNQUOTE) .onFieldsWithData( "\"abc\"", - "\"[\"abc\"]\"", - "\"[\"\\u0041\"]\"", + "\"[\\\"abc\\\"]\"", + "\"[\\\"\\u0041\\\"]\"", "\"\\u0041\"", - "\"[\"\\t\\u0032\"]\"", - "\"[\"This is a \\t test \\n with special characters: \\b \\f \\r \\u0041\"]\"", - "\"\"\"", - "\"\"\ufffa\"", + "\"[\\\"\\t\\u0032\\\"]\"", + "\"[\\\"This is a \\t test \\n with special characters: \\b \\f \\r \\u0041\\\"]\"", + "\"\\\"\"", + "\"\\\"\ufffa\"", "\"a unicode \u2260\"", "\"valid unicode literal \\uD801\\uDC00\"", "[1,2,3]", @@ -1400,7 +1450,17 @@ private static List jsonUnquoteSpecWithInvalidInput() { $("f10").jsonQuote(), "JSON_UNQUOTE(f10)", null, - STRING().nullable())); + STRING().nullable()), + // Trailing garbage after a valid JSON value. The leading '"abc"' must not be + // unquoted on its own, the input is passed through unchanged. + TestSetSpec.forFunction(BuiltInFunctionDefinitions.JSON_UNQUOTE) + .onFieldsWithData("\"abc\" \"def\"") + .andDataTypes(STRING()) + .testResult( + $("f0").jsonUnquote(), + "JSON_UNQUOTE(f0)", + "\"abc\" \"def\"", + STRING())); } private static List jsonArraySpec() { @@ -1543,7 +1603,19 @@ private static List jsonArraySpec() { jsonArray(JsonOnNull.NULL, $("f0")), "JSON_ARRAY(f0 NULL ON NULL)", "[{\"field\\ttab\":\"val4\",\"field\\nline\":\"val3\",\"field\\rreturn\":\"val5\",\"field\\\"quote\":\"val1\",\"field\\\\slash\":\"val2\"}]", - STRING().notNull())); + STRING().notNull()), + // Trailing garbage after a valid JSON value + TestSetSpec.forFunction(BuiltInFunctionDefinitions.JSON_ARRAY) + .onFieldsWithData("{\"a\":1} x") + .andDataTypes(STRING()) + .testSqlRuntimeError( + "JSON_ARRAY(JSON(f0))", + TableRuntimeException.class, + "Invalid JSON string in JSON(value) function") + .testTableApiRuntimeError( + jsonArray(JsonOnNull.NULL, json($("f0"))), + TableRuntimeException.class, + "Invalid JSON string in JSON(value) function")); } /** Pins the local-ref / common-sub-expression handling for JSON construction calls. */ diff --git a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/expressions/ScalarFunctionsTest.scala b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/expressions/ScalarFunctionsTest.scala index f92847ee3a6033..95deca0584423b 100644 --- a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/expressions/ScalarFunctionsTest.scala +++ b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/expressions/ScalarFunctionsTest.scala @@ -724,9 +724,7 @@ class ScalarFunctionsTest extends ScalarTypesTestBase { "This is a \t test \n with special characters: \b \f \r A" ) testSqlApi("JSON_UNQUOTE('\"\"')", "") - testSqlApi("JSON_UNQUOTE('\"\"\"')", "\"") testSqlApi("JSON_UNQUOTE('[]')", "[]") - testSqlApi("JSON_UNQUOTE('\"\"\\ufffa\"')", "\"\ufffa") testSqlApi("JSON_UNQUOTE('{\"key\":1}')", "{\"key\":1}") testSqlApi("JSON_UNQUOTE('true')", "true") } @@ -747,6 +745,10 @@ class ScalarFunctionsTest extends ScalarTypesTestBase { def testJsonUnquoteWithInvalidInput(): Unit = { testSqlApi("JSON_UNQUOTE('\"[1, 2, 3}')", "\"[1, 2, 3}") testSqlApi("JSON_UNQUOTE('\"')", "\"") + // Quoted, but not a valid JSON string literal: the leading '""' is followed by a + // trailing token, so the input is passed through instead of being unquoted. + testSqlApi("JSON_UNQUOTE('\"\"\"')", "\"\"\"") + testSqlApi("JSON_UNQUOTE('\"\"\\ufffa\"')", "\"\"\\ufffa\"") testSqlApi("JSON_UNQUOTE('[}')", "[}") testSqlApi("JSON_UNQUOTE('1\"')", "1\"") testSqlApi("JSON_UNQUOTE('[')", "[") diff --git a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/functions/SqlJsonUtils.java b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/functions/SqlJsonUtils.java index c73e9981339965..f759311b3f759e 100644 --- a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/functions/SqlJsonUtils.java +++ b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/functions/SqlJsonUtils.java @@ -68,7 +68,8 @@ public class SqlJsonUtils { private static final ObjectMapper MAPPER = new ObjectMapper(JSON_FACTORY) .configure(SerializationFeature.ORDER_MAP_ENTRIES_BY_KEYS, true) - .configure(DeserializationFeature.USE_BIG_DECIMAL_FOR_FLOATS, true); + .configure(DeserializationFeature.USE_BIG_DECIMAL_FOR_FLOATS, true) + .configure(DeserializationFeature.FAIL_ON_TRAILING_TOKENS, true); private static final Pattern JSON_PATH_BASE = Pattern.compile( "^\\s*(?strict|lax)\\s+(?.+)$",