James Kyle created SEDONA-66:
--------------------------------
Summary: Coalesce after spatial join results in dropped rows.
Key: SEDONA-66
URL: https://issues.apache.org/jira/browse/SEDONA-66
Project: Apache Sedona
Issue Type: Bug
Affects Versions: 1.0.0
Reporter: James Kyle
Attachments: image-2021-09-30-11-41-43-840.png,
image-2021-09-30-11-43-18-164.png
The code snippet we were using internal that reliably reproduced this error
looks something like
`leftSpatialRDD.spatialJoinFlat(rightSpatialRDD,
ALLOW_OVERLAP).coalesce(NUM_POST_JOIN_PARTITIONS)`
repartition does _not_ drop rows and serves as a workaround..
Below is the summary of our internal conversation that potentially root caused
the problem. I've bolded the final sentence which is the summary of the
potential root cause.
h2. ISSUE 1
{code:java}
val spatialRDD = makeSpatial(ds)
setSpatialPartitioner(spatialRDD)
val joined = spatialRDD.spatialJoinFlat(spatialRDD,
ALLOW_OVERLAP).coalesce(NUM_POST_JOIN_PARTITIONS)
val joinedNew = spatialRDD.spatialJoinFlat(spatialRDD, ALLOW_OVERLAP)
{code}
zipPartitions, and within each partition, it does something like indexed/for
loop element checking
!image-2021-09-30-11-41-43-840.png|width=819,height=454!
then looking at the DAG of executed code, there are something interesting ,
first is SparkSQL inject a Scan operator between Sedona's RDD operation and the
following DataFrame operations (.coalesce), and zipPartitions and coalesce is
put in the same stage, the colocated zipPartitions and coalesce means that,
even the joined RDD from Sedona is supposed to have N partitions, it will still
converge to X partitions where X is specified by coalesce,
!image-2021-09-30-11-43-18-164.png|width=575,height=972!
assumption of self join, and the implementation of coalesce, selfjoin assumes
that the partitions will be 1:1 mapped, however, coalesce is a wrapper of
upstream iterators which doesn't guarantee the 1:1 mapped, *it is highly
possible that the highlight* *match in the screenshot returns a lot of false*
--
This message was sent by Atlassian Jira
(v8.3.4#803005)