Efficient Filter and Refinement for the Hadoop-Based Spatial Join

Jiamin Lu, Ralf Hartmut Güting, Jun Feng · 2016

The spatial join is usually processed in two steps: filter and refinement. The first creates candidate pairs based on the objects' abstractions, like object pairs having overlapping minimum bounding rectangles (MBRs). The second step then further checks each candidate pair whether the objects fit the join condition with their actual shapes. This two-step strategy prevents the disk I/O and CPU overhead spent on retrieving and comparing the non-candidates, hence it is also adopted at the large scale processing, based on the Hadoop platform. However, in order to cope with different data distributions, the Hadoop-based spatial join often processes both above steps in the Reduce stage, causing considerable data migration overhead as all object pairs need to be shuffled over the cluster network. In this paper we propose two novel shuffling approaches for the distributed filter and refinement operations. They are able to retrieve and even transfer only the candidates' actual shapes data in the Reduce tasks, decreasing the migration overhead to the minimum. In our evaluations with a six-node cluster, both approaches outstandingly improve the process efficiency regardless of the join selectivity.

Read the paper · More papers on PaperTik