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)

Reply via email to