[ 
https://issues.apache.org/jira/browse/SPARK-59898?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

ASF GitHub Bot updated SPARK-59898:
-----------------------------------
    Labels: correctness pull-request-available  (was: correctness)

> 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: 1.4.0, 2.4.0, 3.5.0, 4.2.0
>            Reporter: David Mollitor
>            Priority: Minor
>              Labels: correctness, pull-request-available
>
> 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