Josh Rosen created SPARK-58486:
----------------------------------

             Summary: InjectRuntimeFilter builds runtime bloom filters with 
collation-agnostic XxHash64, silently dropping join rows when a collated 
equi-join key is a non-attribute expression
                 Key: SPARK-58486
                 URL: https://issues.apache.org/jira/browse/SPARK-58486
             Project: Spark
          Issue Type: Bug
          Components: SQL
    Affects Versions: 4.2.0, 4.0.1, 4.1.0
            Reporter: Josh Rosen


{{InjectRuntimeFilter}} builds and probes its runtime bloom filter with the 
plain {{XxHash64}} 
([InjectRuntimeFilter.scala#L62-L79|https://github.com/apache/spark/blob/c7b2f1a865cce621d2c702bdfdff5132e251ca74/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/InjectRuntimeFilter.scala#L62-L79]).
 Since SPARK-52828 (4.0.1 / 4.1.0), {{XxHash64}} hashes collated strings by raw 
UTF-8 bytes while the join's equality remains collation-aware, so for a 
collated equi-join whose key is a non-attribute expression, two values that are 
equal under the join's collation but byte-different (e.g. {{'abc'}} vs 
{{'ABC'}} under {{{}UTF8_LCASE{}}}) hash differently, 
{{BloomFilterMightContain}} returns false, and matching probe-side rows are 
filtered out before the join runs: silently missing output rows under default 
configs (runtime bloom filters are on by default).
h3. Reproduction

Reproduced on current master (c7b2f1a865c). Default configs except the two 
thresholds below, which only make the tiny tables eligible for injection (they 
gate sizes, not semantics):
{code:sql}
CREATE TABLE dim(key STRING COLLATE UTF8_LCASE, cat STRING) USING parquet;
INSERT INTO dim VALUES ('ABC', 'x');

CREATE TABLE fact USING parquet AS
SELECT concat('abc', CAST(id AS STRING)) COLLATE UTF8_LCASE AS k FROM 
range(5000);

SET spark.sql.autoBroadcastJoinThreshold=-1;
SET spark.sql.optimizer.runtime.bloomFilter.applicationSideScanSizeThreshold=1;

-- Every fact row satisfies substr(k, 1, 3) = 'abc', and 'abc' = 'ABC' under 
UTF8_LCASE,
-- so this should return 5000:
SELECT count(*) FROM fact f JOIN dim d ON substr(f.k, 1, 3) = d.key WHERE d.cat 
= 'x';
{code}
{noformat}
spark.sql.optimizer.runtime.bloomFilter.enabled=true  (default) -> 0     (WRONG)
spark.sql.optimizer.runtime.bloomFilter.enabled=false           -> 5000  
(correct)
{noformat}
Toggling only {{spark.sql.optimizer.runtime.bloomFilter.enabled}} (which gates 
nothing but {{InjectRuntimeFilter)}} flips the result, and the optimized plan 
of the failing run contains the injected {{might_contain(..., xxhash64(...))}} 
probe (the {{BloomFilterAggregate}} build side sits inside a 
{{{}ScalarSubquery{}}}).
h3. Root cause

SPARK-52828 
([f62724d2d8d|https://github.com/apache/spark/commit/f62724d2d8dc0f405c9d0dedcd0136bf4a4aa3b7],
 first released in 4.1.0; backported to branch-4.0 as 
[542faa70cef|https://github.com/apache/spark/commit/542faa70cefe6526387f2bfaebde546e4e562ace],
 first released in 4.0.1) made {{Murmur3Hash}} / {{XxHash64}} 
collation-agnostic by default (raw-byte hashing for collated strings unless 
{{{}spark.sql.legacy.collationAwareHashFunctions=true{}}}, 
[SQLConf.scala#L1263-L1270|https://github.com/apache/spark/blob/c7b2f1a865cce621d2c702bdfdff5132e251ca74/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala#L1263-L1270])
 and migrated the internal consumers to collation-aware variants – shuffle 
partitioning to {{{}CollationAwareMurmur3Hash{}}}, hash-join keys via 
{{CollationKey}} – but did not update {{{}InjectRuntimeFilter{}}}, which still 
emits plain {{XxHash64}} on both the build side (L64/L66) and the probe side 
(L79).

This is masked in the common case: {{RewriteCollationJoin}} 
([RewriteCollationJoin.scala#L38-L41|https://github.com/apache/spark/blob/c7b2f1a865cce621d2c702bdfdff5132e251ca74/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/RewriteCollationJoin.scala#L38-L41])
 rewrites join equalities to {{CollationKey(l) = CollationKey(r)}} before the 
optimizer runs, so for those keys the bloom filter hashes collation-key bytes 
on both sides and stays correct. But its pattern matches only {{{}Equality(l: 
AttributeReference, r: AttributeReference){}}}: a non-attribute collated key 
({{{}substr(k, 1, 3){}}}, {{{}concat(a, b){}}}, {{{}upper(k){}}}, ...) escapes 
the wrap. The join itself still compares correctly 
({{{}HashJoin.injectCollationKey{}}} normalizes keys at planning, 
[HashJoin.scala#L745-L752|https://github.com/apache/spark/blob/c7b2f1a865cce621d2c702bdfdff5132e251ca74/sql/core/src/main/scala/org/apache/spark/sql/execution/joins/HashJoin.scala#L745-L752]),
 which is why the raw-byte pre-filter produces wrong results.
h3. Related tickets
 * SPARK-52828 is the introducing change; on 4.0.0 and earlier {{XxHash64}} was 
collation-aware and this path was correct.
 * SPARK-57055 originally reported this same defect and was retracted by its 
reporter on the grounds that {{RewriteCollationJoin}} pre-wraps join keys in 
{{{}CollationKey{}}}, making the bloom filter symmetric. That reasoning holds 
only for bare attribute-to-attribute equalities (the rule's pattern); the 
reproduction above – a non-attribute join key – is the case it does not cover. 
That ticket has since been rescoped to documenting the (by-design) byte-wise 
behavior of {{{}DataFrameStatFunctions.bloomFilter{}}}, an adjacent but 
distinct path.
 * SPARK-48000 introduced {{RewriteCollationJoin}} (whose attribute-only 
pattern bounds the exposure); SPARK-50228 moved it ahead of the optimizer, 
which is what protects the bare-attribute case.

 



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to