Coflows for geo-distributed computer sites that communicate via wide area network
A coflow is mapped to a plurality of geo-distributed computer sites that can communicate via wide area network (WAN), where the mapping is subject to one or more location-dependent constraints. Multiple candidate data paths are identified for each of a plurality of source-destination pairs of the plurality of geo-distributed computer sites. A mathematical optimization is performed to find a set of paths from the candidate data paths based on total flow completion time and at least one additional objective of the coflow.
1 . A computer-implemented method, comprising:
accessing a dependency graph structure that describes geo-distributed job objectives and data sources of a coflow, inter-task data volumes of the coflow, and one or more location-dependent constraints of the coflow;
determining, for a stage of a plurality of successive computational stages, a plurality of geo-distributed computer sites based on the one or more location-dependent constraints;
mapping the coflow to the plurality of geo-distributed computer sites, wherein
the plurality of geo-distributed computer sites communicates via wide area network (WAN),
the mapping is based on the one or more location-dependent constraints, and
the mapping includes placing a plurality of tasks at the plurality of geo-distributed computer sites;
identifying multiple candidate data paths for each of a plurality of source-destination pairs of the plurality of geo-distributed computer sites; and
performing a first mathematical optimization to identify a set of paths from the multiple candidate data paths, wherein
the identifying of the set of paths is based on a total flow completion time of the coflow and at least one additional objective of the coflow, and
the at least one additional objective of the coflow is derived from the geo-distributed job objectives.
2 . The computer-implemented method of claim 1 , wherein the at least one additional objective is job-specific.
3 . The computer-implemented method of claim 1 , wherein the at least one additional objective includes at least one of WAN bandwidth utilization overhead, link usage cost, or egress copy overhead.
4 . The computer-implemented method of claim 1 , wherein the mapping comprises:
performing a second mathematical optimization of an objective function based on a proportion of tasks of the plurality of tasks, upload bandwidth, download bandwidth, and data volume for each of the the plurality of geo-distributed computer sites;
quantizing each proportion of the plurality of tasks; and
scaling each quantized proportion of the plurality of tasks to produce an actual number of tasks at each of the plurality of geo-distributed computer sites.
5 . The computer-implemented method of claim 4 , wherein the objective function of the second mathematical optimization is optimized in part with respect to copy time.
6 . The computer-implemented method of claim 1 , wherein the first mathematical optimization comprises:
generating an objective function based on the total flow completion time and the at least one additional objective; and
finding a minimum of the objective function to obtain the set of paths.
7 . The computer-implemented method of claim 6 , wherein:
the coflow includes a plurality of flows,
each flow of the plurality of flows is associated with one of the plurality of source-destination pairs, and
the multiple candidate data paths are assigned to each flow of the plurality of flows;
flow finish time of a given flow, of the plurality of flows, is a summation of path completion times for the assigned multiple candidate data paths of the given flow;
flow finish times of the plurality of flows are equal; and
the total flow completion time equals the flow finish times.
8 . The computer-implemented method of claim 7 , wherein a path completion time of a given path of the multiple candidate data paths is a function of a fraction of volume through the given path and bandwidth assigned to the given path.
9 . The computer-implemented method of claim 8 , wherein:
the bandwidth assigned to the given path is a function of total bandwidth used by the plurality of flows using a link; and
the total bandwidth is limited by residual capacity of the link.
10 . The computer-implemented method of claim 1 , further comprising sending information about the set of paths to a software-driven WAN controller.
11 . The computer-implemented method of claim 1 , further comprising reducing bandwidth of the coflow, wherein
a plurality of flows of the coflow are consumed at its destination at a predicted time based on the reducing of the bandwidth, and
excess bandwidth is created based on the reducing of the bandwidth.
12 . The computer-implemented method of claim 11 wherein the reducing of the bandwidth of the coflow comprises proportionately scaling down bandwidth of the plurality of flows of the coflow.
13 . The computer-implemented method of claim 11 , further comprising:
processing at least one additional coflow; and
reapportioning the excess bandwidth to the at least one additional coflow.
14 . The computer-implemented method of claim 1 , further comprising scheduling the coflow based on a shortest coflow first policy.
15 . A computing device comprising a memory having computer readable instructions, and one or more processors for executing the computer readable instructions to configure the computing device to perform operations comprising:
accessing a dependency graph structure that describes geo-distributed job objectives and data sources of a coflow, inter-task data volumes of the coflow, and one or more location-dependent constraints of the coflow;
determining, for a stage of a plurality of successive computational stages, a plurality of geo-distributed computer sites based on the one or more location-dependent constraints;
mapping the coflow to the plurality of geo-distributed computer sites, wherein
the plurality of geo-distributed computer sites communicates via wide area network (WAN),
the mapping is based on the to one or more location-dependent constraints, and
the mapping includes placing a plurality of tasks at the plurality of geo-distributed computer sites;
identifying multiple candidate data paths for each of a plurality of source-destination pairs of the plurality of geo-distributed computer sites; and
performing a mathematical optimization to identify a set of paths from the multiple candidate data paths, wherein
the identifying of the set of paths is based on a total flow completion time of the coflow and at least one additional objective of the coflow, and
the at least one additional objective of the coflow is derived from the geo-distributed job objectives.
16 . The computing device of claim 15 , wherein the execution of the computer readable instructions configures the computing device to perform the operations further comprising:
processing at least one additional coflow;
reducing bandwidth of the coflow to create excess bandwidth; and
reapportioning the excess bandwidth to the at least one additional coflow, wherein
the at least one additional coflow is processed by the computing device.
17 . A computer program product comprising one or more computer-readable memory devices encoded with data including instructions that, when executed, causes a processor set to perform a method, comprising:
accessing a dependency graph structure that describes geo-distributed job objectives and data sources of a coflow, inter-task data volumes of the coflow, and one or more location-dependent constraints of the coflow;
determining, for a stage of a plurality of successive computational stages, a plurality of geo-distributed computer sites based on the one or more location-dependent constraints:
mapping the coflow to the plurality of geo-distributed computer sites, wherein
the plurality of geo-distributed computer sites communicates via wide area network (WAN),
the mapping is based on one or more location-dependent constraints, and
the mapping includes placing a plurality of tasks at the plurality of geo-distributed computer sites;
identifying multiple candidate data paths for each of a plurality of source-destination pairs of the plurality of geo-distributed computer sites; and
performing a mathematical optimization to identify a set of paths from the multiple candidate data paths, wherein
the identifying of the set of paths is based on a total flow completion time of the coflow and at least one additional objective of the coflow, and
the at least one additional objective of the coflow is derived from the geo-distributed job objectives.