Skip to content

array_contains, arrays_overlap, array_distinct and array_union ignore string collation #6470

Description

@viirya

Describe the bug

On Spark 4.x with the default spark.comet.exec.scalaUDF.codegen.enabled=true, array_contains, arrays_overlap, array_distinct and array_union return wrong results when their array elements are strings with a non-UTF8_BINARY collation.

A cast to a collated string type has no native path, so Comet runs it through the JVM codegen dispatcher. The array(...) above it and the four array functions still run natively, and their native kernels compare strings by raw bytes, so the collation is ignored.

Steps to reproduce

Parquet table with columns _1, _2 and rows ("a", "A"), ("x ", "x"), ("b", "c"):

SELECT array_contains(array(CAST(_1 AS STRING COLLATE UTF8_LCASE)), CAST(_2 AS STRING COLLATE UTF8_LCASE)) FROM t;
SELECT arrays_overlap(array(CAST(_1 AS STRING COLLATE UTF8_LCASE)), array(CAST(_2 AS STRING COLLATE UTF8_LCASE))) FROM t;
  • UTF8_LCASE: Spark returns true for ('a', 'A'), Comet returns false.
  • UTF8_BINARY_RTRIM: Spark returns true for ('x ', 'x'), Comet returns false.
  • UNICODE_CI behaves like UTF8_LCASE.
  • array_distinct(array(a, b)) and array_union(array(a), array(b)) keep both 'a' and 'A' (or 'x ' and 'x') where Spark keeps one.

The plan shows cast dispatched and array plus the array function native.

Expected behavior

The results match Spark. Other collation-sensitive expressions already route collated input through the codegen dispatcher (for example =, IN, array_intersect, array_except, array_max) or fall back (array_position, array_remove).

Additional context

Reproduced on Spark 4.1.

Metadata

Metadata

Assignees

Labels

area:expressionsExpression evaluationbugSomething isn't workingpriority:criticalData corruption, silent wrong results, security issues

Type

No type

Projects

No projects

    Milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions