David Mollitor created SPARK-59898:
--------------------------------------

             Summary: freqItems treats equal BinaryType values as distinct
                 Key: SPARK-59898
                 URL: https://issues.apache.org/jira/browse/SPARK-59898
             Project: Spark
          Issue Type: Improvement
          Components: SQL
    Affects Versions: 4.2.0, 3.5.0, 2.4.0, 1.4.0
            Reporter: David Mollitor


h2. Summary

{{DataFrame.stat.freqItems}} compares {{BinaryType}} values by reference, not 
by content. Equal binary values never match, so their counts never accumulate. 
The result then depends on the number and order of rows rather than on how 
frequent a value is: a value can be returned several times, or be missing even 
when it is in every row. That breaks the documented guarantee of possible false 
positives but no false negatives.
h2. Reproduction
{code:scala}
import spark.implicits._

Seq(11, 12).foreach { n =>
  // Every row holds the same binary value.
  val df = Seq.fill(n)(Array[Byte](1, 2)).toDF("b").coalesce(1)
  val items = df.stat.freqItems(Seq("b"), 
0.5).collect().head.getSeq[Array[Byte]](0)
  println(s"n=$n: ${items.map(_.toSeq)}")
}
// n=11: ArraySeq(ArraySeq(1, 2), ArraySeq(1, 2))   <- the same value twice
// n=12: ArraySeq()                                <- the value in 100% of the 
rows is missing
{code}
The same data as strings ({{{}Seq.fill(n)("ab"){}}}) correctly returns {{[ab]}} 
for every {{{}n{}}}.
h2. Cause

{{{}CollectFrequentItems{}}}, the aggregate behind {{{}freqItems{}}}, keeps its 
counters in a {{mutable.Map[Any, Long]}} keyed by the column value. A 
{{BinaryType}} value is an {{{}Array[Byte]{}}}, and Java arrays use referential 
equality and identity hash codes. Each row's value is a separate array, so 
every row looks like a new value. Its count never exceeds 1, and once the map 
is full the eviction step keeps clearing it.

This dates back to the original implementation in SPARK-7242 (1.4.0).

Other types are not affected. {{{}UTF8String{}}}, {{BinaryView}} (used by 
{{GEOMETRY}} / {{{}GEOGRAPHY{}}}), {{{}UnsafeRow{}}}, {{UnsafeArrayData}} and 
{{GenericArrayData}} all compare by content, including binary nested in a 
struct or array.



--
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