Fine-grained dynamic load balancing in spatial join by work stealing on distributed memory
Fine-grained dynamic load balancing in spatial join by work stealing on distributed memory
复制标题
通过分布式内存上的工作窃取实现空间连接中的细粒度动态负载平衡
DOI:
10.1145/3557915.3560936
复制
发表时间:
2022
期刊:
影响因子:
--
通讯作者:
Zhou, Hui
中科院分区:
文献类型:
--
作者:
Yang, Jie;Puri, Satish;Zhou, Hui
Spatial join is an important operation for combining spatial data. Parallelization is essential for improving spatial join performance. However, load imbalance due to data skew limits the scalability of parallel spatial join. There are many work sharing techniques to address this problem in a parallel environment. One of the techniques is to use data and space partitioning and then scheduling the partitions among threads/processes with the goal of minimizing workload differences across threads/processes. However, load imbalance still exists due to differences in join costs of different pairs of input geometries in the partitions.For the load imbalance problem, we have designed a work stealing spatial join system (WSSJ-DM) on a distributed memory environment. Work stealing is an approach for dynamic load balancing in which an idle processor steals computational tasks from other processors [5]. This is the first work that uses work stealing concept (instead of work sharing) to parallelize spatial join computation on a large compute cluster. We have evaluated the scalability of the system on shared and distributed memory. Our experimental evaluation shows that work stealing is an effective strategy. We compared WSSJ-DM with work sharing implementations of spatial join on a high performance computing environment using partitioned and un-partitioned datasets. Static and dynamic load balancing approaches were used for comparison. We study the effect of memory affinity in work stealing operations involved in spatial join on a multi-core processor.WSSJ-DM performed spatial join usingST_IntersectiononLakes(8.4M polygons) andParks(10M polygons) in 30 seconds using 35 compute nodes on a cluster (1260 CPU cores). A work sharing Master-Worker implementation took 160 seconds in contrast.
登录
查看更多内容
DOI:
10.1007/3-540-60159-7_13
发表时间:
1995-08
期刊:
--
影响因子:
--
作者:
S. Shekhar;S. Ravada;Vipin Kumar;Douglas Chubb;Greg Turner
通讯作者:
S. Shekhar;S. Ravada;Vipin Kumar;Douglas Chubb;Greg Turner
DOI:
--
发表时间:
2016
期刊:
International Symposium on Computing and Networking - Across Practical Development and Theoretical Research
影响因子:
--
作者:
Kouichi Araki;Taiki Shimbo
通讯作者:
Taiki Shimbo
DOI:
10.1109/hipc.2019.00027
发表时间:
2019-12
期刊:
2019 IEEE 26th International Conference on High Performance Computing, Data, and Analytics (HiPC)
影响因子:
--
作者:
Yiming Liu;Jie Yang;S. Puri
通讯作者:
Yiming Liu;Jie Yang;S. Puri
DOI:
10.1109/ccgrid54584.2022.00064
发表时间:
2022
期刊:
2022 22nd IEEE International Symposium on Cluster, Cloud and Internet Computing (CCGrid)
影响因子:
--
作者:
Anmol Paudel;S. Puri
通讯作者:
S. Puri
DOI:
--
发表时间:
2020
期刊:
34th IEEE International Parallel & Distributed Processing Symposium
影响因子:
--
作者:
Yang, Jie;Puri, Satish
通讯作者:
Puri, Satish