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