raminqaf commented on code in PR #29341:
URL: https://github.com/apache/flink/pull/29341#discussion_r4205081771
##########
docs/content.zh/docs/sql/reference/data-types.md:
##########
@@ -1545,6 +1546,137 @@ cases `TRY_PARSE_JSON` returns `NULL`.
A `VARIANT` has no dedicated kind for these values. To keep one, store it as a
JSON string and cast
it back out, for example `CAST(CAST(PARSE_JSON('"Infinity"') AS STRING) AS
FLOAT)`.
+The `PARSE_XML` function maps an XML document to an object with a single field
named after the root
+element. For example:
+
+```sql
+--
{"book":{"#":{"price":[1,2],"title":0},"@pages":"320","price":["12.50","13.50"],"title":"Dune"}}
+PARSE_XML('<book
pages="320"><title>Dune</title><price>12.50</price><price>13.50</price></book>')
+```
+
+Attributes become fields with an `@` prefix, and a child element that occurs
more than once becomes
+an array. The `#` field records the order of the children, as described below.
It can be ignored
+when reading fields.
+
+Values are strings, unless the document declares a type with `xsi:type`, see
below. Since a cast
+from a `VARIANT` never parses a string, cast a value to `STRING` first and
convert it with a regular
+cast:
+
+```sql
+CAST(v['book']['title'] AS STRING) -- 'Dune'
+CAST(v['book']['price'][2] AS STRING) -- '13.50'
+CAST(CAST(v['book']['@pages'] AS STRING) AS INT) -- 320
+```
+
+Each element is mapped by the following rules:
+
+| XML input | Stored
`VARIANT` value |
+|---------------------------------------------------------------------|----------------------------------------------|
Review Comment:
It is nice that we have detailed documentation, but I was wondering if the
PARSE_XML should be explained in `data-types`.
##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/XmlFunctionsITCase.java:
##########
@@ -0,0 +1,207 @@
+/*
+ * 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.planner.functions;
+
+import org.apache.flink.table.api.TableRuntimeException;
+import org.apache.flink.table.functions.BuiltInFunctionDefinitions;
+
+import java.math.BigDecimal;
+import java.util.stream.Stream;
+
+import static org.apache.flink.table.api.DataTypes.BOOLEAN;
+import static org.apache.flink.table.api.DataTypes.DECIMAL;
+import static org.apache.flink.table.api.DataTypes.INT;
+import static org.apache.flink.table.api.DataTypes.STRING;
+import static org.apache.flink.table.api.Expressions.$;
+import static org.apache.flink.table.api.Expressions.call;
+import static org.apache.flink.table.api.Expressions.jsonString;
+import static org.apache.flink.table.api.Expressions.lit;
+import static org.apache.flink.table.api.Expressions.nullOf;
+
+/**
+ * Tests for {@link BuiltInFunctionDefinitions#PARSE_XML} and {@link
+ * BuiltInFunctionDefinitions#TRY_PARSE_XML}.
+ *
+ * <p>The bulk of the XML mapping is covered by {@code XmlToVariantParserTest}.
+ */
+class XmlFunctionsITCase extends BuiltInFunctionTestBase {
+
+ private static final String BOOK =
+ "<book pages=\"320\">\n"
+ + " <title>Dune</title>\n"
+ + " <price xsi:type=\"decimal\">12.50</price>\n"
+ + "</book>";
+
+ private static final String BOOK_JSON =
+ "{\"book\":{\"#\":{\"price\":1,\"title\":0},"
+ + "\"@pages\":\"320\",\"price\":12.5,\"title\":\"Dune\"}}";
+
+ private static final String BOOK_FORCE_ARRAY_JSON =
+ "{\"book\":{\"#\":{\"price\":[1],\"title\":[0]},"
+ +
"\"@pages\":\"320\",\"price\":[12.5],\"title\":[\"Dune\"]}}";
+
+ private static final String INVALID = "<book>";
+
+ private static final String EXTERNAL_ENTITY =
+ "<!DOCTYPE a [<!ENTITY e SYSTEM
\"file:///etc/hosts\">]><a>&e;</a>";
+
+ @Override
+ Stream<TestSetSpec> getTestSetSpecs() {
+ return Stream.of(parseXmlSpec(), tryParseXmlSpec(),
constantFoldingSpec());
+ }
+
+ private static TestSetSpec parseXmlSpec() {
+ return TestSetSpec.forFunction(BuiltInFunctionDefinitions.PARSE_XML)
+ .onFieldsWithData(BOOK, INVALID, EXTERNAL_ENTITY, true, null)
+ .andDataTypes(
+ STRING().notNull(),
+ STRING().notNull(),
+ STRING().notNull(),
+ BOOLEAN().notNull(),
+ BOOLEAN())
+ .testResult(
+ jsonString($("f0").parseXml()),
+ "JSON_STRING(PARSE_XML(f0))",
+ BOOK_JSON,
+ STRING().notNull())
+ .testResult(
+ jsonString($("f0").parseXml(false)),
+ "JSON_STRING(PARSE_XML(f0, FALSE))",
+ BOOK_JSON,
+ STRING().notNull())
+ .testResult(
+ jsonString($("f0").parseXml(true)),
+ "JSON_STRING(PARSE_XML(f0, TRUE))",
+ BOOK_FORCE_ARRAY_JSON,
+ STRING().notNull())
+ .testResult(
+ jsonString(call("PARSE_XML", $("f0"), $("f3"))),
Review Comment:
Can we re-write this?
```suggestion
jsonString($("f0").parseXml($("f3"))),
```
Also it makes it a bit easier for the tests to read if you just pass the
`true` column directly as a literal. Same for the `null` column
##########
flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/functions/XmlToVariantParserTest.java:
##########
@@ -0,0 +1,534 @@
+/*
+ * 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.types.variant.Variant;
+
+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 org.junit.jupiter.params.provider.ValueSource;
+import org.xml.sax.SAXException;
+
+import javax.xml.transform.stream.StreamSource;
+import javax.xml.validation.SchemaFactory;
+
+import java.io.StringReader;
+import java.math.BigDecimal;
+import java.time.Instant;
+import java.time.LocalDate;
+import java.time.LocalDateTime;
+import java.time.LocalTime;
+import java.util.stream.Stream;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatNoException;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import static org.junit.jupiter.params.provider.Arguments.arguments;
+
+/** Tests for {@link XmlToVariantParser}. */
+class XmlToVariantParserTest {
+
+ private final XmlToVariantParser parser = new XmlToVariantParser();
+
+ //
--------------------------------------------------------------------------------------------
+ // Mapping
+ //
--------------------------------------------------------------------------------------------
+
+ @ParameterizedTest(name = "{0}")
+ @MethodSource("mappings")
+ void testMapping(String xml, String expectedJson) {
+ assertThat(parser.parse(xml, false).toJson()).isEqualTo(expectedJson);
+ }
+
+ static Stream<Arguments> mappings() {
+ return Stream.of(
+ arguments("<title>Dune</title>", "{\"title\":\"Dune\"}"),
+ arguments("<title></title>", "{\"title\":\"\"}"),
+ arguments("<title/>", "{\"title\":\"\"}"),
+ arguments(
+ "<book pages=\"320\">\n Dune\n</book>",
+ "{\"book\":{\"$\":\"Dune\",\"@pages\":\"320\"}}"),
+ arguments("<book pages=\"320\"/>",
"{\"book\":{\"@pages\":\"320\"}}"),
+ arguments(
+ "<book>\n"
+ + " <author>Terry Pratchett</author>\n"
+ + " <author>Neil Gaiman</author>\n"
+ + "</book>",
+ "{\"book\":{\"author\":[\"Terry Pratchett\",\"Neil
Gaiman\"]}}"),
+ arguments(
+ "<book>\n <title>Dune</title>\n
<price>12.5</price>\n</book>",
+
"{\"book\":{\"#\":{\"price\":1,\"title\":0},\"price\":\"12.5\",\"title\":\"Dune\"}}"),
+ arguments(
+ "<book>\n"
+ + " This is a description\n"
+ + " <title>Dune</title>\n"
+ + " and this is another text\n"
+ + "</book>",
+ "{\"book\":{\"#\":{\"$\":[0,2],\"title\":1},"
+ + "\"$\":[\"This is a description\",\"and this
is another text\"],"
+ + "\"title\":\"Dune\"}}"),
+ arguments(
+ "<book pages=\"320\">\n"
+ + " <title>Dune</title>\n"
+ + " This is a description\n"
+ + " <price>12.50</price>\n"
+ + " Text in between\n"
+ + " <price currency=\"EUR\">13.50</price>\n"
+ + "</book>",
+
"{\"book\":{\"#\":{\"$\":[1,3],\"price\":[2,4],\"title\":0},"
+ + "\"$\":[\"This is a description\",\"Text in
between\"],"
+ + "\"@pages\":\"320\","
+ +
"\"price\":[\"12.50\",{\"$\":\"13.50\",\"@currency\":\"EUR\"}],"
+ + "\"title\":\"Dune\"}}"),
+ arguments("<a><b><c>deep</c></b></a>",
"{\"a\":{\"b\":{\"c\":\"deep\"}}}"),
+ // Namespace prefixes and declarations are kept as written.
+ arguments(
+ "<ns:a xmlns:ns=\"urn:ns\" xmlns=\"urn:default\"
ns:id=\"1\">"
+ + "<ns:b>x</ns:b></ns:a>",
+
"{\"ns:a\":{\"@ns:id\":\"1\",\"@xmlns\":\"urn:default\","
+ + "\"@xmlns:ns\":\"urn:ns\",\"ns:b\":\"x\"}}"),
+ // Attribute values are not trimmed.
+ arguments("<a b=\" x \"/>", "{\"a\":{\"@b\":\" x \"}}"),
+ // CDATA is text, and it merges with the text around it.
+ arguments("<a><![CDATA[<b>&</b>]]></a>",
"{\"a\":\"<b>&</b>\"}"),
+ arguments("<a>x<![CDATA[y]]>z</a>", "{\"a\":\"xyz\"}"),
+ // Comments and processing instructions are dropped, the text
around them merges.
+ arguments(
+ "<?xml version=\"1.0\"?><!-- c --><a>foo<!-- c
-->bar<?pi x?></a><!-- c -->",
+ "{\"a\":\"foobar\"}"),
+ // The input is a string, so an encoding declared in the
document is ignored.
+ arguments(
+ "<?xml version=\"1.0\"
encoding=\"ISO-8859-1\"?><a>\u00e4</a>",
+ "{\"a\":\"\u00e4\"}"),
+ // Whitespace-only text is dropped, other text is trimmed of
XML whitespace only.
+ arguments("<a>\n <b>x</b>\n</a>", "{\"a\":{\"b\":\"x\"}}"),
+ arguments("<a> \t\r\n </a>", "{\"a\":\"\"}"),
+ arguments("<a>\u00a0x\u00a0</a>", "{\"a\":\"\u00a0x\u00a0\"}"),
+ // Predefined entities and character references are expanded,
in text and in
+ // attribute values.
+ arguments("<a><>&'</a>", "{\"a\":\"<>&'\"}"),
+ arguments("<a>AB</a>", "{\"a\":\"AB\"}"),
+ arguments("<a b=\"<A\"/>", "{\"a\":{\"@b\":\"<A\"}}"),
+ // Entities declared in the DTD are expanded.
+ arguments(
+ "<!DOCTYPE a [<!ENTITY title
\"Dune\">]><a>&title;</a>",
+ "{\"a\":\"Dune\"}"));
+ }
+
+ //
--------------------------------------------------------------------------------------------
+ // XML Schema instance attributes
+ //
--------------------------------------------------------------------------------------------
+
+ @ParameterizedTest(name = "{0}")
+ @MethodSource("xsiMappings")
+ void testXsiMapping(String xml, String expectedJson) {
+ assertThat(parser.parse(xml, false).toJson()).isEqualTo(expectedJson);
+ }
+
+ static Stream<Arguments> xsiMappings() {
+ return Stream.of(
+ // xsi:nil makes an element without other attributes null, and
drops the content.
+ arguments("<a xsi:nil=\"true\"/>", "{\"a\":null}"),
+ arguments("<a xsi:nil=\"1\"></a>", "{\"a\":null}"),
+ arguments("<a xsi:nil=\"true\">x</a>", "{\"a\":null}"),
+ arguments("<a xsi:nil=\"true\"><b>x</b></a>", "{\"a\":null}"),
+ arguments("<a xsi:nil=\" true \" xsi:type=\"int\"/>",
"{\"a\":null}"),
+ arguments("<a xsi:nil=\"true\" xsi:type=\"Dog\"/>",
"{\"a\":null}"),
+ // The declaration of the xsi prefix is dropped.
+ arguments(
+ "<root
xmlns:xsi=\"http://www.w3.org/2001/XMLSchema-instance\">"
+ + "<a xsi:nil=\"true\"/></root>",
+ "{\"root\":{\"a\":null}}"),
+ // xsi:nil that doesn't make the element null is kept, unless
it is false.
+ arguments(
+ "<a xsi:nil=\"true\" id=\"1\"/>",
+ "{\"a\":{\"@id\":\"1\",\"@xsi:nil\":\"true\"}}"),
+ arguments(
+ "<a xsi:nil=\"true\" id=\"1\">x</a>",
+ "{\"a\":{\"@id\":\"1\",\"@xsi:nil\":\"true\"}}"),
+ arguments(
+ "<a xsi:nil=\" 1 \" id=\"1\"/>",
+ "{\"a\":{\"@id\":\"1\",\"@xsi:nil\":\"true\"}}"),
+ arguments(
+ "<a xsi:nil=\"true\" xsi:type=\"string\" id=\"1\"/>",
+
"{\"a\":{\"@id\":\"1\",\"@xsi:nil\":\"true\",\"@xsi:type\":\"string\"}}"),
+ arguments("<a xsi:nil=\"yes\"/>",
"{\"a\":{\"@xsi:nil\":\"yes\"}}"),
+ arguments("<a xsi:nil=\"TRUE\"/>",
"{\"a\":{\"@xsi:nil\":\"TRUE\"}}"),
+ arguments("<a xsi:nil=\"false\">x</a>", "{\"a\":\"x\"}"),
+ arguments("<a xsi:nil=\"0\"/>", "{\"a\":\"\"}"),
+ // xsi:type types the text of an element with text but without
child elements.
+ arguments("<a xsi:type=\"int\">5</a>", "{\"a\":5}"),
+ arguments("<a xsi:type=\"boolean\">true</a>", "{\"a\":true}"),
+ arguments("<a xsi:type=\"decimal\">-1.23</a>",
"{\"a\":-1.23}"),
+ arguments(
+ "<price xsi:type=\"decimal\"
currency=\"EUR\">13.50</price>",
+ "{\"price\":{\"$\":13.5,\"@currency\":\"EUR\"}}"),
+ // xsi:type that doesn't type the text is kept.
+ arguments("<a xsi:type=\"int\"></a>",
"{\"a\":{\"@xsi:type\":\"int\"}}"),
+ arguments("<a xsi:type=\"string\"/>",
"{\"a\":{\"@xsi:type\":\"string\"}}"),
+ arguments(
+ "<a id=\"1\" xsi:type=\"string\"/>",
+ "{\"a\":{\"@id\":\"1\",\"@xsi:type\":\"string\"}}"),
+ arguments(
+ "<a xsi:type=\"int\">five</a>",
+ "{\"a\":{\"$\":\"five\",\"@xsi:type\":\"int\"}}"),
+ arguments(
+ "<a xsi:type=\"xs:token\">x</a>",
+ "{\"a\":{\"$\":\"x\",\"@xsi:type\":\"xs:token\"}}"),
+ arguments(
+ "<animal xsi:type=\"Dog\" name=\"Rex\"/>",
+
"{\"animal\":{\"@name\":\"Rex\",\"@xsi:type\":\"Dog\"}}"),
+ arguments(
+ "<shape
xsi:type=\"Circle\"><radius>1</radius></shape>",
+
"{\"shape\":{\"@xsi:type\":\"Circle\",\"radius\":\"1\"}}"),
+ // Other xsi: attributes are regular attributes.
+ arguments(
+ "<a
xmlns:xsi=\"http://www.w3.org/2001/XMLSchema-instance\" "
+ + "xsi:schemaLocation=\"urn:a a.xsd\">x</a>",
+ "{\"a\":{\"$\":\"x\",\"@xsi:schemaLocation\":\"urn:a
a.xsd\"}}"));
+ }
+
+ @ParameterizedTest(name = "{0}: {1}")
+ @MethodSource("typedValues")
+ void testXsiType(String xsiType, String text, Variant.Type expectedType,
Object expectedValue) {
+ final Variant value = parseTyped(xsiType, text);
+ assertThat(value.getType()).isEqualTo(expectedType);
+ if (expectedValue instanceof BigDecimal) {
+ assertThat(value.getDecimal()).isEqualByComparingTo((BigDecimal)
expectedValue);
+ } else {
+ assertThat(value.get()).isEqualTo(expectedValue);
+ }
+ }
+
+ static Stream<Arguments> typedValues() {
+ return Stream.of(
+ arguments("string", "5", Variant.Type.STRING, "5"),
+ arguments("boolean", "true", Variant.Type.BOOLEAN, true),
+ arguments("boolean", "false", Variant.Type.BOOLEAN, false),
+ arguments("boolean", "1", Variant.Type.BOOLEAN, true),
+ arguments("boolean", "0", Variant.Type.BOOLEAN, false),
+ arguments("byte", "-128", Variant.Type.TINYINT, (byte) -128),
+ arguments("short", "32767", Variant.Type.SMALLINT, (short)
32767),
+ arguments("int", "+5", Variant.Type.INT, 5),
+ arguments("int", " 5 ", Variant.Type.INT, 5),
+ arguments("long", "9223372036854775807", Variant.Type.BIGINT,
Long.MAX_VALUE),
+ arguments(
+ "integer",
+ "123456789012345678901234567890",
+ Variant.Type.DECIMAL,
+ new BigDecimal("123456789012345678901234567890")),
+ arguments("decimal", "12.50", Variant.Type.DECIMAL, new
BigDecimal("12.50")),
+ arguments("decimal", ".5", Variant.Type.DECIMAL, new
BigDecimal("0.5")),
+ arguments("float", "1.5", Variant.Type.FLOAT, 1.5f),
+ arguments("double", "1.5E3", Variant.Type.DOUBLE, 1500.0),
+ arguments("double", "-.5e+2", Variant.Type.DOUBLE, -50.0),
+ arguments("date", "2026-09-23", Variant.Type.DATE,
LocalDate.of(2026, 9, 23)),
+ arguments(
+ "time",
+ "12:30:45.123456",
+ Variant.Type.TIME,
+ LocalTime.of(12, 30, 45, 123_456_000)),
+ // Precision that doesn't fit is dropped, like in a cast.
+ arguments(
+ "time",
+ "12:30:45.123456789",
+ Variant.Type.TIME,
+ LocalTime.of(12, 30, 45, 123_456_000)),
+ arguments(
+ "dateTime",
+ "2300-01-01T00:00:00.000000001",
+ Variant.Type.TIMESTAMP,
+ LocalDateTime.of(2300, 1, 1, 0, 0)),
+ arguments(
+ "dateTime",
+ "2026-09-23T12:30:45",
+ Variant.Type.TIMESTAMP,
+ LocalDateTime.of(2026, 9, 23, 12, 30, 45)),
+ arguments(
+ "dateTime",
+ "2026-09-23T12:30:45.123456789",
+ Variant.Type.TIMESTAMP_NS,
+ LocalDateTime.of(2026, 9, 23, 12, 30, 45,
123_456_789)),
+ arguments(
+ "dateTime",
+ "2026-09-23T12:30:45Z",
+ Variant.Type.TIMESTAMP_LTZ,
+ Instant.parse("2026-09-23T12:30:45Z")),
+ arguments(
+ "dateTime",
+ "2026-09-23T12:30:45.5+02:00",
+ Variant.Type.TIMESTAMP_LTZ,
+ Instant.parse("2026-09-23T10:30:45.5Z")),
+ arguments(
+ "dateTime",
+ "2026-09-23T12:30:45.123456789Z",
+ Variant.Type.TIMESTAMP_LTZ_NS,
+ Instant.parse("2026-09-23T12:30:45.123456789Z")),
+ // Types are matched by their local part.
+ arguments("xs:int", "5", Variant.Type.INT, 5),
+ arguments("xsd:int", "5", Variant.Type.INT, 5));
+ }
+
+ @ParameterizedTest(name = "{0}: {1}")
+ @MethodSource({"untypedValues", "invalidLexicalForms"})
+ void testXsiTypeThatDoesNotApply(String xsiType, String text) {
+ final Variant value = parseTyped(xsiType, text);
+ assertThat(value.getField("$").getString()).isEqualTo(text);
+ assertThat(value.getField("@xsi:type").getString()).isEqualTo(xsiType);
+ }
+
+ static Stream<Arguments> untypedValues() {
+ return Stream.of(
+ arguments("boolean", "yes"),
+ arguments("boolean", "TRUE"),
+ arguments("byte", "128"),
+ arguments("int", "5.0"),
+ arguments("integer", "1.5"),
+ // Longer numbers are not typed.
+ arguments("integer", "1" + "0".repeat(1000)),
+ arguments("decimal", "1e99999"),
+ arguments("decimal",
"1234567890123456789012345678901234567890"),
+ arguments("float", "1e50"),
+ arguments("float", "INF"),
+ arguments("double", "1e99999"),
+ arguments("double", "NaN"),
+ arguments("double", "Infinity"),
+ arguments("date", "2026-02-30"),
+ arguments("date", "2026-09-23Z"),
+ arguments("time", "12:30:45Z"),
+ arguments("dateTime", "2026-09-23"),
+ // Valid in XML Schema, but rejected by the Java parsers.
+ arguments("date", "10000-01-01"),
+ arguments("dateTime", "10000-01-01T00:00:00"),
+ arguments("time", "24:00:00"),
+ arguments("dateTime", "2026-09-23T24:00:00"),
+ arguments("anyURI", "urn:x"),
+ arguments("Int", "5"));
+ }
+
+ /** Values that the Java parsers accept, but XML Schema doesn't. */
+ static Stream<Arguments> invalidLexicalForms() {
+ return Stream.of(
+ arguments("float", "1f"),
+ arguments("float", "1.5F"),
+ arguments("float", "1d"),
+ arguments("double", "1.5D"),
+ arguments("float", "0x1.8p1"),
+ arguments("double", "0x1.8p1"),
+ // Digits from other scripts: Arabic-Indic three, and
fullwidth one and two.
+ arguments("int", "\u0663"), // ٣
+ arguments("long", "\uff11\uff12"), // 12
+ arguments("integer", "\u0663"), // ٣
+ arguments("decimal", "\u0663"), // ٣
+ arguments("decimal", "1e5"),
+ arguments("time", "12:30"),
+ arguments("dateTime", "2026-09-23T12:30"),
+ arguments("dateTime", "2026-09-23t12:30:00"),
+ arguments("dateTime", "2026-09-23T12:30:00z"),
+ arguments("dateTime",
"2026-09-23T12:30:00+01:00[Europe/Paris]"));
+ }
+
+ // The XML Schema validator of the JDK is the reference for which values
are valid. It
+ // implements
+ // XML Schema 1.0, so values that only 1.1 allows, like the year 0000,
can't be typed values
+ // here.
+
+ @ParameterizedTest(name = "{0}: {1}")
+ @MethodSource("typedValues")
+ void testTypedValuesAreValidInXmlSchema(String xsiType, String text)
throws Exception {
+ assertThat(isValidInXmlSchema(xsiType, text)).isTrue();
+ }
+
+ @ParameterizedTest(name = "{0}: {1}")
+ @MethodSource("invalidLexicalForms")
+ void testInvalidLexicalFormsAreInvalidInXmlSchema(String xsiType, String
text)
+ throws Exception {
+ assertThat(isValidInXmlSchema(xsiType, text)).isFalse();
+ }
Review Comment:
Should we use `CsvSource` instead?
##########
flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/functions/XmlToVariantParserTest.java:
##########
@@ -0,0 +1,534 @@
+/*
+ * 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.types.variant.Variant;
+
+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 org.junit.jupiter.params.provider.ValueSource;
+import org.xml.sax.SAXException;
+
+import javax.xml.transform.stream.StreamSource;
+import javax.xml.validation.SchemaFactory;
+
+import java.io.StringReader;
+import java.math.BigDecimal;
+import java.time.Instant;
+import java.time.LocalDate;
+import java.time.LocalDateTime;
+import java.time.LocalTime;
+import java.util.stream.Stream;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatNoException;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import static org.junit.jupiter.params.provider.Arguments.arguments;
+
+/** Tests for {@link XmlToVariantParser}. */
+class XmlToVariantParserTest {
+
+ private final XmlToVariantParser parser = new XmlToVariantParser();
+
+ //
--------------------------------------------------------------------------------------------
+ // Mapping
+ //
--------------------------------------------------------------------------------------------
+
+ @ParameterizedTest(name = "{0}")
+ @MethodSource("mappings")
+ void testMapping(String xml, String expectedJson) {
+ assertThat(parser.parse(xml, false).toJson()).isEqualTo(expectedJson);
+ }
+
+ static Stream<Arguments> mappings() {
+ return Stream.of(
+ arguments("<title>Dune</title>", "{\"title\":\"Dune\"}"),
+ arguments("<title></title>", "{\"title\":\"\"}"),
+ arguments("<title/>", "{\"title\":\"\"}"),
+ arguments(
+ "<book pages=\"320\">\n Dune\n</book>",
+ "{\"book\":{\"$\":\"Dune\",\"@pages\":\"320\"}}"),
+ arguments("<book pages=\"320\"/>",
"{\"book\":{\"@pages\":\"320\"}}"),
+ arguments(
+ "<book>\n"
+ + " <author>Terry Pratchett</author>\n"
+ + " <author>Neil Gaiman</author>\n"
+ + "</book>",
+ "{\"book\":{\"author\":[\"Terry Pratchett\",\"Neil
Gaiman\"]}}"),
+ arguments(
+ "<book>\n <title>Dune</title>\n
<price>12.5</price>\n</book>",
+
"{\"book\":{\"#\":{\"price\":1,\"title\":0},\"price\":\"12.5\",\"title\":\"Dune\"}}"),
+ arguments(
+ "<book>\n"
+ + " This is a description\n"
+ + " <title>Dune</title>\n"
+ + " and this is another text\n"
+ + "</book>",
+ "{\"book\":{\"#\":{\"$\":[0,2],\"title\":1},"
+ + "\"$\":[\"This is a description\",\"and this
is another text\"],"
+ + "\"title\":\"Dune\"}}"),
+ arguments(
+ "<book pages=\"320\">\n"
+ + " <title>Dune</title>\n"
+ + " This is a description\n"
+ + " <price>12.50</price>\n"
+ + " Text in between\n"
+ + " <price currency=\"EUR\">13.50</price>\n"
+ + "</book>",
+
"{\"book\":{\"#\":{\"$\":[1,3],\"price\":[2,4],\"title\":0},"
+ + "\"$\":[\"This is a description\",\"Text in
between\"],"
+ + "\"@pages\":\"320\","
+ +
"\"price\":[\"12.50\",{\"$\":\"13.50\",\"@currency\":\"EUR\"}],"
+ + "\"title\":\"Dune\"}}"),
+ arguments("<a><b><c>deep</c></b></a>",
"{\"a\":{\"b\":{\"c\":\"deep\"}}}"),
+ // Namespace prefixes and declarations are kept as written.
+ arguments(
+ "<ns:a xmlns:ns=\"urn:ns\" xmlns=\"urn:default\"
ns:id=\"1\">"
+ + "<ns:b>x</ns:b></ns:a>",
+
"{\"ns:a\":{\"@ns:id\":\"1\",\"@xmlns\":\"urn:default\","
+ + "\"@xmlns:ns\":\"urn:ns\",\"ns:b\":\"x\"}}"),
+ // Attribute values are not trimmed.
+ arguments("<a b=\" x \"/>", "{\"a\":{\"@b\":\" x \"}}"),
+ // CDATA is text, and it merges with the text around it.
+ arguments("<a><![CDATA[<b>&</b>]]></a>",
"{\"a\":\"<b>&</b>\"}"),
+ arguments("<a>x<![CDATA[y]]>z</a>", "{\"a\":\"xyz\"}"),
+ // Comments and processing instructions are dropped, the text
around them merges.
+ arguments(
+ "<?xml version=\"1.0\"?><!-- c --><a>foo<!-- c
-->bar<?pi x?></a><!-- c -->",
+ "{\"a\":\"foobar\"}"),
+ // The input is a string, so an encoding declared in the
document is ignored.
+ arguments(
+ "<?xml version=\"1.0\"
encoding=\"ISO-8859-1\"?><a>\u00e4</a>",
+ "{\"a\":\"\u00e4\"}"),
+ // Whitespace-only text is dropped, other text is trimmed of
XML whitespace only.
+ arguments("<a>\n <b>x</b>\n</a>", "{\"a\":{\"b\":\"x\"}}"),
+ arguments("<a> \t\r\n </a>", "{\"a\":\"\"}"),
+ arguments("<a>\u00a0x\u00a0</a>", "{\"a\":\"\u00a0x\u00a0\"}"),
+ // Predefined entities and character references are expanded,
in text and in
+ // attribute values.
+ arguments("<a><>&'</a>", "{\"a\":\"<>&'\"}"),
+ arguments("<a>AB</a>", "{\"a\":\"AB\"}"),
+ arguments("<a b=\"<A\"/>", "{\"a\":{\"@b\":\"<A\"}}"),
+ // Entities declared in the DTD are expanded.
+ arguments(
+ "<!DOCTYPE a [<!ENTITY title
\"Dune\">]><a>&title;</a>",
+ "{\"a\":\"Dune\"}"));
+ }
+
+ //
--------------------------------------------------------------------------------------------
+ // XML Schema instance attributes
+ //
--------------------------------------------------------------------------------------------
+
+ @ParameterizedTest(name = "{0}")
+ @MethodSource("xsiMappings")
+ void testXsiMapping(String xml, String expectedJson) {
+ assertThat(parser.parse(xml, false).toJson()).isEqualTo(expectedJson);
+ }
+
+ static Stream<Arguments> xsiMappings() {
+ return Stream.of(
+ // xsi:nil makes an element without other attributes null, and
drops the content.
+ arguments("<a xsi:nil=\"true\"/>", "{\"a\":null}"),
+ arguments("<a xsi:nil=\"1\"></a>", "{\"a\":null}"),
+ arguments("<a xsi:nil=\"true\">x</a>", "{\"a\":null}"),
+ arguments("<a xsi:nil=\"true\"><b>x</b></a>", "{\"a\":null}"),
+ arguments("<a xsi:nil=\" true \" xsi:type=\"int\"/>",
"{\"a\":null}"),
+ arguments("<a xsi:nil=\"true\" xsi:type=\"Dog\"/>",
"{\"a\":null}"),
+ // The declaration of the xsi prefix is dropped.
+ arguments(
+ "<root
xmlns:xsi=\"http://www.w3.org/2001/XMLSchema-instance\">"
+ + "<a xsi:nil=\"true\"/></root>",
+ "{\"root\":{\"a\":null}}"),
+ // xsi:nil that doesn't make the element null is kept, unless
it is false.
+ arguments(
+ "<a xsi:nil=\"true\" id=\"1\"/>",
+ "{\"a\":{\"@id\":\"1\",\"@xsi:nil\":\"true\"}}"),
+ arguments(
+ "<a xsi:nil=\"true\" id=\"1\">x</a>",
+ "{\"a\":{\"@id\":\"1\",\"@xsi:nil\":\"true\"}}"),
+ arguments(
+ "<a xsi:nil=\" 1 \" id=\"1\"/>",
+ "{\"a\":{\"@id\":\"1\",\"@xsi:nil\":\"true\"}}"),
+ arguments(
+ "<a xsi:nil=\"true\" xsi:type=\"string\" id=\"1\"/>",
+
"{\"a\":{\"@id\":\"1\",\"@xsi:nil\":\"true\",\"@xsi:type\":\"string\"}}"),
+ arguments("<a xsi:nil=\"yes\"/>",
"{\"a\":{\"@xsi:nil\":\"yes\"}}"),
+ arguments("<a xsi:nil=\"TRUE\"/>",
"{\"a\":{\"@xsi:nil\":\"TRUE\"}}"),
+ arguments("<a xsi:nil=\"false\">x</a>", "{\"a\":\"x\"}"),
+ arguments("<a xsi:nil=\"0\"/>", "{\"a\":\"\"}"),
+ // xsi:type types the text of an element with text but without
child elements.
+ arguments("<a xsi:type=\"int\">5</a>", "{\"a\":5}"),
+ arguments("<a xsi:type=\"boolean\">true</a>", "{\"a\":true}"),
+ arguments("<a xsi:type=\"decimal\">-1.23</a>",
"{\"a\":-1.23}"),
+ arguments(
+ "<price xsi:type=\"decimal\"
currency=\"EUR\">13.50</price>",
+ "{\"price\":{\"$\":13.5,\"@currency\":\"EUR\"}}"),
+ // xsi:type that doesn't type the text is kept.
+ arguments("<a xsi:type=\"int\"></a>",
"{\"a\":{\"@xsi:type\":\"int\"}}"),
+ arguments("<a xsi:type=\"string\"/>",
"{\"a\":{\"@xsi:type\":\"string\"}}"),
+ arguments(
+ "<a id=\"1\" xsi:type=\"string\"/>",
+ "{\"a\":{\"@id\":\"1\",\"@xsi:type\":\"string\"}}"),
+ arguments(
+ "<a xsi:type=\"int\">five</a>",
+ "{\"a\":{\"$\":\"five\",\"@xsi:type\":\"int\"}}"),
+ arguments(
+ "<a xsi:type=\"xs:token\">x</a>",
+ "{\"a\":{\"$\":\"x\",\"@xsi:type\":\"xs:token\"}}"),
+ arguments(
+ "<animal xsi:type=\"Dog\" name=\"Rex\"/>",
+
"{\"animal\":{\"@name\":\"Rex\",\"@xsi:type\":\"Dog\"}}"),
+ arguments(
+ "<shape
xsi:type=\"Circle\"><radius>1</radius></shape>",
+
"{\"shape\":{\"@xsi:type\":\"Circle\",\"radius\":\"1\"}}"),
+ // Other xsi: attributes are regular attributes.
+ arguments(
+ "<a
xmlns:xsi=\"http://www.w3.org/2001/XMLSchema-instance\" "
+ + "xsi:schemaLocation=\"urn:a a.xsd\">x</a>",
+ "{\"a\":{\"$\":\"x\",\"@xsi:schemaLocation\":\"urn:a
a.xsd\"}}"));
+ }
+
+ @ParameterizedTest(name = "{0}: {1}")
+ @MethodSource("typedValues")
+ void testXsiType(String xsiType, String text, Variant.Type expectedType,
Object expectedValue) {
+ final Variant value = parseTyped(xsiType, text);
+ assertThat(value.getType()).isEqualTo(expectedType);
+ if (expectedValue instanceof BigDecimal) {
+ assertThat(value.getDecimal()).isEqualByComparingTo((BigDecimal)
expectedValue);
+ } else {
+ assertThat(value.get()).isEqualTo(expectedValue);
+ }
+ }
+
+ static Stream<Arguments> typedValues() {
+ return Stream.of(
+ arguments("string", "5", Variant.Type.STRING, "5"),
+ arguments("boolean", "true", Variant.Type.BOOLEAN, true),
+ arguments("boolean", "false", Variant.Type.BOOLEAN, false),
+ arguments("boolean", "1", Variant.Type.BOOLEAN, true),
+ arguments("boolean", "0", Variant.Type.BOOLEAN, false),
+ arguments("byte", "-128", Variant.Type.TINYINT, (byte) -128),
+ arguments("short", "32767", Variant.Type.SMALLINT, (short)
32767),
+ arguments("int", "+5", Variant.Type.INT, 5),
+ arguments("int", " 5 ", Variant.Type.INT, 5),
+ arguments("long", "9223372036854775807", Variant.Type.BIGINT,
Long.MAX_VALUE),
+ arguments(
+ "integer",
+ "123456789012345678901234567890",
+ Variant.Type.DECIMAL,
+ new BigDecimal("123456789012345678901234567890")),
+ arguments("decimal", "12.50", Variant.Type.DECIMAL, new
BigDecimal("12.50")),
+ arguments("decimal", ".5", Variant.Type.DECIMAL, new
BigDecimal("0.5")),
+ arguments("float", "1.5", Variant.Type.FLOAT, 1.5f),
+ arguments("double", "1.5E3", Variant.Type.DOUBLE, 1500.0),
+ arguments("double", "-.5e+2", Variant.Type.DOUBLE, -50.0),
+ arguments("date", "2026-09-23", Variant.Type.DATE,
LocalDate.of(2026, 9, 23)),
+ arguments(
+ "time",
+ "12:30:45.123456",
+ Variant.Type.TIME,
+ LocalTime.of(12, 30, 45, 123_456_000)),
+ // Precision that doesn't fit is dropped, like in a cast.
+ arguments(
+ "time",
+ "12:30:45.123456789",
+ Variant.Type.TIME,
+ LocalTime.of(12, 30, 45, 123_456_000)),
+ arguments(
+ "dateTime",
+ "2300-01-01T00:00:00.000000001",
+ Variant.Type.TIMESTAMP,
+ LocalDateTime.of(2300, 1, 1, 0, 0)),
+ arguments(
+ "dateTime",
+ "2026-09-23T12:30:45",
+ Variant.Type.TIMESTAMP,
+ LocalDateTime.of(2026, 9, 23, 12, 30, 45)),
+ arguments(
+ "dateTime",
+ "2026-09-23T12:30:45.123456789",
+ Variant.Type.TIMESTAMP_NS,
+ LocalDateTime.of(2026, 9, 23, 12, 30, 45,
123_456_789)),
+ arguments(
+ "dateTime",
+ "2026-09-23T12:30:45Z",
+ Variant.Type.TIMESTAMP_LTZ,
+ Instant.parse("2026-09-23T12:30:45Z")),
+ arguments(
+ "dateTime",
+ "2026-09-23T12:30:45.5+02:00",
+ Variant.Type.TIMESTAMP_LTZ,
+ Instant.parse("2026-09-23T10:30:45.5Z")),
+ arguments(
+ "dateTime",
+ "2026-09-23T12:30:45.123456789Z",
+ Variant.Type.TIMESTAMP_LTZ_NS,
+ Instant.parse("2026-09-23T12:30:45.123456789Z")),
+ // Types are matched by their local part.
+ arguments("xs:int", "5", Variant.Type.INT, 5),
+ arguments("xsd:int", "5", Variant.Type.INT, 5));
+ }
+
+ @ParameterizedTest(name = "{0}: {1}")
+ @MethodSource({"untypedValues", "invalidLexicalForms"})
+ void testXsiTypeThatDoesNotApply(String xsiType, String text) {
+ final Variant value = parseTyped(xsiType, text);
+ assertThat(value.getField("$").getString()).isEqualTo(text);
+ assertThat(value.getField("@xsi:type").getString()).isEqualTo(xsiType);
+ }
+
+ static Stream<Arguments> untypedValues() {
+ return Stream.of(
+ arguments("boolean", "yes"),
+ arguments("boolean", "TRUE"),
+ arguments("byte", "128"),
+ arguments("int", "5.0"),
+ arguments("integer", "1.5"),
+ // Longer numbers are not typed.
+ arguments("integer", "1" + "0".repeat(1000)),
+ arguments("decimal", "1e99999"),
+ arguments("decimal",
"1234567890123456789012345678901234567890"),
+ arguments("float", "1e50"),
+ arguments("float", "INF"),
+ arguments("double", "1e99999"),
+ arguments("double", "NaN"),
+ arguments("double", "Infinity"),
+ arguments("date", "2026-02-30"),
+ arguments("date", "2026-09-23Z"),
+ arguments("time", "12:30:45Z"),
+ arguments("dateTime", "2026-09-23"),
+ // Valid in XML Schema, but rejected by the Java parsers.
+ arguments("date", "10000-01-01"),
+ arguments("dateTime", "10000-01-01T00:00:00"),
+ arguments("time", "24:00:00"),
+ arguments("dateTime", "2026-09-23T24:00:00"),
+ arguments("anyURI", "urn:x"),
+ arguments("Int", "5"));
+ }
+
+ /** Values that the Java parsers accept, but XML Schema doesn't. */
+ static Stream<Arguments> invalidLexicalForms() {
+ return Stream.of(
+ arguments("float", "1f"),
+ arguments("float", "1.5F"),
+ arguments("float", "1d"),
+ arguments("double", "1.5D"),
+ arguments("float", "0x1.8p1"),
+ arguments("double", "0x1.8p1"),
+ // Digits from other scripts: Arabic-Indic three, and
fullwidth one and two.
+ arguments("int", "\u0663"), // ٣
+ arguments("long", "\uff11\uff12"), // 12
+ arguments("integer", "\u0663"), // ٣
+ arguments("decimal", "\u0663"), // ٣
+ arguments("decimal", "1e5"),
+ arguments("time", "12:30"),
+ arguments("dateTime", "2026-09-23T12:30"),
+ arguments("dateTime", "2026-09-23t12:30:00"),
+ arguments("dateTime", "2026-09-23T12:30:00z"),
+ arguments("dateTime",
"2026-09-23T12:30:00+01:00[Europe/Paris]"));
+ }
+
+ // The XML Schema validator of the JDK is the reference for which values
are valid. It
+ // implements
+ // XML Schema 1.0, so values that only 1.1 allows, like the year 0000,
can't be typed values
+ // here.
+
+ @ParameterizedTest(name = "{0}: {1}")
+ @MethodSource("typedValues")
+ void testTypedValuesAreValidInXmlSchema(String xsiType, String text)
throws Exception {
+ assertThat(isValidInXmlSchema(xsiType, text)).isTrue();
+ }
+
+ @ParameterizedTest(name = "{0}: {1}")
+ @MethodSource("invalidLexicalForms")
+ void testInvalidLexicalFormsAreInvalidInXmlSchema(String xsiType, String
text)
+ throws Exception {
+ assertThat(isValidInXmlSchema(xsiType, text)).isFalse();
+ }
+
+ /** The validator checks the text against the xsi:type of the element,
without a schema. */
+ private static boolean isValidInXmlSchema(String xsiType, String text)
throws Exception {
+ final String xml =
+ String.format(
+ "<a
xmlns:xsi=\"http://www.w3.org/2001/XMLSchema-instance\""
+ + "
xmlns:xs=\"http://www.w3.org/2001/XMLSchema\""
+ + "
xmlns:xsd=\"http://www.w3.org/2001/XMLSchema\""
+ + " xsi:type=\"%s\">%s</a>",
+ xsiType.contains(":") ? xsiType : "xs:" + xsiType,
text);
+ try {
+ SchemaFactory.newDefaultInstance()
+ .newSchema()
+ .newValidator()
+ .validate(new StreamSource(new StringReader(xml)));
+ return true;
+ } catch (SAXException e) {
+ // Any other error, e.g. an unknown type, is a mistake in the test.
+ if (!e.getMessage().startsWith("cvc-datatype-valid")) {
+ throw e;
+ }
+ return false;
+ }
+ }
+
+ private Variant parseTyped(String xsiType, String text) {
+ return parser.parse(String.format("<a xsi:type=\"%s\">%s</a>",
xsiType, text), false)
+ .getField("a");
+ }
+
+ //
--------------------------------------------------------------------------------------------
+ // forceArray
+ //
--------------------------------------------------------------------------------------------
Review Comment:
Use `@Nested` to separate tests from each other. You can then move each
helper that is specific to that test family into the nested class
##########
flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/functions/SqlXmlUtilsTest.java:
##########
Review Comment:
How useful are these test? Should we keep them? I think we have enough tests
that cover these logics
##########
flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/api/internal/BaseExpressions.java:
##########
@@ -1420,6 +1422,51 @@ public OutType tryParseJson(boolean allowDuplicateKeys) {
unresolvedCall(TRY_PARSE_JSON, toExpr(),
valueLiteral(allowDuplicateKeys)));
}
+ /**
+ * Parses an XML string into a value of type {@link DataTypes#VARIANT()}.
If the XML string is
+ * invalid, an error is thrown. To return {@code NULL} instead of an
error, use {@link
+ * #tryParseXml()}.
+ *
+ * <p>This is a shortcut for {@code parseXml(false)}. See {@link
#parseXml(boolean)}.
+ */
+ public OutType parseXml() {
+ return toApiSpecificExpression(unresolvedCall(PARSE_XML,
objectToExpression(toExpr())));
+ }
+
+ /**
+ * Parses an XML string into a value of type {@link DataTypes#VARIANT()}.
If the XML string is
+ * invalid, an error is thrown. To return {@code NULL} instead of an
error, use {@link
+ * #tryParseXml(boolean)}.
+ *
+ * <p>If {@code forceArray} is {@code true}, every child element is stored
as an array, even if
Review Comment:
The Javadoc and the Python docstring only mention child elements. With
`force_array`, the text and the `#` positions become arrays too. Could you use
the wording from the SQL docs?
##########
flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/functions/XmlToVariantParserTest.java:
##########
@@ -0,0 +1,534 @@
+/*
+ * 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.types.variant.Variant;
+
+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 org.junit.jupiter.params.provider.ValueSource;
+import org.xml.sax.SAXException;
+
+import javax.xml.transform.stream.StreamSource;
+import javax.xml.validation.SchemaFactory;
+
+import java.io.StringReader;
+import java.math.BigDecimal;
+import java.time.Instant;
+import java.time.LocalDate;
+import java.time.LocalDateTime;
+import java.time.LocalTime;
+import java.util.stream.Stream;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatNoException;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import static org.junit.jupiter.params.provider.Arguments.arguments;
+
+/** Tests for {@link XmlToVariantParser}. */
+class XmlToVariantParserTest {
+
+ private final XmlToVariantParser parser = new XmlToVariantParser();
+
+ //
--------------------------------------------------------------------------------------------
+ // Mapping
+ //
--------------------------------------------------------------------------------------------
+
+ @ParameterizedTest(name = "{0}")
+ @MethodSource("mappings")
+ void testMapping(String xml, String expectedJson) {
+ assertThat(parser.parse(xml, false).toJson()).isEqualTo(expectedJson);
+ }
+
+ static Stream<Arguments> mappings() {
+ return Stream.of(
+ arguments("<title>Dune</title>", "{\"title\":\"Dune\"}"),
+ arguments("<title></title>", "{\"title\":\"\"}"),
+ arguments("<title/>", "{\"title\":\"\"}"),
+ arguments(
+ "<book pages=\"320\">\n Dune\n</book>",
+ "{\"book\":{\"$\":\"Dune\",\"@pages\":\"320\"}}"),
+ arguments("<book pages=\"320\"/>",
"{\"book\":{\"@pages\":\"320\"}}"),
+ arguments(
+ "<book>\n"
+ + " <author>Terry Pratchett</author>\n"
+ + " <author>Neil Gaiman</author>\n"
+ + "</book>",
+ "{\"book\":{\"author\":[\"Terry Pratchett\",\"Neil
Gaiman\"]}}"),
+ arguments(
+ "<book>\n <title>Dune</title>\n
<price>12.5</price>\n</book>",
+
"{\"book\":{\"#\":{\"price\":1,\"title\":0},\"price\":\"12.5\",\"title\":\"Dune\"}}"),
+ arguments(
+ "<book>\n"
+ + " This is a description\n"
+ + " <title>Dune</title>\n"
+ + " and this is another text\n"
+ + "</book>",
+ "{\"book\":{\"#\":{\"$\":[0,2],\"title\":1},"
+ + "\"$\":[\"This is a description\",\"and this
is another text\"],"
+ + "\"title\":\"Dune\"}}"),
+ arguments(
+ "<book pages=\"320\">\n"
+ + " <title>Dune</title>\n"
+ + " This is a description\n"
+ + " <price>12.50</price>\n"
+ + " Text in between\n"
+ + " <price currency=\"EUR\">13.50</price>\n"
+ + "</book>",
+
"{\"book\":{\"#\":{\"$\":[1,3],\"price\":[2,4],\"title\":0},"
+ + "\"$\":[\"This is a description\",\"Text in
between\"],"
+ + "\"@pages\":\"320\","
+ +
"\"price\":[\"12.50\",{\"$\":\"13.50\",\"@currency\":\"EUR\"}],"
+ + "\"title\":\"Dune\"}}"),
+ arguments("<a><b><c>deep</c></b></a>",
"{\"a\":{\"b\":{\"c\":\"deep\"}}}"),
+ // Namespace prefixes and declarations are kept as written.
+ arguments(
+ "<ns:a xmlns:ns=\"urn:ns\" xmlns=\"urn:default\"
ns:id=\"1\">"
+ + "<ns:b>x</ns:b></ns:a>",
+
"{\"ns:a\":{\"@ns:id\":\"1\",\"@xmlns\":\"urn:default\","
+ + "\"@xmlns:ns\":\"urn:ns\",\"ns:b\":\"x\"}}"),
+ // Attribute values are not trimmed.
+ arguments("<a b=\" x \"/>", "{\"a\":{\"@b\":\" x \"}}"),
+ // CDATA is text, and it merges with the text around it.
+ arguments("<a><![CDATA[<b>&</b>]]></a>",
"{\"a\":\"<b>&</b>\"}"),
+ arguments("<a>x<![CDATA[y]]>z</a>", "{\"a\":\"xyz\"}"),
+ // Comments and processing instructions are dropped, the text
around them merges.
+ arguments(
+ "<?xml version=\"1.0\"?><!-- c --><a>foo<!-- c
-->bar<?pi x?></a><!-- c -->",
+ "{\"a\":\"foobar\"}"),
+ // The input is a string, so an encoding declared in the
document is ignored.
+ arguments(
+ "<?xml version=\"1.0\"
encoding=\"ISO-8859-1\"?><a>\u00e4</a>",
+ "{\"a\":\"\u00e4\"}"),
+ // Whitespace-only text is dropped, other text is trimmed of
XML whitespace only.
+ arguments("<a>\n <b>x</b>\n</a>", "{\"a\":{\"b\":\"x\"}}"),
+ arguments("<a> \t\r\n </a>", "{\"a\":\"\"}"),
+ arguments("<a>\u00a0x\u00a0</a>", "{\"a\":\"\u00a0x\u00a0\"}"),
+ // Predefined entities and character references are expanded,
in text and in
+ // attribute values.
+ arguments("<a><>&'</a>", "{\"a\":\"<>&'\"}"),
+ arguments("<a>AB</a>", "{\"a\":\"AB\"}"),
+ arguments("<a b=\"<A\"/>", "{\"a\":{\"@b\":\"<A\"}}"),
+ // Entities declared in the DTD are expanded.
+ arguments(
+ "<!DOCTYPE a [<!ENTITY title
\"Dune\">]><a>&title;</a>",
+ "{\"a\":\"Dune\"}"));
+ }
+
+ //
--------------------------------------------------------------------------------------------
+ // XML Schema instance attributes
+ //
--------------------------------------------------------------------------------------------
+
+ @ParameterizedTest(name = "{0}")
+ @MethodSource("xsiMappings")
+ void testXsiMapping(String xml, String expectedJson) {
+ assertThat(parser.parse(xml, false).toJson()).isEqualTo(expectedJson);
+ }
+
+ static Stream<Arguments> xsiMappings() {
+ return Stream.of(
+ // xsi:nil makes an element without other attributes null, and
drops the content.
+ arguments("<a xsi:nil=\"true\"/>", "{\"a\":null}"),
+ arguments("<a xsi:nil=\"1\"></a>", "{\"a\":null}"),
+ arguments("<a xsi:nil=\"true\">x</a>", "{\"a\":null}"),
+ arguments("<a xsi:nil=\"true\"><b>x</b></a>", "{\"a\":null}"),
+ arguments("<a xsi:nil=\" true \" xsi:type=\"int\"/>",
"{\"a\":null}"),
+ arguments("<a xsi:nil=\"true\" xsi:type=\"Dog\"/>",
"{\"a\":null}"),
+ // The declaration of the xsi prefix is dropped.
+ arguments(
+ "<root
xmlns:xsi=\"http://www.w3.org/2001/XMLSchema-instance\">"
+ + "<a xsi:nil=\"true\"/></root>",
+ "{\"root\":{\"a\":null}}"),
+ // xsi:nil that doesn't make the element null is kept, unless
it is false.
+ arguments(
+ "<a xsi:nil=\"true\" id=\"1\"/>",
+ "{\"a\":{\"@id\":\"1\",\"@xsi:nil\":\"true\"}}"),
+ arguments(
+ "<a xsi:nil=\"true\" id=\"1\">x</a>",
+ "{\"a\":{\"@id\":\"1\",\"@xsi:nil\":\"true\"}}"),
+ arguments(
+ "<a xsi:nil=\" 1 \" id=\"1\"/>",
+ "{\"a\":{\"@id\":\"1\",\"@xsi:nil\":\"true\"}}"),
+ arguments(
+ "<a xsi:nil=\"true\" xsi:type=\"string\" id=\"1\"/>",
+
"{\"a\":{\"@id\":\"1\",\"@xsi:nil\":\"true\",\"@xsi:type\":\"string\"}}"),
+ arguments("<a xsi:nil=\"yes\"/>",
"{\"a\":{\"@xsi:nil\":\"yes\"}}"),
+ arguments("<a xsi:nil=\"TRUE\"/>",
"{\"a\":{\"@xsi:nil\":\"TRUE\"}}"),
+ arguments("<a xsi:nil=\"false\">x</a>", "{\"a\":\"x\"}"),
+ arguments("<a xsi:nil=\"0\"/>", "{\"a\":\"\"}"),
+ // xsi:type types the text of an element with text but without
child elements.
+ arguments("<a xsi:type=\"int\">5</a>", "{\"a\":5}"),
+ arguments("<a xsi:type=\"boolean\">true</a>", "{\"a\":true}"),
+ arguments("<a xsi:type=\"decimal\">-1.23</a>",
"{\"a\":-1.23}"),
+ arguments(
+ "<price xsi:type=\"decimal\"
currency=\"EUR\">13.50</price>",
+ "{\"price\":{\"$\":13.5,\"@currency\":\"EUR\"}}"),
+ // xsi:type that doesn't type the text is kept.
+ arguments("<a xsi:type=\"int\"></a>",
"{\"a\":{\"@xsi:type\":\"int\"}}"),
+ arguments("<a xsi:type=\"string\"/>",
"{\"a\":{\"@xsi:type\":\"string\"}}"),
+ arguments(
+ "<a id=\"1\" xsi:type=\"string\"/>",
+ "{\"a\":{\"@id\":\"1\",\"@xsi:type\":\"string\"}}"),
+ arguments(
+ "<a xsi:type=\"int\">five</a>",
+ "{\"a\":{\"$\":\"five\",\"@xsi:type\":\"int\"}}"),
+ arguments(
+ "<a xsi:type=\"xs:token\">x</a>",
+ "{\"a\":{\"$\":\"x\",\"@xsi:type\":\"xs:token\"}}"),
+ arguments(
+ "<animal xsi:type=\"Dog\" name=\"Rex\"/>",
+
"{\"animal\":{\"@name\":\"Rex\",\"@xsi:type\":\"Dog\"}}"),
+ arguments(
+ "<shape
xsi:type=\"Circle\"><radius>1</radius></shape>",
+
"{\"shape\":{\"@xsi:type\":\"Circle\",\"radius\":\"1\"}}"),
+ // Other xsi: attributes are regular attributes.
+ arguments(
+ "<a
xmlns:xsi=\"http://www.w3.org/2001/XMLSchema-instance\" "
+ + "xsi:schemaLocation=\"urn:a a.xsd\">x</a>",
+ "{\"a\":{\"$\":\"x\",\"@xsi:schemaLocation\":\"urn:a
a.xsd\"}}"));
+ }
+
+ @ParameterizedTest(name = "{0}: {1}")
+ @MethodSource("typedValues")
+ void testXsiType(String xsiType, String text, Variant.Type expectedType,
Object expectedValue) {
+ final Variant value = parseTyped(xsiType, text);
+ assertThat(value.getType()).isEqualTo(expectedType);
+ if (expectedValue instanceof BigDecimal) {
+ assertThat(value.getDecimal()).isEqualByComparingTo((BigDecimal)
expectedValue);
+ } else {
+ assertThat(value.get()).isEqualTo(expectedValue);
+ }
+ }
+
+ static Stream<Arguments> typedValues() {
+ return Stream.of(
+ arguments("string", "5", Variant.Type.STRING, "5"),
+ arguments("boolean", "true", Variant.Type.BOOLEAN, true),
+ arguments("boolean", "false", Variant.Type.BOOLEAN, false),
+ arguments("boolean", "1", Variant.Type.BOOLEAN, true),
+ arguments("boolean", "0", Variant.Type.BOOLEAN, false),
+ arguments("byte", "-128", Variant.Type.TINYINT, (byte) -128),
+ arguments("short", "32767", Variant.Type.SMALLINT, (short)
32767),
+ arguments("int", "+5", Variant.Type.INT, 5),
+ arguments("int", " 5 ", Variant.Type.INT, 5),
+ arguments("long", "9223372036854775807", Variant.Type.BIGINT,
Long.MAX_VALUE),
+ arguments(
+ "integer",
+ "123456789012345678901234567890",
+ Variant.Type.DECIMAL,
+ new BigDecimal("123456789012345678901234567890")),
+ arguments("decimal", "12.50", Variant.Type.DECIMAL, new
BigDecimal("12.50")),
+ arguments("decimal", ".5", Variant.Type.DECIMAL, new
BigDecimal("0.5")),
+ arguments("float", "1.5", Variant.Type.FLOAT, 1.5f),
+ arguments("double", "1.5E3", Variant.Type.DOUBLE, 1500.0),
+ arguments("double", "-.5e+2", Variant.Type.DOUBLE, -50.0),
+ arguments("date", "2026-09-23", Variant.Type.DATE,
LocalDate.of(2026, 9, 23)),
+ arguments(
+ "time",
+ "12:30:45.123456",
+ Variant.Type.TIME,
+ LocalTime.of(12, 30, 45, 123_456_000)),
+ // Precision that doesn't fit is dropped, like in a cast.
+ arguments(
+ "time",
+ "12:30:45.123456789",
+ Variant.Type.TIME,
+ LocalTime.of(12, 30, 45, 123_456_000)),
+ arguments(
+ "dateTime",
+ "2300-01-01T00:00:00.000000001",
+ Variant.Type.TIMESTAMP,
+ LocalDateTime.of(2300, 1, 1, 0, 0)),
+ arguments(
+ "dateTime",
+ "2026-09-23T12:30:45",
+ Variant.Type.TIMESTAMP,
+ LocalDateTime.of(2026, 9, 23, 12, 30, 45)),
+ arguments(
+ "dateTime",
+ "2026-09-23T12:30:45.123456789",
+ Variant.Type.TIMESTAMP_NS,
+ LocalDateTime.of(2026, 9, 23, 12, 30, 45,
123_456_789)),
+ arguments(
+ "dateTime",
+ "2026-09-23T12:30:45Z",
+ Variant.Type.TIMESTAMP_LTZ,
+ Instant.parse("2026-09-23T12:30:45Z")),
+ arguments(
+ "dateTime",
+ "2026-09-23T12:30:45.5+02:00",
+ Variant.Type.TIMESTAMP_LTZ,
+ Instant.parse("2026-09-23T10:30:45.5Z")),
+ arguments(
+ "dateTime",
+ "2026-09-23T12:30:45.123456789Z",
+ Variant.Type.TIMESTAMP_LTZ_NS,
+ Instant.parse("2026-09-23T12:30:45.123456789Z")),
+ // Types are matched by their local part.
+ arguments("xs:int", "5", Variant.Type.INT, 5),
+ arguments("xsd:int", "5", Variant.Type.INT, 5));
+ }
+
+ @ParameterizedTest(name = "{0}: {1}")
+ @MethodSource({"untypedValues", "invalidLexicalForms"})
+ void testXsiTypeThatDoesNotApply(String xsiType, String text) {
+ final Variant value = parseTyped(xsiType, text);
+ assertThat(value.getField("$").getString()).isEqualTo(text);
+ assertThat(value.getField("@xsi:type").getString()).isEqualTo(xsiType);
+ }
+
+ static Stream<Arguments> untypedValues() {
Review Comment:
The name is a bit misleading. They are actually typed as far as I can read
##########
flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/functions/XmlToVariantParserTest.java:
##########
@@ -0,0 +1,534 @@
+/*
+ * 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.types.variant.Variant;
+
+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 org.junit.jupiter.params.provider.ValueSource;
+import org.xml.sax.SAXException;
+
+import javax.xml.transform.stream.StreamSource;
+import javax.xml.validation.SchemaFactory;
+
+import java.io.StringReader;
+import java.math.BigDecimal;
+import java.time.Instant;
+import java.time.LocalDate;
+import java.time.LocalDateTime;
+import java.time.LocalTime;
+import java.util.stream.Stream;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatNoException;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import static org.junit.jupiter.params.provider.Arguments.arguments;
+
+/** Tests for {@link XmlToVariantParser}. */
+class XmlToVariantParserTest {
+
+ private final XmlToVariantParser parser = new XmlToVariantParser();
+
+ //
--------------------------------------------------------------------------------------------
+ // Mapping
+ //
--------------------------------------------------------------------------------------------
+
+ @ParameterizedTest(name = "{0}")
+ @MethodSource("mappings")
+ void testMapping(String xml, String expectedJson) {
+ assertThat(parser.parse(xml, false).toJson()).isEqualTo(expectedJson);
+ }
+
+ static Stream<Arguments> mappings() {
+ return Stream.of(
+ arguments("<title>Dune</title>", "{\"title\":\"Dune\"}"),
+ arguments("<title></title>", "{\"title\":\"\"}"),
+ arguments("<title/>", "{\"title\":\"\"}"),
+ arguments(
+ "<book pages=\"320\">\n Dune\n</book>",
+ "{\"book\":{\"$\":\"Dune\",\"@pages\":\"320\"}}"),
+ arguments("<book pages=\"320\"/>",
"{\"book\":{\"@pages\":\"320\"}}"),
+ arguments(
+ "<book>\n"
+ + " <author>Terry Pratchett</author>\n"
+ + " <author>Neil Gaiman</author>\n"
+ + "</book>",
+ "{\"book\":{\"author\":[\"Terry Pratchett\",\"Neil
Gaiman\"]}}"),
+ arguments(
+ "<book>\n <title>Dune</title>\n
<price>12.5</price>\n</book>",
+
"{\"book\":{\"#\":{\"price\":1,\"title\":0},\"price\":\"12.5\",\"title\":\"Dune\"}}"),
+ arguments(
+ "<book>\n"
+ + " This is a description\n"
+ + " <title>Dune</title>\n"
+ + " and this is another text\n"
+ + "</book>",
+ "{\"book\":{\"#\":{\"$\":[0,2],\"title\":1},"
+ + "\"$\":[\"This is a description\",\"and this
is another text\"],"
+ + "\"title\":\"Dune\"}}"),
+ arguments(
+ "<book pages=\"320\">\n"
+ + " <title>Dune</title>\n"
+ + " This is a description\n"
+ + " <price>12.50</price>\n"
+ + " Text in between\n"
+ + " <price currency=\"EUR\">13.50</price>\n"
+ + "</book>",
+
"{\"book\":{\"#\":{\"$\":[1,3],\"price\":[2,4],\"title\":0},"
+ + "\"$\":[\"This is a description\",\"Text in
between\"],"
+ + "\"@pages\":\"320\","
+ +
"\"price\":[\"12.50\",{\"$\":\"13.50\",\"@currency\":\"EUR\"}],"
+ + "\"title\":\"Dune\"}}"),
+ arguments("<a><b><c>deep</c></b></a>",
"{\"a\":{\"b\":{\"c\":\"deep\"}}}"),
+ // Namespace prefixes and declarations are kept as written.
+ arguments(
+ "<ns:a xmlns:ns=\"urn:ns\" xmlns=\"urn:default\"
ns:id=\"1\">"
+ + "<ns:b>x</ns:b></ns:a>",
+
"{\"ns:a\":{\"@ns:id\":\"1\",\"@xmlns\":\"urn:default\","
+ + "\"@xmlns:ns\":\"urn:ns\",\"ns:b\":\"x\"}}"),
+ // Attribute values are not trimmed.
+ arguments("<a b=\" x \"/>", "{\"a\":{\"@b\":\" x \"}}"),
+ // CDATA is text, and it merges with the text around it.
+ arguments("<a><![CDATA[<b>&</b>]]></a>",
"{\"a\":\"<b>&</b>\"}"),
+ arguments("<a>x<![CDATA[y]]>z</a>", "{\"a\":\"xyz\"}"),
+ // Comments and processing instructions are dropped, the text
around them merges.
+ arguments(
+ "<?xml version=\"1.0\"?><!-- c --><a>foo<!-- c
-->bar<?pi x?></a><!-- c -->",
+ "{\"a\":\"foobar\"}"),
+ // The input is a string, so an encoding declared in the
document is ignored.
+ arguments(
+ "<?xml version=\"1.0\"
encoding=\"ISO-8859-1\"?><a>\u00e4</a>",
+ "{\"a\":\"\u00e4\"}"),
+ // Whitespace-only text is dropped, other text is trimmed of
XML whitespace only.
+ arguments("<a>\n <b>x</b>\n</a>", "{\"a\":{\"b\":\"x\"}}"),
+ arguments("<a> \t\r\n </a>", "{\"a\":\"\"}"),
+ arguments("<a>\u00a0x\u00a0</a>", "{\"a\":\"\u00a0x\u00a0\"}"),
+ // Predefined entities and character references are expanded,
in text and in
+ // attribute values.
+ arguments("<a><>&'</a>", "{\"a\":\"<>&'\"}"),
+ arguments("<a>AB</a>", "{\"a\":\"AB\"}"),
+ arguments("<a b=\"<A\"/>", "{\"a\":{\"@b\":\"<A\"}}"),
+ // Entities declared in the DTD are expanded.
+ arguments(
+ "<!DOCTYPE a [<!ENTITY title
\"Dune\">]><a>&title;</a>",
+ "{\"a\":\"Dune\"}"));
+ }
+
+ //
--------------------------------------------------------------------------------------------
+ // XML Schema instance attributes
+ //
--------------------------------------------------------------------------------------------
+
+ @ParameterizedTest(name = "{0}")
+ @MethodSource("xsiMappings")
+ void testXsiMapping(String xml, String expectedJson) {
+ assertThat(parser.parse(xml, false).toJson()).isEqualTo(expectedJson);
+ }
+
+ static Stream<Arguments> xsiMappings() {
+ return Stream.of(
+ // xsi:nil makes an element without other attributes null, and
drops the content.
+ arguments("<a xsi:nil=\"true\"/>", "{\"a\":null}"),
+ arguments("<a xsi:nil=\"1\"></a>", "{\"a\":null}"),
+ arguments("<a xsi:nil=\"true\">x</a>", "{\"a\":null}"),
+ arguments("<a xsi:nil=\"true\"><b>x</b></a>", "{\"a\":null}"),
+ arguments("<a xsi:nil=\" true \" xsi:type=\"int\"/>",
"{\"a\":null}"),
+ arguments("<a xsi:nil=\"true\" xsi:type=\"Dog\"/>",
"{\"a\":null}"),
+ // The declaration of the xsi prefix is dropped.
+ arguments(
+ "<root
xmlns:xsi=\"http://www.w3.org/2001/XMLSchema-instance\">"
+ + "<a xsi:nil=\"true\"/></root>",
+ "{\"root\":{\"a\":null}}"),
+ // xsi:nil that doesn't make the element null is kept, unless
it is false.
+ arguments(
+ "<a xsi:nil=\"true\" id=\"1\"/>",
+ "{\"a\":{\"@id\":\"1\",\"@xsi:nil\":\"true\"}}"),
+ arguments(
+ "<a xsi:nil=\"true\" id=\"1\">x</a>",
+ "{\"a\":{\"@id\":\"1\",\"@xsi:nil\":\"true\"}}"),
+ arguments(
+ "<a xsi:nil=\" 1 \" id=\"1\"/>",
+ "{\"a\":{\"@id\":\"1\",\"@xsi:nil\":\"true\"}}"),
+ arguments(
+ "<a xsi:nil=\"true\" xsi:type=\"string\" id=\"1\"/>",
+
"{\"a\":{\"@id\":\"1\",\"@xsi:nil\":\"true\",\"@xsi:type\":\"string\"}}"),
+ arguments("<a xsi:nil=\"yes\"/>",
"{\"a\":{\"@xsi:nil\":\"yes\"}}"),
+ arguments("<a xsi:nil=\"TRUE\"/>",
"{\"a\":{\"@xsi:nil\":\"TRUE\"}}"),
+ arguments("<a xsi:nil=\"false\">x</a>", "{\"a\":\"x\"}"),
+ arguments("<a xsi:nil=\"0\"/>", "{\"a\":\"\"}"),
+ // xsi:type types the text of an element with text but without
child elements.
+ arguments("<a xsi:type=\"int\">5</a>", "{\"a\":5}"),
+ arguments("<a xsi:type=\"boolean\">true</a>", "{\"a\":true}"),
+ arguments("<a xsi:type=\"decimal\">-1.23</a>",
"{\"a\":-1.23}"),
+ arguments(
+ "<price xsi:type=\"decimal\"
currency=\"EUR\">13.50</price>",
+ "{\"price\":{\"$\":13.5,\"@currency\":\"EUR\"}}"),
+ // xsi:type that doesn't type the text is kept.
+ arguments("<a xsi:type=\"int\"></a>",
"{\"a\":{\"@xsi:type\":\"int\"}}"),
+ arguments("<a xsi:type=\"string\"/>",
"{\"a\":{\"@xsi:type\":\"string\"}}"),
+ arguments(
+ "<a id=\"1\" xsi:type=\"string\"/>",
+ "{\"a\":{\"@id\":\"1\",\"@xsi:type\":\"string\"}}"),
+ arguments(
+ "<a xsi:type=\"int\">five</a>",
+ "{\"a\":{\"$\":\"five\",\"@xsi:type\":\"int\"}}"),
+ arguments(
+ "<a xsi:type=\"xs:token\">x</a>",
+ "{\"a\":{\"$\":\"x\",\"@xsi:type\":\"xs:token\"}}"),
+ arguments(
+ "<animal xsi:type=\"Dog\" name=\"Rex\"/>",
+
"{\"animal\":{\"@name\":\"Rex\",\"@xsi:type\":\"Dog\"}}"),
+ arguments(
+ "<shape
xsi:type=\"Circle\"><radius>1</radius></shape>",
+
"{\"shape\":{\"@xsi:type\":\"Circle\",\"radius\":\"1\"}}"),
+ // Other xsi: attributes are regular attributes.
+ arguments(
+ "<a
xmlns:xsi=\"http://www.w3.org/2001/XMLSchema-instance\" "
+ + "xsi:schemaLocation=\"urn:a a.xsd\">x</a>",
+ "{\"a\":{\"$\":\"x\",\"@xsi:schemaLocation\":\"urn:a
a.xsd\"}}"));
+ }
+
+ @ParameterizedTest(name = "{0}: {1}")
+ @MethodSource("typedValues")
+ void testXsiType(String xsiType, String text, Variant.Type expectedType,
Object expectedValue) {
+ final Variant value = parseTyped(xsiType, text);
+ assertThat(value.getType()).isEqualTo(expectedType);
+ if (expectedValue instanceof BigDecimal) {
+ assertThat(value.getDecimal()).isEqualByComparingTo((BigDecimal)
expectedValue);
+ } else {
+ assertThat(value.get()).isEqualTo(expectedValue);
+ }
+ }
+
+ static Stream<Arguments> typedValues() {
+ return Stream.of(
+ arguments("string", "5", Variant.Type.STRING, "5"),
+ arguments("boolean", "true", Variant.Type.BOOLEAN, true),
+ arguments("boolean", "false", Variant.Type.BOOLEAN, false),
+ arguments("boolean", "1", Variant.Type.BOOLEAN, true),
+ arguments("boolean", "0", Variant.Type.BOOLEAN, false),
+ arguments("byte", "-128", Variant.Type.TINYINT, (byte) -128),
+ arguments("short", "32767", Variant.Type.SMALLINT, (short)
32767),
+ arguments("int", "+5", Variant.Type.INT, 5),
+ arguments("int", " 5 ", Variant.Type.INT, 5),
+ arguments("long", "9223372036854775807", Variant.Type.BIGINT,
Long.MAX_VALUE),
+ arguments(
+ "integer",
+ "123456789012345678901234567890",
+ Variant.Type.DECIMAL,
+ new BigDecimal("123456789012345678901234567890")),
+ arguments("decimal", "12.50", Variant.Type.DECIMAL, new
BigDecimal("12.50")),
Review Comment:
are these allowed? +0.0 or -0.0
##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/functions/XmlToVariantParser.java:
##########
@@ -0,0 +1,656 @@
+/*
+ * 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.annotation.Internal;
+import org.apache.flink.annotation.VisibleForTesting;
+import org.apache.flink.types.variant.BinaryVariantInternalBuilder;
+import org.apache.flink.types.variant.BinaryVariantInternalBuilder.FieldEntry;
+import org.apache.flink.types.variant.Variant;
+import org.apache.flink.types.variant.VariantTypeException;
+
+import javax.annotation.Nullable;
+import javax.xml.XMLConstants;
+import javax.xml.stream.XMLInputFactory;
+import javax.xml.stream.XMLStreamConstants;
+import javax.xml.stream.XMLStreamException;
+import javax.xml.stream.XMLStreamReader;
+
+import java.io.StringReader;
+import java.math.BigDecimal;
+import java.time.DateTimeException;
+import java.time.Instant;
+import java.time.LocalDate;
+import java.time.LocalDateTime;
+import java.time.LocalTime;
+import java.time.OffsetDateTime;
+import java.time.ZoneOffset;
+import java.time.format.DateTimeFormatter;
+import java.time.temporal.ChronoField;
+import java.time.temporal.TemporalAccessor;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.function.Consumer;
+import java.util.regex.Pattern;
+
+import static java.util.Map.entry;
+import static
org.apache.flink.types.variant.BinaryVariantInternalBuilder.toVariantDecimal;
+import static org.apache.flink.types.variant.BinaryVariantUtil.SIZE_LIMIT;
+import static
org.apache.flink.types.variant.BinaryVariantUtil.microsSinceEpoch;
+import static org.apache.flink.types.variant.BinaryVariantUtil.nanosSinceEpoch;
+
+/**
+ * Parses XML into a {@link Variant} for {@code PARSE_XML} and {@code
TRY_PARSE_XML}. The mapping is
+ * described in the documentation of the {@code VARIANT} data type.
+ *
+ * <p>Parsing has two steps. First, {@link #read} reads the document into a
tree of {@link
+ * XmlElement}s, which group the children of each element by name and apply
{@code xsi:nil} and
+ * {@code xsi:type}. Then, {@link #encode} writes the tree into a variant. The
variant can't be
+ * written while reading, because the occurrences of a repeated child become
one array, and they
+ * don't have to be next to each other in the document.
+ *
+ * <p>Instances are not thread-safe.
+ */
+@Internal
+public final class XmlToVariantParser {
+
+ private static final String ATTRIBUTE_PREFIX = "@";
+ private static final String TEXT_KEY = "$";
+ private static final String ORDER_KEY = "#";
+
+ private static final String XSI_NIL = "xsi:nil";
+ private static final String XSI_TYPE = "xsi:type";
+ private static final String XSI_NAMESPACE_DECLARATION = "xmlns:xsi";
+
+ private final XMLInputFactory inputFactory = createInputFactory();
+
+ @VisibleForTesting
+ Variant parse(String xml, boolean forceArray) {
+ return encode(read(xml), forceArray);
+ }
+
+ /**
+ * Reads the document into the tree of its root element. Calls on the same
document can share
+ * the tree, since {@link #encode} doesn't change it.
+ */
+ XmlElement read(String xml) {
+ try {
+ return readDocument(xml);
+ } catch (XMLStreamException e) {
+ // Only the message is kept, since the location of the exception
is not serializable.
+ throw new IllegalArgumentException(e.getMessage());
+ }
+ }
+
+ static Variant encode(XmlElement root, boolean forceArray) {
+ return new VariantEncoder(forceArray).encodeDocument(root);
+ }
+
+ private static XMLInputFactory createInputFactory() {
+ // The JDK implementation, regardless of other StAX implementations on
the classpath.
+ final XMLInputFactory factory = XMLInputFactory.newDefaultFactory();
+ factory.setProperty(XMLInputFactory.IS_NAMESPACE_AWARE, false);
+ factory.setProperty(XMLInputFactory.IS_COALESCING, true);
+ // Internal entities need DTD support. External entities are enabled
only so that they
+ // reach the resolver and fail, instead of being dropped silently. The
JDK refuses all
+ // protocols in case a load bypasses the resolver.
+ factory.setProperty(XMLInputFactory.SUPPORT_DTD, true);
+ factory.setProperty(XMLInputFactory.IS_SUPPORTING_EXTERNAL_ENTITIES,
true);
+ factory.setProperty(XMLConstants.ACCESS_EXTERNAL_DTD, "");
+ factory.setXMLResolver(
+ (publicId, systemId, baseUri, namespace) -> {
+ throw new XMLStreamException(
+ String.format(
+ "External entities and DTDs are not
supported: %s", systemId));
+ });
+ // Set here, so that they don't depend on the JDK version or JVM-wide
settings. Apart from
+ // the depth and the entity sizes, these are the defaults of JDK 11 to
21.
+ // Each element can add an object and an array to the variant, so 500
levels give at most
+ // the 1000 levels that PARSE_JSON allows.
+ setLimit(factory, "maxElementDepth", 500);
+ setLimit(factory, "entityExpansionLimit", 64_000);
+ setLimit(factory, "totalEntitySizeLimit", SIZE_LIMIT);
+ setLimit(factory, "maxGeneralEntitySizeLimit", SIZE_LIMIT);
+ setLimit(factory, "maxParameterEntitySizeLimit", 1_000_000);
+ setLimit(factory, "elementAttributeLimit", 10_000);
+ setLimit(factory, "maxXMLNameLimit", 1_000);
+ return factory;
+ }
+
+ private static void setLimit(XMLInputFactory factory, String limit, int
value) {
+ factory.setProperty("http://www.oracle.com/xml/jaxp/properties/" +
limit, value);
+ }
+
+ private XmlElement readDocument(String xml) throws XMLStreamException {
+ final XMLStreamReader reader = inputFactory.createXMLStreamReader(new
StringReader(xml));
+ try {
+ final String version = reader.getVersion();
+ if (!XmlVersion.isSupported(version)) {
+ throw new XMLStreamException(
+ String.format("XML %s documents are not supported.",
version));
+ }
+ while (reader.next() != XMLStreamConstants.START_ELEMENT) {
+ // Skip the prolog.
+ }
+ final XmlElement root = readElement(reader);
+ // Lets the reader reject content after the root element.
+ while (reader.hasNext()) {
+ reader.next();
+ }
+ return root;
+ } finally {
+ reader.close();
+ }
+ }
+
+ private static XmlElement readElement(XMLStreamReader reader) throws
XMLStreamException {
+ final XmlElement element =
+ new XmlElement(qualifiedName(reader.getPrefix(),
reader.getLocalName()));
+ readAttributes(reader, element);
+ final StringBuilder textRun = new StringBuilder();
+ for (int event = reader.next();
+ event != XMLStreamConstants.END_ELEMENT;
+ event = reader.next()) {
+ switch (event) {
+ case XMLStreamConstants.START_ELEMENT:
+ // A child element ends the current text run.
+ element.addTextRun(textRun);
+ element.addChild(readElement(reader));
+ break;
+ case XMLStreamConstants.CHARACTERS:
+ case XMLStreamConstants.CDATA:
+ case XMLStreamConstants.SPACE:
+ textRun.append(reader.getText());
+ break;
+ default:
+ // Comments and processing instructions don't end the text
run.
+ break;
+ }
+ }
+ element.addTextRun(textRun);
+ element.end();
+ return element;
+ }
+
+ private static void readAttributes(XMLStreamReader reader, XmlElement
element) {
+ for (int i = 0; i < reader.getAttributeCount(); i++) {
+ element.addAttribute(
+ qualifiedName(reader.getAttributePrefix(i),
reader.getAttributeLocalName(i)),
+ reader.getAttributeValue(i));
+ }
+ }
+
+ // The reader splits attribute names into prefix and local name, even if
not namespace-aware.
+ private static String qualifiedName(@Nullable String prefix, String
localName) {
+ return prefix == null || prefix.isEmpty() ? localName : prefix + ':' +
localName;
+ }
+
+ /**
+ * The XML versions that can be parsed. XML 1.1 isn't supported, since it
allows characters that
+ * XML 1.0 can't represent, e.g. most control characters. Accepting only
XML 1.0 keeps every
+ * result representable as XML 1.0.
+ */
+ private enum XmlVersion {
+ XML_1_0("1.0");
+
+ private final String version;
+
+ XmlVersion(String version) {
+ this.version = version;
+ }
+
+ /** A document without an XML declaration is XML 1.0. */
+ static boolean isSupported(@Nullable String version) {
+ return version == null
+ || Arrays.stream(values())
+ .anyMatch(supported ->
supported.version.equals(version));
+ }
+ }
+
+ //
--------------------------------------------------------------------------------------------
+
+ /** Parses the text of an element with {@code xsi:type}. */
+ private static final class XmlSchemaTypes {
+
+ // Parsing a BigDecimal is quadratic in the number of digits, and no
value of the types
+ // below needs that many characters.
+ private static final int MAX_LENGTH = 1000;
+
+ // The lexical spaces of XML Schema 1.1. The Java parsers accept more,
e.g. 1f for a float,
+ // non-ASCII digits, an exponent in a decimal, or a time without
seconds.
+
+ // https://www.w3.org/TR/xmlschema11-2/#integer, which byte, short,
int, and long inherit.
+ private static final Pattern INTEGER_FORM =
Pattern.compile("[\\-+]?[0-9]+");
+ // https://www.w3.org/TR/xmlschema11-2/#decimal
+ private static final Pattern DECIMAL_FORM =
+ Pattern.compile("(\\+|-)?([0-9]+(\\.[0-9]*)?|\\.[0-9]+)");
+ // https://www.w3.org/TR/xmlschema11-2/#float, the same form as for
double.
+ private static final Pattern FLOAT_FORM =
+ Pattern.compile(
+
"(\\+|-)?([0-9]+(\\.[0-9]*)?|\\.[0-9]+)([Ee](\\+|-)?[0-9]+)?|(\\+|-)?INF|NaN");
+ // https://www.w3.org/TR/xmlschema11-2/#date, without the time zone.
+ private static final String DATE_FORM =
+
"-?([1-9][0-9]{3,}|0[0-9]{3})-(0[1-9]|1[0-2])-(0[1-9]|[12][0-9]|3[01])";
+ // https://www.w3.org/TR/xmlschema11-2/#time, without the time zone.
+ private static final String TIME_FORM =
+
"(([01][0-9]|2[0-3]):[0-5][0-9]:[0-5][0-9](\\.[0-9]+)?|(24:00:00(\\.0+)?))";
+ // https://www.w3.org/TR/xmlschema11-2/#dateTime, the time zone of
date, time, and dateTime.
+ private static final String TIMEZONE_FORM =
+ "(Z|(\\+|-)((0[0-9]|1[0-3]):[0-5][0-9]|14:00))?";
+ private static final Map<String, Pattern> LEXICAL_FORMS =
+ Map.ofEntries(
+ entry("byte", INTEGER_FORM),
+ entry("short", INTEGER_FORM),
+ entry("int", INTEGER_FORM),
+ entry("long", INTEGER_FORM),
+ entry("integer", INTEGER_FORM),
+ entry("decimal", DECIMAL_FORM),
+ entry("float", FLOAT_FORM),
+ entry("double", FLOAT_FORM),
+ entry("date", Pattern.compile(DATE_FORM +
TIMEZONE_FORM)),
+ entry("time", Pattern.compile(TIME_FORM +
TIMEZONE_FORM)),
+ entry(
+ "dateTime",
+ Pattern.compile(DATE_FORM + "T" + TIME_FORM +
TIMEZONE_FORM)));
+
+ private XmlSchemaTypes() {}
+
+ /**
+ * Returns how to write the text as a value of the type, or null if
the type is unknown or
+ * the text isn't a valid value of it. The value is written later by
the {@link
+ * VariantEncoder}, but whether the type applies has to be known while
reading.
+ */
+ static @Nullable Consumer<BinaryVariantInternalBuilder> parse(String
text, String xsiType) {
+ final String type = localPart(xsiType.trim());
+ if (!hasLexicalForm(type, text)) {
+ return null;
+ }
+ try {
+ switch (type) {
+ case "string":
+ return builder -> builder.appendString(text);
+ case "boolean":
+ {
+ final Boolean value = parseBoolean(text);
+ return value == null ? null : builder ->
builder.appendBoolean(value);
+ }
+ case "byte":
+ {
+ final byte value = Byte.parseByte(text);
+ return builder -> builder.appendByte(value);
+ }
+ case "short":
+ {
+ final short value = Short.parseShort(text);
+ return builder -> builder.appendShort(value);
+ }
+ case "int":
+ {
+ final int value = Integer.parseInt(text);
+ return builder -> builder.appendInt(value);
+ }
+ case "long":
+ {
+ final long value = Long.parseLong(text);
+ return builder -> builder.appendLong(value);
+ }
+ case "integer":
+ case "decimal":
+ {
+ // The variant has no unbounded integer type, so
an integer is stored
+ // as a decimal.
+ final BigDecimal value = toVariantDecimal(new
BigDecimal(text));
+ return builder -> builder.appendDecimal(value);
+ }
+ case "float":
+ {
+ // Like in PARSE_JSON, NaN and infinity aren't
typed. A number out of
+ // range parses as infinity.
+ final float value = Float.parseFloat(text);
+ return Float.isFinite(value)
+ ? builder -> builder.appendFloat(value)
+ : null;
+ }
+ case "double":
+ {
+ final double value = Double.parseDouble(text);
+ return Double.isFinite(value)
+ ? builder -> builder.appendDouble(value)
+ : null;
+ }
+ case "date":
+ {
+ final int days =
Math.toIntExact(LocalDate.parse(text).toEpochDay());
+ return builder -> builder.appendDate(days);
+ }
+ case "time":
+ {
+ // A variant time has microsecond precision, so
nanoseconds are dropped
+ // like in a cast.
+ final long micros =
LocalTime.parse(text).toNanoOfDay() / 1_000;
+ return builder -> builder.appendTime(micros);
+ }
+ case "dateTime":
+ return
dateTime(DateTimeFormatter.ISO_DATE_TIME.parse(text));
+ default:
+ return null;
+ }
+ } catch (NumberFormatException
+ | DateTimeException
+ | ArithmeticException
+ | VariantTypeException e) {
+ return null;
+ }
+ }
+
+ private static String localPart(String type) {
+ return type.substring(type.lastIndexOf(':') + 1);
+ }
+
+ private static boolean hasLexicalForm(String type, String text) {
+ final Pattern form = LEXICAL_FORMS.get(type);
+ return form == null || (text.length() <= MAX_LENGTH &&
form.matcher(text).matches());
+ }
+
+ static @Nullable Boolean parseBoolean(String text) {
+ switch (text) {
+ case "true":
+ case "1":
+ return true;
+ case "false":
+ case "0":
+ return false;
+ default:
+ return null;
+ }
+ }
+
+ private static Consumer<BinaryVariantInternalBuilder>
dateTime(TemporalAccessor dateTime) {
+ final boolean hasOffset =
dateTime.isSupported(ChronoField.OFFSET_SECONDS);
+ final Instant instant =
+ hasOffset
+ ? OffsetDateTime.from(dateTime).toInstant()
+ :
LocalDateTime.from(dateTime).toInstant(ZoneOffset.UTC);
+ // Nanoseconds are kept if the timestamp is within the range of
nanosecond timestamps,
+ // from 1677 to 2262. Otherwise, they are dropped like in a cast.
+ if (instant.getNano() % 1_000 != 0) {
+ try {
+ final long nanos = nanosSinceEpoch(instant);
+ return hasOffset
+ ? builder -> builder.appendTimestampLtzNanos(nanos)
+ : builder -> builder.appendTimestampNanos(nanos);
+ } catch (VariantTypeException e) {
+ // Out of range, so the timestamp is written with
microseconds below.
+ }
+ }
+ final long micros = microsSinceEpoch(instant);
+ return hasOffset
+ ? builder -> builder.appendTimestampLtz(micros)
+ : builder -> builder.appendTimestamp(micros);
+ }
+ }
+
+ //
--------------------------------------------------------------------------------------------
+
+ /** An element as read from the document. */
+ static final class XmlElement {
+
+ final String name;
+
+ final Map<String, String> attributes = new LinkedHashMap<>();
+
+ // The child elements, grouped by name, and the text runs. Each knows
its position among
+ // the content, for the # field.
+ final Map<String, List<XmlElement>> children = new LinkedHashMap<>();
+ final List<TextRun> textRuns = new ArrayList<>();
+ private int contentCount;
+
+ // The position among the content of the parent element.
+ int position;
+
+ // Set by end() from xsi:nil and xsi:type.
+ boolean isNull;
+ @Nullable Consumer<BinaryVariantInternalBuilder> typedText;
+
+ private boolean xsiNil;
+ private @Nullable String xsiType;
+
+ XmlElement(String name) {
+ this.name = name;
+ }
+
+ void addAttribute(String name, String value) {
+ switch (name) {
+ case XSI_NAMESPACE_DECLARATION:
+ // The xsi prefix is recognized without the declaration.
+ break;
+ case XSI_NIL:
+ readXsiNil(value);
+ break;
+ case XSI_TYPE:
+ xsiType = value;
+ break;
+ default:
+ attributes.put(name, value);
+ }
+ }
+
+ /**
+ * xsi:nil="true" marks the element as nil, and xsi:nil="false" has no
effect. Any other
+ * value isn't a boolean, so it is kept as a regular attribute.
+ */
+ private void readXsiNil(String value) {
+ final Boolean nil = XmlSchemaTypes.parseBoolean(value.trim());
+ if (nil == null) {
+ attributes.put(XSI_NIL, value);
+ } else {
+ xsiNil = nil;
+ }
+ }
+
+ /** Adds the text run, unless it is whitespace only, and clears it for
the next one. */
+ void addTextRun(StringBuilder textRun) {
+ final String trimmed = textRun.toString().trim();
+ if (!trimmed.isEmpty()) {
+ textRuns.add(new TextRun(trimmed, contentCount++));
+ }
+ textRun.setLength(0);
+ }
+
+ void addChild(XmlElement child) {
+ child.position = contentCount++;
+ children.computeIfAbsent(child.name, name -> new
ArrayList<>()).add(child);
+ }
+
+ /**
+ * Applies xsi:nil and xsi:type once the element is read, since they
depend on all of it.
+ */
+ void end() {
+ if (xsiNil) {
+ applyXsiNil();
+ } else if (xsiType != null) {
+ applyXsiType();
+ }
+ }
+
+ /**
+ * xsi:nil="true" drops the content. The element is null, unless it
has other attributes.
+ * Then it keeps xsi:nil and xsi:type as attributes.
+ */
+ private void applyXsiNil() {
+ children.clear();
+ textRuns.clear();
+ if (attributes.isEmpty()) {
+ isNull = true;
+ } else {
+ attributes.put(XSI_NIL, "true");
+ if (xsiType != null) {
+ attributes.put(XSI_TYPE, xsiType);
+ }
+ }
+ }
+
+ /**
+ * xsi:type types the text of an element with text but without child
elements. Otherwise, or
+ * if the text isn't a valid value of the type, xsi:type is kept as an
attribute.
+ */
+ private void applyXsiType() {
+ if (children.isEmpty() && !textRuns.isEmpty()) {
+ typedText = XmlSchemaTypes.parse(text(), xsiType);
+ }
+ if (typedText == null) {
+ attributes.put(XSI_TYPE, xsiType);
+ }
+ }
+
+ /** Returns the text of an element without child elements. */
+ String text() {
+ return textRuns.isEmpty() ? "" : textRuns.get(0).text;
+ }
+ }
+
+ /** A text run of an element, which child elements separate from the next
one. */
+ private static final class TextRun {
+
+ final String text;
+
+ // The position among the content of the element.
+ final int position;
+
+ TextRun(String text, int position) {
+ this.text = text;
+ this.position = position;
+ }
+ }
+
+ //
--------------------------------------------------------------------------------------------
+
+ /**
+ * Writes a tree of {@link XmlElement}s into a variant.
+ *
+ * <p>The builder writes the values of an object or array first, and its
header last, in {@code
+ * finishWritingObject} or {@code finishWritingArray}. {@link #addField}
records where each
+ * field of an object starts, for the header.
+ */
+ private static final class VariantEncoder {
+
+ private final BinaryVariantInternalBuilder builder =
+ new BinaryVariantInternalBuilder(false);
+ private final boolean forceArray;
+
+ VariantEncoder(boolean forceArray) {
+ this.forceArray = forceArray;
+ }
+
+ /** The root element is a single field, never an array, since a
document has only one. */
+ Variant encodeDocument(XmlElement root) {
+ final int start = builder.getWritePos();
+ final ArrayList<FieldEntry> fields = new ArrayList<>(1);
+ addField(fields, start, root.name);
+ encodeElement(root);
+ builder.finishWritingObject(start, fields);
+ return builder.build();
+ }
+
+ /** An element without attributes and child elements is its text, any
other an object. */
+ private void encodeElement(XmlElement element) {
+ if (element.isNull) {
+ builder.appendNull();
+ } else if (element.attributes.isEmpty() &&
element.children.isEmpty()) {
+ encodeText(element, element.text());
+ } else {
+ encodeObject(element);
+ }
+ }
+
+ private void encodeObject(XmlElement element) {
+ final int start = builder.getWritePos();
+ final ArrayList<FieldEntry> fields = new ArrayList<>();
+ for (Map.Entry<String, String> attribute :
element.attributes.entrySet()) {
+ addField(fields, start, ATTRIBUTE_PREFIX + attribute.getKey());
+ builder.appendString(attribute.getValue());
+ }
+ for (Map.Entry<String, List<XmlElement>> children :
element.children.entrySet()) {
+ addField(fields, start, children.getKey());
+ encodeItems(children.getValue(), this::encodeElement);
+ }
+ if (!element.textRuns.isEmpty()) {
+ addField(fields, start, TEXT_KEY);
+ encodeItems(element.textRuns, textRun -> encodeText(element,
textRun.text));
+ }
+ // The # field records the order of the content across keys.
Within a key, it is the
+ // order of the array, so with a single key, there is nothing to
record.
+ final int contentKeys = element.children.size() +
(element.textRuns.isEmpty() ? 0 : 1);
+ if (contentKeys > 1) {
+ addField(fields, start, ORDER_KEY);
+ encodeOrder(element);
+ }
+ builder.finishWritingObject(start, fields);
+ }
+
+ private void encodeOrder(XmlElement element) {
+ final int start = builder.getWritePos();
+ final ArrayList<FieldEntry> fields = new ArrayList<>();
+ for (Map.Entry<String, List<XmlElement>> children :
element.children.entrySet()) {
+ addField(fields, start, children.getKey());
+ encodeItems(children.getValue(), child ->
builder.appendNumeric(child.position));
+ }
+ if (!element.textRuns.isEmpty()) {
+ addField(fields, start, TEXT_KEY);
+ encodeItems(element.textRuns, textRun ->
builder.appendNumeric(textRun.position));
+ }
+ builder.finishWritingObject(start, fields);
+ }
+
+ /**
+ * Writes the items of a key: a single item as is, and several items
as an array. With
+ * {@code forceArray}, a single item is an array as well.
+ */
+ private <T> void encodeItems(List<T> items, Consumer<T> encodeItem) {
+ if (items.size() == 1 && !forceArray) {
+ encodeItem.accept(items.get(0));
+ return;
+ }
+ final int start = builder.getWritePos();
+ final ArrayList<Integer> offsets = new ArrayList<>(items.size());
+ for (T item : items) {
+ offsets.add(builder.getWritePos() - start);
+ encodeItem.accept(item);
+ }
+ builder.finishWritingArray(start, offsets);
+ }
+
+ private void addField(ArrayList<FieldEntry> fields, int start, String
key) {
+ fields.add(new FieldEntry(key, builder.addKey(key),
builder.getWritePos() - start));
+ }
Review Comment:
This can be extracted to the builder and be reused:
```java
/**
* Adds a field to the object that starts at {@code start}. Call it right
before writing the
* value, since the field records the current write position as its offset.
*/
public void addField(ArrayList<FieldEntry> fields, int start, String key) {
fields.add(new FieldEntry(key, addKey(key), writePos - start));
}
```
If you re-base you can replace this in the `toVariantConverter` too
##########
flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/functions/XmlToVariantParserTest.java:
##########
@@ -0,0 +1,534 @@
+/*
+ * 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.types.variant.Variant;
+
+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 org.junit.jupiter.params.provider.ValueSource;
+import org.xml.sax.SAXException;
+
+import javax.xml.transform.stream.StreamSource;
+import javax.xml.validation.SchemaFactory;
+
+import java.io.StringReader;
+import java.math.BigDecimal;
+import java.time.Instant;
+import java.time.LocalDate;
+import java.time.LocalDateTime;
+import java.time.LocalTime;
+import java.util.stream.Stream;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatNoException;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import static org.junit.jupiter.params.provider.Arguments.arguments;
+
+/** Tests for {@link XmlToVariantParser}. */
+class XmlToVariantParserTest {
+
+ private final XmlToVariantParser parser = new XmlToVariantParser();
+
+ //
--------------------------------------------------------------------------------------------
+ // Mapping
+ //
--------------------------------------------------------------------------------------------
+
+ @ParameterizedTest(name = "{0}")
+ @MethodSource("mappings")
+ void testMapping(String xml, String expectedJson) {
+ assertThat(parser.parse(xml, false).toJson()).isEqualTo(expectedJson);
+ }
+
+ static Stream<Arguments> mappings() {
+ return Stream.of(
+ arguments("<title>Dune</title>", "{\"title\":\"Dune\"}"),
+ arguments("<title></title>", "{\"title\":\"\"}"),
+ arguments("<title/>", "{\"title\":\"\"}"),
+ arguments(
+ "<book pages=\"320\">\n Dune\n</book>",
+ "{\"book\":{\"$\":\"Dune\",\"@pages\":\"320\"}}"),
+ arguments("<book pages=\"320\"/>",
"{\"book\":{\"@pages\":\"320\"}}"),
+ arguments(
+ "<book>\n"
+ + " <author>Terry Pratchett</author>\n"
+ + " <author>Neil Gaiman</author>\n"
+ + "</book>",
+ "{\"book\":{\"author\":[\"Terry Pratchett\",\"Neil
Gaiman\"]}}"),
+ arguments(
+ "<book>\n <title>Dune</title>\n
<price>12.5</price>\n</book>",
+
"{\"book\":{\"#\":{\"price\":1,\"title\":0},\"price\":\"12.5\",\"title\":\"Dune\"}}"),
+ arguments(
+ "<book>\n"
+ + " This is a description\n"
+ + " <title>Dune</title>\n"
+ + " and this is another text\n"
+ + "</book>",
+ "{\"book\":{\"#\":{\"$\":[0,2],\"title\":1},"
+ + "\"$\":[\"This is a description\",\"and this
is another text\"],"
+ + "\"title\":\"Dune\"}}"),
+ arguments(
+ "<book pages=\"320\">\n"
+ + " <title>Dune</title>\n"
+ + " This is a description\n"
+ + " <price>12.50</price>\n"
+ + " Text in between\n"
+ + " <price currency=\"EUR\">13.50</price>\n"
+ + "</book>",
+
"{\"book\":{\"#\":{\"$\":[1,3],\"price\":[2,4],\"title\":0},"
+ + "\"$\":[\"This is a description\",\"Text in
between\"],"
+ + "\"@pages\":\"320\","
+ +
"\"price\":[\"12.50\",{\"$\":\"13.50\",\"@currency\":\"EUR\"}],"
+ + "\"title\":\"Dune\"}}"),
+ arguments("<a><b><c>deep</c></b></a>",
"{\"a\":{\"b\":{\"c\":\"deep\"}}}"),
+ // Namespace prefixes and declarations are kept as written.
+ arguments(
+ "<ns:a xmlns:ns=\"urn:ns\" xmlns=\"urn:default\"
ns:id=\"1\">"
+ + "<ns:b>x</ns:b></ns:a>",
+
"{\"ns:a\":{\"@ns:id\":\"1\",\"@xmlns\":\"urn:default\","
+ + "\"@xmlns:ns\":\"urn:ns\",\"ns:b\":\"x\"}}"),
+ // Attribute values are not trimmed.
+ arguments("<a b=\" x \"/>", "{\"a\":{\"@b\":\" x \"}}"),
+ // CDATA is text, and it merges with the text around it.
+ arguments("<a><![CDATA[<b>&</b>]]></a>",
"{\"a\":\"<b>&</b>\"}"),
+ arguments("<a>x<![CDATA[y]]>z</a>", "{\"a\":\"xyz\"}"),
+ // Comments and processing instructions are dropped, the text
around them merges.
+ arguments(
+ "<?xml version=\"1.0\"?><!-- c --><a>foo<!-- c
-->bar<?pi x?></a><!-- c -->",
+ "{\"a\":\"foobar\"}"),
+ // The input is a string, so an encoding declared in the
document is ignored.
+ arguments(
+ "<?xml version=\"1.0\"
encoding=\"ISO-8859-1\"?><a>\u00e4</a>",
+ "{\"a\":\"\u00e4\"}"),
+ // Whitespace-only text is dropped, other text is trimmed of
XML whitespace only.
+ arguments("<a>\n <b>x</b>\n</a>", "{\"a\":{\"b\":\"x\"}}"),
+ arguments("<a> \t\r\n </a>", "{\"a\":\"\"}"),
+ arguments("<a>\u00a0x\u00a0</a>", "{\"a\":\"\u00a0x\u00a0\"}"),
+ // Predefined entities and character references are expanded,
in text and in
+ // attribute values.
+ arguments("<a><>&'</a>", "{\"a\":\"<>&'\"}"),
+ arguments("<a>AB</a>", "{\"a\":\"AB\"}"),
+ arguments("<a b=\"<A\"/>", "{\"a\":{\"@b\":\"<A\"}}"),
+ // Entities declared in the DTD are expanded.
+ arguments(
+ "<!DOCTYPE a [<!ENTITY title
\"Dune\">]><a>&title;</a>",
+ "{\"a\":\"Dune\"}"));
+ }
+
+ //
--------------------------------------------------------------------------------------------
+ // XML Schema instance attributes
+ //
--------------------------------------------------------------------------------------------
+
+ @ParameterizedTest(name = "{0}")
+ @MethodSource("xsiMappings")
+ void testXsiMapping(String xml, String expectedJson) {
+ assertThat(parser.parse(xml, false).toJson()).isEqualTo(expectedJson);
+ }
+
+ static Stream<Arguments> xsiMappings() {
+ return Stream.of(
+ // xsi:nil makes an element without other attributes null, and
drops the content.
+ arguments("<a xsi:nil=\"true\"/>", "{\"a\":null}"),
+ arguments("<a xsi:nil=\"1\"></a>", "{\"a\":null}"),
+ arguments("<a xsi:nil=\"true\">x</a>", "{\"a\":null}"),
+ arguments("<a xsi:nil=\"true\"><b>x</b></a>", "{\"a\":null}"),
+ arguments("<a xsi:nil=\" true \" xsi:type=\"int\"/>",
"{\"a\":null}"),
+ arguments("<a xsi:nil=\"true\" xsi:type=\"Dog\"/>",
"{\"a\":null}"),
+ // The declaration of the xsi prefix is dropped.
+ arguments(
+ "<root
xmlns:xsi=\"http://www.w3.org/2001/XMLSchema-instance\">"
+ + "<a xsi:nil=\"true\"/></root>",
+ "{\"root\":{\"a\":null}}"),
+ // xsi:nil that doesn't make the element null is kept, unless
it is false.
+ arguments(
+ "<a xsi:nil=\"true\" id=\"1\"/>",
+ "{\"a\":{\"@id\":\"1\",\"@xsi:nil\":\"true\"}}"),
+ arguments(
+ "<a xsi:nil=\"true\" id=\"1\">x</a>",
+ "{\"a\":{\"@id\":\"1\",\"@xsi:nil\":\"true\"}}"),
+ arguments(
+ "<a xsi:nil=\" 1 \" id=\"1\"/>",
+ "{\"a\":{\"@id\":\"1\",\"@xsi:nil\":\"true\"}}"),
+ arguments(
+ "<a xsi:nil=\"true\" xsi:type=\"string\" id=\"1\"/>",
+
"{\"a\":{\"@id\":\"1\",\"@xsi:nil\":\"true\",\"@xsi:type\":\"string\"}}"),
+ arguments("<a xsi:nil=\"yes\"/>",
"{\"a\":{\"@xsi:nil\":\"yes\"}}"),
+ arguments("<a xsi:nil=\"TRUE\"/>",
"{\"a\":{\"@xsi:nil\":\"TRUE\"}}"),
+ arguments("<a xsi:nil=\"false\">x</a>", "{\"a\":\"x\"}"),
+ arguments("<a xsi:nil=\"0\"/>", "{\"a\":\"\"}"),
+ // xsi:type types the text of an element with text but without
child elements.
+ arguments("<a xsi:type=\"int\">5</a>", "{\"a\":5}"),
+ arguments("<a xsi:type=\"boolean\">true</a>", "{\"a\":true}"),
+ arguments("<a xsi:type=\"decimal\">-1.23</a>",
"{\"a\":-1.23}"),
+ arguments(
+ "<price xsi:type=\"decimal\"
currency=\"EUR\">13.50</price>",
+ "{\"price\":{\"$\":13.5,\"@currency\":\"EUR\"}}"),
+ // xsi:type that doesn't type the text is kept.
+ arguments("<a xsi:type=\"int\"></a>",
"{\"a\":{\"@xsi:type\":\"int\"}}"),
+ arguments("<a xsi:type=\"string\"/>",
"{\"a\":{\"@xsi:type\":\"string\"}}"),
+ arguments(
+ "<a id=\"1\" xsi:type=\"string\"/>",
+ "{\"a\":{\"@id\":\"1\",\"@xsi:type\":\"string\"}}"),
+ arguments(
+ "<a xsi:type=\"int\">five</a>",
+ "{\"a\":{\"$\":\"five\",\"@xsi:type\":\"int\"}}"),
+ arguments(
+ "<a xsi:type=\"xs:token\">x</a>",
+ "{\"a\":{\"$\":\"x\",\"@xsi:type\":\"xs:token\"}}"),
+ arguments(
+ "<animal xsi:type=\"Dog\" name=\"Rex\"/>",
+
"{\"animal\":{\"@name\":\"Rex\",\"@xsi:type\":\"Dog\"}}"),
+ arguments(
+ "<shape
xsi:type=\"Circle\"><radius>1</radius></shape>",
+
"{\"shape\":{\"@xsi:type\":\"Circle\",\"radius\":\"1\"}}"),
+ // Other xsi: attributes are regular attributes.
+ arguments(
+ "<a
xmlns:xsi=\"http://www.w3.org/2001/XMLSchema-instance\" "
+ + "xsi:schemaLocation=\"urn:a a.xsd\">x</a>",
+ "{\"a\":{\"$\":\"x\",\"@xsi:schemaLocation\":\"urn:a
a.xsd\"}}"));
+ }
+
+ @ParameterizedTest(name = "{0}: {1}")
+ @MethodSource("typedValues")
+ void testXsiType(String xsiType, String text, Variant.Type expectedType,
Object expectedValue) {
+ final Variant value = parseTyped(xsiType, text);
+ assertThat(value.getType()).isEqualTo(expectedType);
+ if (expectedValue instanceof BigDecimal) {
+ assertThat(value.getDecimal()).isEqualByComparingTo((BigDecimal)
expectedValue);
+ } else {
+ assertThat(value.get()).isEqualTo(expectedValue);
+ }
+ }
+
+ static Stream<Arguments> typedValues() {
+ return Stream.of(
+ arguments("string", "5", Variant.Type.STRING, "5"),
+ arguments("boolean", "true", Variant.Type.BOOLEAN, true),
+ arguments("boolean", "false", Variant.Type.BOOLEAN, false),
+ arguments("boolean", "1", Variant.Type.BOOLEAN, true),
+ arguments("boolean", "0", Variant.Type.BOOLEAN, false),
+ arguments("byte", "-128", Variant.Type.TINYINT, (byte) -128),
+ arguments("short", "32767", Variant.Type.SMALLINT, (short)
32767),
+ arguments("int", "+5", Variant.Type.INT, 5),
+ arguments("int", " 5 ", Variant.Type.INT, 5),
+ arguments("long", "9223372036854775807", Variant.Type.BIGINT,
Long.MAX_VALUE),
+ arguments(
+ "integer",
+ "123456789012345678901234567890",
+ Variant.Type.DECIMAL,
+ new BigDecimal("123456789012345678901234567890")),
+ arguments("decimal", "12.50", Variant.Type.DECIMAL, new
BigDecimal("12.50")),
+ arguments("decimal", ".5", Variant.Type.DECIMAL, new
BigDecimal("0.5")),
+ arguments("float", "1.5", Variant.Type.FLOAT, 1.5f),
+ arguments("double", "1.5E3", Variant.Type.DOUBLE, 1500.0),
+ arguments("double", "-.5e+2", Variant.Type.DOUBLE, -50.0),
+ arguments("date", "2026-09-23", Variant.Type.DATE,
LocalDate.of(2026, 9, 23)),
+ arguments(
+ "time",
+ "12:30:45.123456",
+ Variant.Type.TIME,
+ LocalTime.of(12, 30, 45, 123_456_000)),
+ // Precision that doesn't fit is dropped, like in a cast.
+ arguments(
+ "time",
+ "12:30:45.123456789",
+ Variant.Type.TIME,
+ LocalTime.of(12, 30, 45, 123_456_000)),
+ arguments(
+ "dateTime",
+ "2300-01-01T00:00:00.000000001",
+ Variant.Type.TIMESTAMP,
+ LocalDateTime.of(2300, 1, 1, 0, 0)),
+ arguments(
+ "dateTime",
+ "2026-09-23T12:30:45",
+ Variant.Type.TIMESTAMP,
+ LocalDateTime.of(2026, 9, 23, 12, 30, 45)),
+ arguments(
+ "dateTime",
+ "2026-09-23T12:30:45.123456789",
+ Variant.Type.TIMESTAMP_NS,
+ LocalDateTime.of(2026, 9, 23, 12, 30, 45,
123_456_789)),
+ arguments(
+ "dateTime",
+ "2026-09-23T12:30:45Z",
+ Variant.Type.TIMESTAMP_LTZ,
+ Instant.parse("2026-09-23T12:30:45Z")),
+ arguments(
+ "dateTime",
+ "2026-09-23T12:30:45.5+02:00",
+ Variant.Type.TIMESTAMP_LTZ,
+ Instant.parse("2026-09-23T10:30:45.5Z")),
+ arguments(
+ "dateTime",
+ "2026-09-23T12:30:45.123456789Z",
+ Variant.Type.TIMESTAMP_LTZ_NS,
+ Instant.parse("2026-09-23T12:30:45.123456789Z")),
+ // Types are matched by their local part.
+ arguments("xs:int", "5", Variant.Type.INT, 5),
+ arguments("xsd:int", "5", Variant.Type.INT, 5));
+ }
+
+ @ParameterizedTest(name = "{0}: {1}")
+ @MethodSource({"untypedValues", "invalidLexicalForms"})
+ void testXsiTypeThatDoesNotApply(String xsiType, String text) {
+ final Variant value = parseTyped(xsiType, text);
+ assertThat(value.getField("$").getString()).isEqualTo(text);
+ assertThat(value.getField("@xsi:type").getString()).isEqualTo(xsiType);
+ }
+
+ static Stream<Arguments> untypedValues() {
+ return Stream.of(
+ arguments("boolean", "yes"),
+ arguments("boolean", "TRUE"),
+ arguments("byte", "128"),
+ arguments("int", "5.0"),
+ arguments("integer", "1.5"),
+ // Longer numbers are not typed.
+ arguments("integer", "1" + "0".repeat(1000)),
+ arguments("decimal", "1e99999"),
+ arguments("decimal",
"1234567890123456789012345678901234567890"),
+ arguments("float", "1e50"),
+ arguments("float", "INF"),
+ arguments("double", "1e99999"),
+ arguments("double", "NaN"),
+ arguments("double", "Infinity"),
+ arguments("date", "2026-02-30"),
+ arguments("date", "2026-09-23Z"),
+ arguments("time", "12:30:45Z"),
+ arguments("dateTime", "2026-09-23"),
+ // Valid in XML Schema, but rejected by the Java parsers.
+ arguments("date", "10000-01-01"),
+ arguments("dateTime", "10000-01-01T00:00:00"),
+ arguments("time", "24:00:00"),
+ arguments("dateTime", "2026-09-23T24:00:00"),
+ arguments("anyURI", "urn:x"),
+ arguments("Int", "5"));
+ }
+
+ /** Values that the Java parsers accept, but XML Schema doesn't. */
+ static Stream<Arguments> invalidLexicalForms() {
+ return Stream.of(
+ arguments("float", "1f"),
+ arguments("float", "1.5F"),
+ arguments("float", "1d"),
+ arguments("double", "1.5D"),
+ arguments("float", "0x1.8p1"),
+ arguments("double", "0x1.8p1"),
+ // Digits from other scripts: Arabic-Indic three, and
fullwidth one and two.
+ arguments("int", "\u0663"), // ٣
+ arguments("long", "\uff11\uff12"), // 12
+ arguments("integer", "\u0663"), // ٣
+ arguments("decimal", "\u0663"), // ٣
+ arguments("decimal", "1e5"),
+ arguments("time", "12:30"),
+ arguments("dateTime", "2026-09-23T12:30"),
Review Comment:
```java
arguments("dateTime", "2026-09-23T12:30.1234567890123456789"),
```
What happens with a case like this?
##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/codegen/XmlReadReuseTest.java:
##########
@@ -0,0 +1,256 @@
+/*
+ * 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.planner.codegen;
+
+import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+import org.apache.flink.table.api.EnvironmentSettings;
+import org.apache.flink.table.api.TableEnvironment;
+import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
+import org.apache.flink.table.api.config.ExecutionConfigOptions;
+import org.apache.flink.table.api.config.OptimizerConfigOptions;
+import org.apache.flink.table.api.config.TableConfigOptions;
+import org.apache.flink.table.codesplit.JavaCodeSplitter;
+import org.apache.flink.table.planner.factories.TestValuesTableFactory;
+import org.apache.flink.table.runtime.functions.SqlXmlUtils;
+import org.apache.flink.types.Row;
+
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.time.Instant;
+import java.util.Arrays;
+import java.util.List;
+import java.util.regex.Pattern;
+import java.util.stream.Collectors;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/** Tests that the {@code PARSE_XML} and {@code TRY_PARSE_XML} calls on the
same input share it. */
+class XmlReadReuseTest {
+
+ private static final Pattern READ_PATTERN =
+ Pattern.compile(Pattern.quote(SqlXmlUtils.class.getCanonicalName()
+ ".read("));
+ private static final Pattern PARSE_XML_PATTERN =
+ Pattern.compile(Pattern.quote(SqlXmlUtils.class.getCanonicalName()
+ ".parseXml("));
+
+ private StreamTableEnvironment tEnv;
+
+ @BeforeEach
+ void setUp() {
+ tEnv =
+ StreamTableEnvironment.create(
+ StreamExecutionEnvironment.getExecutionEnvironment(),
+ EnvironmentSettings.inStreamingMode());
+ tEnv.createTemporaryView(
+ "xml_src",
+ tEnv.fromValues(
+ Row.of(
+ "<book><title>Dune</title></book>",
+ "<a>1</a>",
+ "{\"a\":1}",
+ "<book>"),
+ Row.of(
+ "<book><title>Emma</title></book>",
+ "<a>2</a>",
+ "<a/>",
+ "<a>"))
+ .as("xml", "other_xml", "json", "invalid_xml"));
+ }
+
+ @Test
+ void testCallsOnTheSameInputReadItOnce() {
+ final String sql =
+ "SELECT CAST(PARSE_XML(xml)['book']['title'] AS STRING), "
+ + "JSON_STRING(PARSE_XML(xml, TRUE)), "
+ + "JSON_STRING(TRY_PARSE_XML(xml, TRUE)) FROM xml_src";
+ assertThat(collect(sql))
+ .containsExactlyInAnyOrder(
+ Row.of(
+ "Dune",
+ "{\"book\":{\"title\":[\"Dune\"]}}",
+ "{\"book\":{\"title\":[\"Dune\"]}}"),
+ Row.of(
+ "Emma",
+ "{\"book\":{\"title\":[\"Emma\"]}}",
+ "{\"book\":{\"title\":[\"Emma\"]}}"));
+ assertThat(countReads(sql)).isOne();
+ }
+
+ @Test
+ void testCallsOnDifferentInputsReadEachOnce() {
+ final String sql =
+ "SELECT JSON_STRING(PARSE_XML(xml)),
JSON_STRING(TRY_PARSE_XML(xml, TRUE)), "
+ + "JSON_STRING(PARSE_XML(other_xml)), "
+ + "JSON_STRING(TRY_PARSE_XML(other_xml, TRUE)) FROM
xml_src";
+ assertThat(collect(sql))
+ .containsExactlyInAnyOrder(
+ Row.of(
+ "{\"book\":{\"title\":\"Dune\"}}",
+ "{\"book\":{\"title\":[\"Dune\"]}}",
+ "{\"a\":\"1\"}",
+ "{\"a\":\"1\"}"),
+ Row.of(
+ "{\"book\":{\"title\":\"Emma\"}}",
+ "{\"book\":{\"title\":[\"Emma\"]}}",
+ "{\"a\":\"2\"}",
+ "{\"a\":\"2\"}"));
+ assertThat(countReads(sql)).isEqualTo(2);
+ }
+
+ @Test
+ void testFirstCallInBranchThatIsNotTaken() {
+ // The first call is never taken, so the second one has to read the
document.
+ final String sql =
+ "SELECT CASE WHEN CHARACTER_LENGTH(other_xml) > 100 "
+ + "THEN CAST(PARSE_XML(xml)['book']['title'] AS
STRING) ELSE 'none' END, "
+ + "CAST(PARSE_XML(xml, TRUE)['book']['title'][1] AS
STRING) FROM xml_src";
+ assertThat(collect(sql))
+ .containsExactlyInAnyOrder(Row.of("none", "Dune"),
Row.of("none", "Emma"));
+ assertThat(countCalls(sql, PARSE_XML_PATTERN))
+ .as("the branch with the first call must survive the
optimizer")
+ .isEqualTo(2);
+ assertThat(countReads(sql)).isOne();
+ }
+
+ @Test
+ void testInvalidDocumentIsReadOnce() {
+ final String sql =
+ "SELECT JSON_STRING(TRY_PARSE_XML(invalid_xml)), "
+ + "JSON_STRING(TRY_PARSE_XML(invalid_xml, TRUE)) FROM
xml_src";
+ assertThat(collect(sql)).containsExactlyInAnyOrder(Row.of(null, null),
Row.of(null, null));
+ assertThat(countReads(sql)).isOne();
+ }
+
+ @Test
+ void testJsonFunctionOnTheSameInput() {
+ // The JSON functions share their parsed input as well, which must not
be mixed up.
+ final String sql =
+ "SELECT JSON_VALUE(json, '$.a'),
JSON_STRING(TRY_PARSE_XML(json)) FROM xml_src";
+ assertThat(collect(sql))
+ .containsExactlyInAnyOrder(Row.of("1", null), Row.of(null,
"{\"a\":\"\"}"));
+ assertThat(countReads(sql)).isOne();
+ }
+
+ @Test
+ void testReuseSurvivesCodeSplitting() {
+ tEnv.getConfig()
+ .set(TableConfigOptions.MAX_LENGTH_GENERATED_CODE, 1)
+ .set(TableConfigOptions.MAX_MEMBERS_GENERATED_CODE, 1);
+ final String sql =
+ "SELECT CAST(PARSE_XML(xml)['book']['title'] AS STRING), "
+ + "JSON_STRING(PARSE_XML(xml, TRUE)), "
+ + "JSON_STRING(TRY_PARSE_XML(xml, TRUE)) FROM xml_src";
+ // Only the split code is compiled, so correct results show that the
reuse survives it.
+ assertThat(collect(sql))
+ .containsExactlyInAnyOrder(
+ Row.of(
+ "Dune",
+ "{\"book\":{\"title\":[\"Dune\"]}}",
+ "{\"book\":{\"title\":[\"Dune\"]}}"),
+ Row.of(
+ "Emma",
+ "{\"book\":{\"title\":[\"Emma\"]}}",
+ "{\"book\":{\"title\":[\"Emma\"]}}"));
+ final List<String> splitCodes =
+ GeneratedCodeTestUtils.generatedClassCodes(tEnv, sql).stream()
+ .map(code -> JavaCodeSplitter.split(code, 1, 1))
+ .collect(Collectors.toList());
+ assertThat(GeneratedCodeTestUtils.generatedClassCodes(tEnv, sql))
+ .as("the limits must actually split a generated class")
+ .anySatisfy(
+ code -> assertThat(JavaCodeSplitter.split(code, 1,
1)).isNotEqualTo(code));
+ assertThat(
+ splitCodes.stream()
+ .mapToInt(
+ code ->
+
GeneratedCodeTestUtils.countMatches(
+ READ_PATTERN, code))
+ .sum())
+ .isOne();
+ }
+
+ @Test
+ void testDocumentIsReadForEachRowInMatchRecognize() {
+ // A matches the first three rows, so a read for each row gives 1 + 2
+ 3.
+ final String dataId =
+ TestValuesTableFactory.registerData(
+ Arrays.asList(
+ Row.of(1000, "<n>1</n>",
Instant.ofEpochMilli(1000L)),
+ Row.of(2000, "<n>2</n>",
Instant.ofEpochMilli(2000L)),
+ Row.of(3000, "<n>3</n>",
Instant.ofEpochMilli(3000L)),
+ Row.of(9000, "<n>9</n>",
Instant.ofEpochMilli(9000L))));
+ tEnv.executeSql(
+ "CREATE TABLE events (f0 INT, f1 STRING, ts TIMESTAMP_LTZ(3), "
+ + "WATERMARK FOR ts AS ts) WITH ('connector' =
'values', 'data-id' = '"
+ + dataId
+ + "', 'bounded' = 'true')");
+ final String sql =
+ "SELECT total FROM events MATCH_RECOGNIZE ("
+ + " ORDER BY ts"
+ + " MEASURES SUM(CAST(CAST(PARSE_XML(A.f1)['n'] AS
STRING) AS INT)) AS total"
+ + " AFTER MATCH SKIP PAST LAST ROW"
+ + " PATTERN (A+ B)"
+ + " DEFINE A AS A.f0 < 9000, B AS B.f0 >= 9000)";
+ assertThat(collect(sql)).containsExactly(Row.of(6));
Review Comment:
Should we call at the end `TestValuesTableFactory.clearAllData()`?
--
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]