SquirrelJoin: network-aware distributed join processing with lazy partitioning
File(s) p1250-rupprecht.pdf (531.44 KB)
Published version
Author(s)
Rupprecht, L
Culhane, WJ
Pietzuch, P
Type
Journal Article
Abstract
To execute distributed joins in parallel on compute clusters, systems
partition and exchange data records between workers. With large
datasets, workers spend a considerable amount of time transferring
data over the network. When compute clusters are shared among
multiple applications, workers must compete for network bandwidth
with other applications. These variances in the available network
bandwidth lead to
network skew
, which causes straggling workers
to prolong the join completion time.
We describe
SquirrelJoin
, a distributed join processing technique
that uses
lazy partitioning
to adapt to transient network skew in
clusters. Workers maintain in-memory
lazy partitions
to withhold a
subset of records, i.e. not sending them immediately to other work-
ers for processing. Lazy partitions are then assigned dynamically
to other workers based on network conditions: each worker takes
periodic throughput measurements to estimate its completion time,
and lazy partitions are allocated as to minimise the join completion
time. We implement SquirrelJoin as part of the Apache Flink dis-
tributed dataflow framework and show that, under transient network
contention in a shared compute cluster, SquirrelJoin speeds up join
completion times by up to 2.9× with only a small, fixed overhead.
partition and exchange data records between workers. With large
datasets, workers spend a considerable amount of time transferring
data over the network. When compute clusters are shared among
multiple applications, workers must compete for network bandwidth
with other applications. These variances in the available network
bandwidth lead to
network skew
, which causes straggling workers
to prolong the join completion time.
We describe
SquirrelJoin
, a distributed join processing technique
that uses
lazy partitioning
to adapt to transient network skew in
clusters. Workers maintain in-memory
lazy partitions
to withhold a
subset of records, i.e. not sending them immediately to other work-
ers for processing. Lazy partitions are then assigned dynamically
to other workers based on network conditions: each worker takes
periodic throughput measurements to estimate its completion time,
and lazy partitions are allocated as to minimise the join completion
time. We implement SquirrelJoin as part of the Apache Flink dis-
tributed dataflow framework and show that, under transient network
contention in a shared compute cluster, SquirrelJoin speeds up join
completion times by up to 2.9× with only a small, fixed overhead.
Date Issued
2017-08-01
Date Acceptance
2017-05-16
Citation
Proceedings of the VLDB Endowment, 2017, 10 (11), pp.1250-1261
ISSN
2150-8097
Publisher
VLDB Endowment
Start Page
1250
End Page
1261
Journal / Book Title
Proceedings of the VLDB Endowment
Volume
10
Issue
11
Copyright Statement
© 2017 VLDB Endowment. This work is licensed under the Creative Commons Attribution-NonCommercial-NoDerivatives 4.0 International License. To view a copy of this license, visit http://creativecommons.org/licenses/by-nc-nd/4.0/. For any use beyond those covered by this license, obtain permission by emailing info@vldb.org.
Sponsor
Engineering & Physical Science Research Council (EPSRC)
Grant Number
EP/K032968/1
Publication Status
Published
Coverage Spatial
Munich, Germany
