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

David Mollitor updated SPARK-59898:
-----------------------------------
    Description: 
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.

  was:
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.


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