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]