raminqaf opened a new pull request, #29411:
URL: https://github.com/apache/flink/pull/29411
## What is the purpose of the change
Flink compares VARIANT values by their binary encoding, but one value has
many encodings. Object keys can be stored in any order, a field read from an
object keeps the key dictionary of its parent, and a number keeps its storage
width. So queries that group, deduplicate, join or compare VARIANT values
return wrong results in streaming, without an error.
```sql
-- t(s STRING) has three rows: {"a":1,"b":2}, {"b":2,"a":1}, {"a":1,"c":3}
SELECT a, COUNT(*) FROM (SELECT PARSE_JSON(s)['a'] AS a FROM t) GROUP BY a;
```
| Query | Before,
streaming | After
|
|----------------------------------------------------------|-------------------------------------------|------------------------------------------------------------------------------------------------------------------------------------------------------|
| the `GROUP BY` above | three groups,
`+I[1, 1]` each | `Column 'a' of type VARIANT cannot be used as a
grouping key, because the type has no equality. Cast the value to a comparable
type first.` |
| `... WHERE PARSE_JSON(s)['a'] = PARSE_JSON('1')` | 0 rows
| `An expression of type VARIANT cannot be used in a
comparison with '=', because the type has no equality. Cast the value to a
comparable type first.` |
| `ARRAY_CONTAINS(ARRAY[PARSE_JSON(s)['a']], PARSE_JSON('1'))` | `false`
| `Type 'VARIANT' should support 'EQUALS' comparison
with itself.`
|
| `... GROUP BY CAST(PARSE_JSON(s)['a'] AS INT)` | one group,
count 3 | unchanged
|
Batch already rejected most of these queries late, during code generation,
with `Type 'VARIANT' cannot be ordered, so it cannot be used as a key for
sorting, grouping, or joining`. It now gets the same error as streaming.
This follows Spark and Databricks, which also reject VARIANT where equality
or ordering is needed. Users cast to a concrete type first, as in the last row.
## Brief change log
- `LogicalTypeChecks#isComparableKeyType` decides which types can be
grouping, join, partition and sort keys and operands of comparisons. VARIANT is
not, also when nested in ROW, ARRAY or MAP.
- A new `KeyTypeValidator` checks the logical plan before optimization, so
it covers SQL and the Table API. It rejects such types in grouping keys,
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 and IN subqueries. Joins are covered through
the comparisons in their conditions.
- `LogicalTypeChecks#areComparable` treats VARIANT as not comparable.
ARRAY_CONTAINS, ARRAY_DISTINCT, ARRAY_POSITION, ARRAY_REMOVE, ARRAY_UNION,
ARRAY_EXCEPT, ARRAY_INTERSECT and MAP_CONTAINS_KEY declare in their input type
strategy that they compare elements or keys. Their runtime already compares
through `isEqual` evaluators, so for other types the error only moves from code
generation to validation.
- The VARIANT section of the data types page lists where VARIANT cannot be
compared.
The check runs on the logical plan rather than on exec nodes. It sees the
keys the query asks for, and keys the optimizer adds later, for example during
decorrelation, only match a value with copies of itself. A job restored from a
compiled plan skips the check and keeps running. The exception is a plan that
calls one of the functions above on VARIANT elements or keys: these functions
build their equality check again during code generation, so such a job now
fails when it starts.
Not in this PR: VARIANT as a PRIMARY KEY, a MAP key or a MULTISET element.
Equality of MAP and MULTISET values as grouping keys has a related streaming
bug and needs its own issue.
## Verifying this change
This change added tests and can be verified as follows:
- `KeyTypeValidatorTest` runs rejected and allowed queries for each clause
in streaming and batch mode, plus Table API cases.
- `VariantSemanticTest` turns the reproductions from the JIRA issue into
failing-SQL checks and keeps the CAST workaround as a passing program.
- `CollectionFunctionsITCase` and `MapFunctionITCase` cover VARIANT elements
and keys for each affected function in SQL and the Table API.
- `ComparableInputTypeStrategyTest`,
`EqualsComparableElementArgumentTypeStrategyTest` and `LogicalTypeChecksTest`
cover the type checks.
## Does this pull request potentially affect one of the following parts:
- Dependencies (does it add or upgrade a dependency): no
- The public API, i.e., is any changed class annotated with
`@Public(Evolving)`: yes. `BuiltInFunctionDefinitions` changes the input type
strategy of the eight functions above. No signature changes. Queries that used
VARIANT in these places now fail, which needs a release note.
- The serializers: no
- The runtime per-record code paths (performance sensitive): no
- Anything that affects deployment or recovery: JobManager (and its
components), Checkpointing, Kubernetes/Yarn, ZooKeeper: no
- The S3 file system connector: no
## Documentation
- Does this pull request introduce a new feature? no
- If yes, how is the feature documented? docs, the VARIANT section of the
data types page
---
##### Was generative AI tooling used to co-author this PR?
- [ ] Yes (please specify the tool below)
Generated-by: [Tool Name and Version]
--
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]