[
https://issues.apache.org/jira/browse/FLINK-40910?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Ramin Gharib updated FLINK-40910:
---------------------------------
Release Note: VARIANT values can no longer be grouped, deduplicated,
joined, sorted or compared. Flink compared them by their binary encoding, and
one VARIANT value has many encodings, so these queries returned wrong results
in streaming. This affects grouping keys and SELECT DISTINCT, DISTINCT
aggregates, UNION, INTERSECT and EXCEPT except UNION ALL, ORDER BY, the
PARTITION BY and ORDER BY of OVER windows, MATCH_RECOGNIZE and table arguments,
comparison operators, IN and join conditions, and ARRAY_CONTAINS,
ARRAY_DISTINCT, ARRAY_POSITION, ARRAY_REMOVE, ARRAY_UNION, ARRAY_EXCEPT,
ARRAY_INTERSECT and MAP_CONTAINS_KEY on VARIANT elements or keys. Such queries
now fail before they run. Cast the value to a concrete type first, for example
CAST(v['a'] AS INT). A job restored from a compiled plan keeps running, unless
it calls one of these functions on VARIANT elements or keys. Such a job fails
when it starts.
> Reject VARIANT where SQL needs equality. Byte comparison gives wrong results
> in GROUP BY, DISTINCT, joins and =
> ---------------------------------------------------------------------------------------------------------------
>
> Key: FLINK-40910
> URL: https://issues.apache.org/jira/browse/FLINK-40910
> Project: Flink
> Issue Type: Bug
> Components: Table SQL / Planner, Table SQL / Runtime
> Reporter: Ramin Gharib
> Assignee: Ramin Gharib
> Priority: Major
>
> h2. Problem
> Flink compares VARIANT values by their binary encoding. One logical value has
> many encodings. Queries that group, deduplicate, join or compare on VARIANT
> therefore return wrong results, without any error.
> {code:sql}
> -- t(s STRING) has three rows:
> -- row 1: \{"a":1,"b":2}
> -- row 2: \{"b":2,"a":1}
> -- row 3: \{"a":1,"c":3}
> SELECT a, COUNT(*) FROM (SELECT PARSE_JSON(s)['a'] AS a FROM t) GROUP BY a;
> {code}
> {noformat}
> Query Expected
> Actual
> GROUP BY PARSE_JSON(s)['a'] one group, count
> 3 three groups, count 1 each
> COUNT(DISTINCT PARSE_JSON(s)) over row 1 and row 2 1
> 2
> WHERE PARSE_JSON(s)['a'] = PARSE_JSON('1') 3 rows
> 0 rows
> l JOIN r ON PARSE_JSON(l.s)['id'] = PARSE_JSON(r.s)['id'] 1 row
> 0 rows
> with l = \{"id":7,"x":1} and r = \{"y":2,"id":7}
> {noformat}
> The actual results come from VariantSemanticTest on master 5afe863af62.
> Example output for the GROUP BY:
> {noformat}
> Expecting actual:
> ["+I[1, 1]", "+I[1, 1]", "+I[1, 1]"]
> to contain exactly in any order:
> ["+I[1, 1]", "-U[1, 1]", "+U[1, 2]", "-U[1, 2]", "+U[1, 3]"]
> {noformat}
> h2. Cause
> A VARIANT is two byte arrays. The metadata is a dictionary of key names. The
> value refers to keys by their position in that dictionary. Flink treats two
> VARIANTs as equal only if both arrays are byte-equal:
> * Grouping, DISTINCT and join keys are BinaryRowData.
> \{{AbstractBinaryWriter#writeVariant}} writes the dictionary and the value
> bytes into the key, and BinaryRowData hashes and compares those bytes.
> * \{{=}} calls \{{BinaryVariant#equals}}, which compares the value, the
> dictionary and the position.
> One value has many valid encodings:
> * The dictionary lists keys in the order they first appear. Row 1 and row 2
> get different dictionaries.
> * A field taken by field access keeps the full dictionary of its parent
> object. The number 1 from row 1 differs from the number 1 from row 3 and from
> PARSE_JSON('1').
> * Numbers keep their storage width. PARSE_JSON('1') is a TINYINT, CAST(CAST(1
> AS BIGINT) AS VARIANT) is a BIGINT, and PARSE_JSON('1.0') is a DECIMAL. All
> three print as 1.
> The variant encoding spec allows all of these. For the same reason, the
> Iceberg spec defines no hash for variant.
> h2. Proposal
> Reject VARIANT where SQL needs equality, like ORDER BY on VARIANT is rejected
> today (\{{TypeCheckUtils#isComparable}}). Users cast to a concrete type
> first. That already works:
> {code:sql}
> SELECT a, COUNT(*) FROM (SELECT CAST(PARSE_JSON(s)['a'] AS INT) AS a FROM t)
> GROUP BY a;
> -- one group, count 3
> {code}
> Reject VARIANT, also when nested in ROW, ARRAY or MAP, in:
> * GROUP BY keys, including GROUPING SETS, ROLLUP, CUBE and window aggregations
> * DISTINCT aggregates and SELECT DISTINCT
> * equi-join keys
> * =, <>, IN, IS DISTINCT FROM and IS NOT DISTINCT FROM between VARIANT
> operands
> * UNION, INTERSECT and EXCEPT without ALL
> * PARTITION BY of OVER windows, Top-N, deduplication, process table functions
> and MATCH_RECOGNIZE
> Still allowed: VARIANT as an aggregate argument such as COUNT(v),
> FIRST_VALUE(v) and LAST_VALUE(v), IS NULL, IS NOT NULL, field access and
> casts.
> The error should name the clause and suggest a cast, for example:
> {noformat}
> VARIANT cannot be used as a grouping key because VARIANT values have no
> equality. Cast the value to a concrete type first, for example CAST(v['a'] AS
> INT).
> {noformat}
> Open question: VARIANT as a MAP key. MAP lookups use the same byte equality.
> h2. Other systems
> ||System||Behavior||
> |Spark / Databricks|VARIANT is not comparable since SPARK-47569. GROUP BY
> fails with GROUP_EXPRESSION_TYPE_IS_NOT_ORDERABLE. DISTINCT and set
> operations fail. PARTITION BY fails with PARTITION_BY_VARIANT. Databricks
> docs: "The VARIANT data type cannot be used for comparisons, grouping,
> ordering, and set operations."|
> |BigQuery|JSON is neither groupable nor comparable, and cannot be a partition
> or cluster key. Users extract with JSON_VALUE first.|
> |Apache Iceberg|The spec defines no hash for variant, because equivalent
> values have several representations.|
> Sources:
> * https://docs.databricks.com/aws/en/semi-structured/variant
> *
> https://github.com/apache/spark/blob/master/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/ExprUtils.scala
> (checkValidGroupingExprs)
> *
> https://github.com/apache/spark/blob/master/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/CheckAnalysis.scala
> (variantColumnInSetOperation, PARTITION_BY_VARIANT)
> * https://cloud.google.com/bigquery/docs/reference/standard-sql/data-types
> * https://github.com/apache/iceberg/blob/main/format/spec.md
> h2. Compatibility
> * Flink 2.1 to 2.3 accept these queries. After the change they fail at
> validation. This needs a release note.
> * Existing tests assert the old behavior and must change: VariantSemanticTest
> VARIANT_AS_AGG_KEY, and the COUNT(DISTINCT v) in BUILTIN_AGG and
> BUILTIN_AGG_WITH_RETRACTION.
> * Open question: a job restored from a compiled plan skips SQL validation.
> Should it fail too, or keep running?
> * Supporting equality later stays compatible, because it only accepts more
> queries.
> h2. Alternatives considered
> * Write a field value with only the keys it uses. This fixes scalar fields,
> but not objects or number widths.
> * A canonical encoding: sorted dictionary, no unused keys, smallest number
> width. This changes the bytes of existing state, and data from other writers
> still differs because the spec allows several encodings.
> * Semantic equality and hashing on decoded values. Correct, but a larger
> design: how 1 compares to 1.0, NaN, and the timestamp kinds. It can follow
> later.
> h2. Tests
> Four VariantSemanticTest programs reproduce the bug and fail on master as
> shown above: variant-field-as-agg-key, variant-object-as-distinct-key,
> variant-field-equals-variant and variant-field-as-join-key. With the fix they
> become failing-SQL checks for the new error. A fifth program,
> variant-field-cast-as-agg-key, checks the CAST workaround and passes today.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)