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