Prosecution Insights
Last updated: October 02, 2026
Application No. 18/654,953

SYSTEMS AND METHODS OF RESOURCE CONFIGURATION OPTIMIZATION FOR MACHINE LEARNING WORKLOADS

Non-Final OA §102§103§112
Filed
May 03, 2024
Priority
Mar 11, 2021 — continuation of 12/001,511
Examiner
LU, HWEI-MIN
Art Unit
Tech Center
Assignee
Hewlett Packard Enterprise Development L.P.
OA Round
1 (Non-Final)
63%
Grant Probability
Moderate
1-2
OA Rounds
6m
Est. Remaining
99%
With Interview

Examiner Intelligence

Grants 63% of resolved cases
63%
Career Allowance Rate
152 granted / 240 resolved
+3.3% vs TC avg
Strong +40% interview lift
Without
With
+40.2%
Interview Lift
resolved cases with interview
Typical timeline
2y 11m
Avg Prosecution
26 currently pending
Career history
264
Total Applications
across all art units

Statute-Specific Performance

§101
9.6%
-30.4% vs TC avg
§103
50.4%
+10.4% vs TC avg
§102
11.0%
-29.0% vs TC avg
§112
28.9%
-11.1% vs TC avg
Black line = Tech Center average estimate • Based on career data from 240 resolved cases

Office Action

§102 §103 §112
DETAILED ACTION Notice of Pre-AIA or AIA Status The present application, filed on or after March 16, 2013, is being examined under the first inventor to file provisions of the AIA . This office action is in responsive to communication(s): original application filed on 05/03/2024, said application claims a priority filing date of 03/11/2021. Claims pending. Claims 21 and 33 are independent. Drawings The drawings are objected to as failing to comply with 37 CFR 1.84(p)(4) because (1) reference character “406” has been used to designate both "DT-CONFIGURATION GENERATOR" in FIG. 4 and "DT-configuration search space" in ¶¶ [0043]-[0044]; (2) reference character “700” has been used to designate both "computing component" in ¶ [0062] with FIG. 7 and "selection of an algorithm for hyperparameter searching/tuning" in ¶ [0069]; (3) reference character “702” has been used to designate both "hardware processor" in ¶¶ [0062]-[0063] and [0065]-[0068] with FIG. 7 and "model hyperparameters" in ¶ [0069]; (4) reference character “704” has been used to designate both "machine-readable storage medium" in ¶¶ [0062]-[0064] with FIG. 1 and "; (5) reference character “706” has been used to designate both "determine a plurality of computing resource configurations used to perform machine learning model training jobs" or "resource configurations can be specified for a TensorFlow training job" in ¶¶ [0065] an [0070] with FIG. 7 and "resource configurations" in ¶¶ [0070]-[0071]; (6) reference character “708” has been used to designate both " and "optimize resource configurations to achieve the quickest completion (training) time" in ¶ [0070]; (6) reference character “710” has been used to designate both "determine whether a stopping criterion is satisfied" in ¶ [0066] with FIG. 7 and "Bayesian optimization techniques may be applied" in ¶ [0071]; and (7) reference character “712” has been used to designate both "select one of the plurality of computing resource configurations" in ¶ [0068] with FIG. 7 and "all potential resource configurations" in [0071]. Corrected drawing sheets in compliance with 37 CFR 1.121(d) are required in reply to the Office action to avoid abandonment of the application. Any amended replacement drawing sheet should include all of the figures appearing on the immediate prior version of the sheet, even if only one figure is being amended. Each drawing sheet submitted after the filing date of an application must be labeled in the top margin as either “Replacement Sheet” or “New Sheet” pursuant to 37 CFR 1.121(d). If the changes are not accepted by the examiner, the applicant will be notified and informed of any required corrective action in the next Office action. The objection to the drawings will not be held in abeyance. The drawings are objected to as failing to comply with 37 CFR 1.84(p)(5) because they do not include the following reference sign(s) mentioned in the description: (1) . Corrected drawing sheets in compliance with 37 CFR 1.121(d) are required in reply to the Office action to avoid abandonment of the application. Any amended replacement drawing sheet should include all of the figures appearing on the immediate prior version of the sheet, even if only one figure is being amended. Each drawing sheet submitted after the filing date of an application must be labeled in the top margin as either “Replacement Sheet” or “New Sheet” pursuant to 37 CFR 1.121(d). If the changes are not accepted by the examiner, the applicant will be notified and informed of any required corrective action in the next Office action. The objection to the drawings will not be held in abeyance. Specification The disclosure is objected to because of the following informalities: in . Appropriate correction is required. The use of the term "Bluetooth" in ¶ [0039], which is a trade name or a mark used in commerce, has been noted in this application. The term should be accompanied by the generic terminology; furthermore the term should be capitalized wherever it appears or, where appropriate, include a proper symbol indicating use in commerce such as ™, SM , or ® following the term. Although the use of trade names and marks used in commerce (i.e., trademarks, service marks, certification marks, and collective marks) are permissible in patent applications, the proprietary nature of the marks should be respected and every effort made to prevent their use in any manner which might adversely affect their validity as commercial marks. Claim Objections Claims 21, 28, and 33 are objected to because of the following informalities: in Claim 21, line 5, "… evaluating and select one or more optimal DT configurations …" appears to be "… evaluating and selecting one or more optimal DT configurations …"; in Claim 28, lines 2-4, "… re-allocating at least one of the parameter servers and the worker nodes in real-time … when at least one of the parameter servers and the worker nodes are operating …" appears to be "… re-allocating at least one of the parameter servers and the worker nodes in real-time … when the at least one of the parameter servers and the worker nodes are operating …" (see also 112 Rejections to Claim 28); in Claim 33, lines 3-4, "… executing training jobs to determine … wherein execution of the training jobs occurs in accordance with …" appears to be "… executing training jobs to determine … wherein the execution of the training jobs occurs in accordance with …". Appropriate correction is required. Claim Rejections - 35 USC § 112 The following is a quotation of 35 U.S.C. 112(b): (b) CONCLUSION.—The specification shall conclude with one or more claims particularly pointing out and distinctly claiming the subject matter which the inventor or a joint inventor regards as the invention. The following is a quotation of 35 U.S.C. 112 (pre-AIA ), second paragraph: The specification shall conclude with one or more claims particularly pointing out and distinctly claiming the subject matter which the applicant regards as his invention. Claims 25-40 are rejected under 35 U.S.C. 112(b) or 35 U.S.C. 112 (pre-AIA ), second paragraph, as being indefinite for failing to particularly point out and distinctly claim the subject matter which the inventor or a joint inventor (or for applications subject to pre-AIA 35 U.S.C. 112, the applicant), regards as the invention. Claim 25 recites the limitation "... selecting one of " in lines 2-3, which rendering the claim indefinite because " . Claims 26-32 are rejected for fully incorporating the deficiency of their respective base claims. Claim 28 recites the limitation "... adaptively re-allocating at least one of " in l, which rendering the claim indefinite because "… defining a DT search space of tuples representing . Claims 29-30 are rejected for fully incorporating the deficiency of their respective base claims. Claim 33 recites the limitation "... selecting the model hyperparameter values, the use of which in the execution of the training jobs results in ..." in in lines 9-10. There is insufficient antecedent basis for the limitation "the use of which in the execution of the training jobs" in the claim. Clarification is required. Claims 34-40 are rejected for fully incorporating the deficiency of their respective base claims. Claim 38 recites the limitation "... evaluating the monitored resource usage metrics against ..." in lines 1-2. There is insufficient antecedent basis for the limitation "the monitored resource usage metrics" in the claim. Since "monitoring resource usage metrics" is recited in Claim 37 and not in Claim 36, for examination purpose, consider Claim 38 as depending on Claim 37 instead of Claim 36. Claim 40 is rejected for fully incorporating the deficiency of their respective base claims. Claim 40 recites the limitation "… after the re-allocation of the one or more resources" in lines 2-3. There is insufficient antecedent basis for this limitation in the claim. Since "re-allocating one or more resources" is recited in Claim 39 and not in Claim 38, for examination purpose, consider Claim 40 as depending on Claim 39 instead of Claim 38. Claim Rejections - 35 USC § 102 In the event the determination of the status of the application as subject to AIA 35 U.S.C. 102 and 103 (or as subject to pre-AIA 35 U.S.C. 102 and 103) is incorrect, any correction of the statutory basis (i.e., changing from AIA to pre-AIA ) for the rejection will not be considered a new ground of rejection if the prior art relied upon, and the rationale supporting the rejection, would be the same under either status. The following is a quotation of the appropriate paragraphs of 35 U.S.C. 102 that form the basis for the rejections under this section made in this Office action: A person shall be entitled to a patent unless – (a)(1) the claimed invention was patented, described in a printed publication, or in public use, on sale, or otherwise available to the public before the effective filing date of the claimed invention. Claims 21-24 are rejected under 35 U.S.C. 102(a)(1) as being anticipated by Qiao et al. ("Pollux: Co-adaptive Cluster Scheduling for Goodput-Optimized Deep Learning", arXiv:2008.12260, Aug. 27, 2020, pp. 1-16), hereinafter Qiao. Independent Claim 21 Qiao discloses a computer-implemented method comprising: generating a plurality of distributed training (DT) configurations defining a DT search space of tuples representing a number of parameter servers and a number of worker nodes based on hyperparameter tuning training jobs specifying hyperparameter test values (Qiao, Section 1 with FIG. 1 of Pages 1-2: DL jobs are resource-intensive and long-running, demanding distributed execution using expensive hardware devices (e.g. GPUs or TPUs) in order to complete within reasonable amounts of time. However, existing schedulers require users submitting jobs to also specify training parameters (i.e., hyper-parameters) that, if set incorrectly, can greatly degrade job performance and resource efficiency. Of these training parameters (i.e., hyper-parameters), the batch size and learning rate of a DL job are strongly dependent on its allocation of resources, making them particularly difficult to decide in advance in shared-resource environments. Furthermore, an allocation of resources that can be efficiently utilized by a DL job not only depends on the structure of the model being trained, but also on the batch size and learning rate. This co-dependence between the resources, batch size, and learning rate creates a complex web of considerations a user must make in order to configure their job for efficient execution and resource utilization. Fundamentally, an efficiently configured DL job strikes a balance between two often opposing desires: (1) system throughput, the number of training examples processed per wall-clock time, and (2) statistical efficiency, the amount of progress made per training example processed. A larger batch size enables higher utilization of more resources (e.g. larger number of GPUs). However, when the batch size is increased, the learning rate must be re-tuned. Otherwise, statistical efficiency will decrease so that the total training time will not be any shorter, wasting the additionally allocated GPUs. Even with an optimally-tuned learning rate, increasing the batch size results in faster-decreasing statistical efficiency. For every distinct allocation of GPUs, there is potentially a different batch size that best balances increasing system throughput with decreasing statistical efficiency, as illustrated in Fig. 1b. Thus, the best choice of batch size and learning rate depends on the resource allocation, which in turn depends on competition from other jobs sharing the cluster. In turn, the best choice of resource allocation depends on the chosen batch size. The batch size, learning rate, and therefore the best resource allocation, all depend on the current training progress of the job. Therefore, we argue that the choice of resource allocations, batch sizes, and learning rates are best made collectively and dynamically by a knowledgeable cluster scheduler. This paper presents Pollux, a hybrid resource scheduler that co-adaptively allocates resources while tuning the batch size and learning rate for every DL job in a shared cluster. We design and implement a scheduling architecture that locally tunes the batch size and learning rate for each DL job, and globally optimizes cluster-wide resource allocations using a genetic algorithm. Both components actively cooperate with each other, and operate based on a common goal of goodput maximization. In cloud environments, Pollux can provision the right amount of resources at the right time, based on job training progress, to maximize statistical efficiency across the entire lifetime of a large DL job. Section 2 in Pages 2-4: Data-parallelism is a popular method of distributed execution for DL training. The model parameters w(t) are replicated across a set of distributed GPUs 1, …, K, and each mini-batch M(t) is divided into equal-sized partitions per node, M 1 ( t ) , … , M K ( t ) . Each GPU k computes a local gradient estimate g ^ k ( t ) using its own partition, as shown in Eqn. 4. These local gradient estimates are then averaged across all replicas to obtain the desired g ^ ( t ) , as defined by Eqn. 3. Finally, each node applies the same update using g ^ ( t ) to obtain the new model parameters w(t+1), as defined by Eqn. 2. The batch size used during training determines the upper limit on the system throughput. On the other hand, Tsync is typically dependent on the size of the gradients and model parameters, rather than the batch size or number of replicas. Due to Amdahl’s Law, no matter how many replicas are used, the run-time of each training iteration is lower bounded by Tsync. To overcome this scalability limitation, it is desirable to increase the batch size, which allows a larger proportion of time to be spent computing the local gradient estimates as opposed to synchronizing gradients and parameters over the network. Thus, using a larger batch size enables higher system throughput when scaling to more data-parallel replicas. The statistical efficiency of DL training can be defined as the amount of training progress made per unit of data processed. Parameters such as batch size and learning rate influence this quantity. For example, when the batch size is increased, the efficiency will decrease (by an amount that depends on how the learning rate is comparatively scaled). We express the statistical efficiency in terms of a quantity called the gradient noise scale. We then describe adaptive scaling rules that have been developed to set the learning rate with respect to the batch size, to achieve high statistical efficiency. In Pollux, the gradient noise scale will be used to define a measure of statistical efficiency for a given choice of batch size, which will allow us to evaluate and optimize the goodput. Relatedly, adaptive scaling rules will be used to dynamically choose an appropriate learning rate with respect to the current batch size during this optimization. A knowledgeable scheduler like Pollux can adaptively scale training parameters for the right DL jobs at the right times, to better accommodate for their model-dependent and time-varying levels of efficiency. If the batch size is increased, the learning rate must also be scaled up to maintain a high statistical efficiency. Correctly scaling the learning rate can improve statistical efficiency when training with large batch sizes, and result in orders-of-magnitude improvements to DL scalability and job completion time. Instead of scaling the learning rate by a constant factor for each batch size, AdaScale scales the learning rate adaptively based on the gradient noise scale φ t . AdaScale scales the learning rate by the factor rt as shown in Eq. 5. For resource scheduling, the most important characteristic of AdaScale is its predictability. One iteration of AdaScale training with batch size m is approximately equivalent to rt iterations of training with the original batch size m0. This property can be leveraged to measure statistical efficiency during training, and to predict its value before scaling to different batch sizes. Section 3 with FIGS. 2-3 in Pages 4-5: Definition 3.1. (Goodput) The goodput of a DL job at iteration t is the product between its system throughput and its statistical efficiency at iteration t. GOODPUTt (a, m)=THROUGHPUT(a, m)´EFFICIENCYt (m). a ∈ R N is an allocation vector, where a n is the number of GPUs allocated from node n, and m is the batch size. An initial batch size m0 and learning rate η0 are selected by the user when submitting their job. As the job runs, Pollux profiles its execution to learn and refine predictive models for both THROUGHPUT and EFFICIENCY. Using these predictive models, Pollux periodically re-tunes a and m for each job, according to cluster-wide resource availability and performance. The learning rate η is re-tuned using AdaScale. EFFICIENCYt (m) in Eqn. 7 can be used by Pollux to predict the statistical efficiency at a different batch size without needing to train using that batch size ahead of time. To model and predict the system throughput for data-parallel DL, we aim to predict the time spent per training iteration, Titer, given an allocation vector a and batch size m, and then calculate the throughput as Eqn. 8. Fig. 3 shows an example of our THROUGHPUT function fit to measured throughput values for a range of resource allocations and batch sizes. Section 4 with FIGS. 4-5 in Pages 5-8: At a job-level granularity, Pollux dynamically tunes the batch size and learning rate for best utilization of the allocated resources. Pollux’s design consists of a PolluxAgent which measures the gradient noise scale and system throughput for each job, and tunes its batch size and learning rate for efficient utilization of its current allocated resources. An instance of PolluxAgent is started with each training job. During training, it continually measures the job’s gradient noise scale and system throughput, and reports them to PolluxSched at a fixed interval. It also uses this information to determine the most efficient batch size for its job given its current resource allocations, and adapts its job’s learning rate to this batch size using AdaScale. PolluxAgent measures the time taken per iteration, Titer, and records the triple (a, m, Titer) for all combinations of resource allocations a and batch size m encountered during its lifetime. PolluxAgent determines the most efficient batch size according to Eqn. 13, where a is the job’s current resource allocation. To perform this maximization efficiently, we observe that GOODPUT(a, m) is a unimodal function of m, and use golden-section search to find its maximum; once a new batch size is found, the job will use it for its subsequent training iterations, using AdaScale to adapt its learning rate appropriately. As the job’s statistical efficiency changes over time, PolluxAgent will periodically re-evaluate the most efficient batch size and learning rate. A in Eqn. 14 is an allocation matrix with each row Aj being the placement vector for a job j, thus Ajn is the number of GPUs on node n allocated to job j. Our genetic algorithm operates on a population of distinct allocation matrices (see Fig. 5). During each generation of the algorithm, existing allocation matrices are first randomly mutated, then crossed over to produce offspring allocation matrices, and finally modified to satisfy node resource constraints. After reporting to PolluxSched, PolluxAgent updates the job’s batch size, by optimizing its now up-to-date goodput function (Eqn. 6) with its currently allocated resources. Section 5 in Pages 8-12: We configured PolluxSched to use a 60s scheduling interval. During each scheduling interval, we run the genetic algorithm for 100 generations using a population of 100 allocation matrices. We set GPUTIME_THRES to 4 GPU-hours with λ = 0:5, and RESTART_PENALTY to 0.25. We configured PolluxAgent to report up-to-date system throughput parameters and gradient statistics every 30s. The primary workload used in our evaluation consists of 160 job submissions randomly sampled from an 8-hour period which includes the peak rate of job submissions during an average 24-hour day. We conduct experiments using a cluster consisting of 16 nodes and 64 GPUs. Each node is an AWS EC2 g4dn.12xlarge instance with 4 Tesla T4 GPUs, 48 vCPUs, 192GB memory, and a 900GB SSD. All instances are launched within the same placement group); evaluating and select one or more optimal DT configurations from the plurality of DT configurations; running one or more machine learning (ML) jobs on the selected one or more optimal DT configurations to determine an optimal resource allocation comprising preferred numbers of parameter servers and worker nodes for the one or more ML jobs; and determining a best DT configuration based on the determined optimal resource allocation (Qiao, Abstract of Page 1: Pollux improves scheduling performance in deep learning (DL) clusters by adaptively co-optimizing inter-dependent factors both at the per-job level and at the cluster-wide level. Pollux simultaneously considers both aspects. By observing each job during training, Pollux models how their goodput (system throughput combined with statistical efficiency) would change by adding or removing resources. Leveraging these models, Pollux dynamically (re-)assigns resources to maximize cluster-wide goodput, while continually optimizing each DL job to better utilize those resources. Section 1 of Pages 1-2: This paper presents Pollux, a hybrid resource scheduler that co-adaptively allocates resources while tuning the batch size and learning rate for every DL job in a shared cluster. We provide a formulation of goodput for DL jobs, which is a performance metric that takes into account both system throughput and statistical efficiency. A model of goodput can be learned by observing the throughput and statistical behavior during training, and used for predicting the performance of DL jobs given different resource allocations and batch sizes. We design and implement a scheduling architecture that locally tunes the batch size and learning rate for each DL job, and globally optimizes cluster-wide resource allocations using a genetic algorithm. Both components actively cooperate with each other, and operate based on a common goal of goodput maximization. Pollux not only adapts to the changes in statistical efficiency over time, but can leverage it to reduce the cost of training large models. In cloud environments, Pollux can provision the right amount of resources at the right time, based on job training progress, to maximize statistical efficiency across the entire lifetime of a large DL job. Section 2 in Pages 2-4: The system throughput of DL training can be defined as the number of training samples processed per unit of wall-clock time. When a DL job is distributed across several nodes, its system throughput is determined by several factors, including (1) the allocation and placement of resources assigned to the job, (2) the method of distributed execution and synchronization, and (3) the batch size used by the SGD algorithm. The run-time of each training iteration is determined by two main components. First, the time spent computing the local gradient estimates g ^ k ( t ) , which we denote by Tgrad. Second, the time spent averaging the local gradients and synchronizing the model parameters across all job replicas, which we denote by Tsync. Tsync can be influenced by the placement of replicas, and is typically smaller when the replicas are co-located within the same physical node or rack, rather than spread across different nodes or racks. Section 3 in Pages 4-5: Goodput predictions are leveraged by Pollux to jointly optimize cluster-wide resource allocations and batch sizes. To model and predict the system throughput for data-parallel DL, we aim to predict the time spent per training iteration, Titer, given an allocation vector a and batch size m, and then calculate the throughput as Eqn. 8. We start by separately modeling Tgrad, the time in each iteration spent computing local gradient estimates, and Tsync, the time in each iteration spent averaging gradient estimates and synchronizing model parameters across all GPUs. We model Tgrad as Tgrad ( a , m)=αgrad+βgrad·m/K, where m is the overall batch size, K = ∑ n a n is the number of allocated GPUs, and αgrad and βgrad are learnable parameters. We model Tsync as Eqn. 10, where N is the number of physical nodes occupied by at least one replica. α s y n c l o c a l and β s y n c l o c a l are the constant and retrogression parameters for when all processes are co-located onto the same node. α s y n c n o d e and β s y n c n o d e are the analogous parameters for when at least two process are located on different nodes. Modern DL frameworks can partially overlap Tgrad and Tsync by overlapping gradient computation with network communication. The degree of this overlap depends on structures in the specific DL model being trained, like the ordering and sizes of its layers. To capture the overlap between Tgrad and Tsync, we model Titer as Eqn. 11, where γ≥1 is a learnable parameter. Eqn. 11 has the property that Titer=Tgrad+Tsync when γ =1, and smoothly transitions towards Titer =max(Tgrad, Tsync) as γ→∞. Section 4 with FIG. 4 in Pages 5-8: At the cluster-wide granularity, Pollux dynamically re-allocates resources, driven by the goodput of all jobs sharing the cluster. Pollux’s design consists of two primary components. First, a PolluxAgent runs together with each job, which measures the gradient noise scale and system throughput for that job. PolluxAgent periodically reports the goodput function of its job to the PolluxSched. Second, the PolluxSched periodically optimizes the resource allocations for all jobs in the cluster, taking into account the current statistical efficiency for each job. It uses the goodput function to predict a job’s training performance when allocated different resources. PolluxAgent and PolluxSched co-adapt to each other. While PolluxAgent adapts each training job to make efficient use of its allocated resources, PolluxSched dynamically re-allocates each job’s resources, taking into account the PolluxAgent’s ability to tune its job. Fig. 4 illustrates Pollux’s co-adaptive architecture (Fig. 4c), compared with existing schedulers which are either nonadaptive (Fig. 4a), or resource-adaptive without being involved with statistical efficiency (Fig. 4b). In Sec. 3.2, we defined the system throughput parameters θsys of a training job as the 7-tuple which are required to construct the THROUGHPUT function. Together with the gradient noise scale φ t and initial batch size m0, the triple (θsys, φ t , m0) specifies the GOODPUT function. While m0 is a constant configuration provided by the user, and φ t can be computed according to Sec. 3.1, θsys is estimated by fitting the THROUGHPUT function to observed throughput values collected about the job during training. We minimize the root mean squared logarithmic error (RMSLE) between Eqn. 11 and the collected data triples, using L-BFGS-B. We set constraints for each α and β parameter to be non-negative, and γ to be in the range [1, 10]. PolluxAgent then reports the updated values of θsys and φ t to PolluxSched. To ensure that Pollux finds efficient resource allocations through systematic exploration, we impose several priors which bias θsys towards the belief that throughput scales perfectly with more resources, until such resource configurations are explored. Each job starts with a single GPU, and is initially assumed to scale perfectly to more GPUs. PolluxSched is then encouraged to allocate more GPUs and/or nodes to the job, naturally as part of its resource optimization (Sec. 4.2), until the PolluxAgent can estimate θsys more accurately. Finally, to prevent a job from being immediately scaled out to arbitrarily many GPUs, we restrict the maximum number of GPUs which can be allocated to at most twice the maximum number of GPUs the job has been allocated in its lifetime. The PolluxSched periodically allocates (or re-allocates) resources for every job in the cluster. To determine a set of efficient cluster-wide resource allocations, it uses a genetic algorithm to maximize a fitness function which is defined as a weighted mean across speedups for each job in Eqn. 14. maximizing FITNESS causes PolluxSched to allocate more resources to jobs that achieve a high SPEEDUP when provided with many GPUs (i.e. jobs that scale well). Although PolluxSched can work well with all weights wj set to 1, we provide the ability to re-weight jobs as an additional tuning knob for cluster operators. We define the weight wj of job j as Eqn. 16. The weight of a job is 1 if its current total GPU-time is at most GPUTIME_THRES, and decays gradually thereafter. A constant population size of distinct allocation matrices is maintained after each generation by discarding the allocation matrices with the lowest objective values (according to Eqn. 14). After several generations of the genetic algorithm, the allocation matrix with the highest fitness score is applied to the jobs running in the cluster. Each time a job is re-allocated to a different set of GPUs, it will need to save a checkpoint of its current model parameters, and restart using its new allocation of GPUs. To prevent an excessive number of re-allocations, when PolluxSched evaluates the fitness function for a given allocation matrix, it applies a penalty for every job that needs to restart. When multiple distributed DL jobs share a single node, their network usage while synchronizing gradients and model parameters may interfere with each other, causing both jobs to slow down. PolluxSched mitigates this issue by disallowing different distributed jobs (each using GPUs across multiple nodes) from sharing the same node. In cloud computing environments, GPU nodes can be dynamically provisioned during periods of high demand to decrease job completion time, as well as released during periods of low demand to decrease cost of cluster resources. Since many distributed DL jobs tend to increase in statistical efficiency as training progresses, it may be more cost-effective to provision more cloud resources during the later iterations of a large training job, rather than earlier on. In order to decide when to request or release nodes in the cloud, we first define a measure of cluster resource utility for a given allocation matrix as Eqn. 17. A cluster operator defines a LOW_UTIL_THRES and a HIGH_UTIL_THRES. PolluxSched will request additional nodes (when UTILITY(A) is high) or release existing nodes (when UTILITY(A) is low) until UTILITY(A) falls within this range, where A is the current allocations applied in the cluster. The cluster operator also defines a MIN_NODES and MAX_NODES, which limit the size of the cluster. PolluxSched performs a binary search for the desired number of nodes, with the assumption that UTILITY decreases with increasing numbers of nodes. At each step of the binary search, PolluxSched runs its genetic algorithm to evaluate the UTILITY of the cluster size being queried. In the end, the cluster size with a UTILITY value closest to (LOW_UTIL_THRES+HIGH_UTIL_THRES)=2 is selected. PolluxSched is implemented as a service in Kubernetes [2]. At a fixed time interval, PolluxSched runs a fixed number of generations of its genetic algorithm. It then applies the best allocation matrix to the cluster, by creating and terminating Kubernetes Pods which run the job replicas. Although only the allocation matrix with the highest fitness score is applied to the cluster, the entire population is saved and used to bootstrap the genetic algorithm in the next scheduling interval. Section 5 with Table 1 in Pages 8-12: the primary value of Pollux is automatically choosing the right resource allocations, and adapting the batch size and learning rate of each job to best utilize those resources. We built a discrete-time cluster simulator to study the effects of the techniques proposed in Pollux. Our simulator reproduces the system throughput of the jobs in our workload, using different batch sizes, numbers of GPUs, and different placements of GPUs across nodes. It also reproduces the gradient noise scale across the lifetime of each job, enabling statistical efficiency to be calculated using the simulator. We use a cluster simulator in order to evaluate a broader set of workloads and settings. Our simulator is constructed by measuring the performance and gradient statistics of each model in Table 1, under many different resource and batch size configurations, and re-playing them for each simulated job. This way, we are able to simulate both the system throughput and statistical efficiency of the jobs in our workload. We measured the time per training iteration for all allocations of GPUs in a 4-node cluster with 4 GPUs each, removing symmetric placements. For each allocation, we used a range of batch sizes, spaced geometrically by factors of ≈ 2 , up to the largest batch size which fit into the total GPU memory. We did the same for all valid combinations of 4 to 16 nodes and 1 to 4 GPUs per node, where the same number of GPUs is allocated from each node. The data we collected about each job enables our simulator to reproduce several system effects, including the performance impact of different GPU placements. In cloud environments, computing resources can be obtained and released as required, and users pay for the duration they hold onto those resources. Since the statistical efficiency of DL training typically increases over time, larger batch sizes can be utilized more effectively later on in training. Pollux measures the statistical efficiency during training, and is able to provision more resources to accelerate large-batch training when it is the most impactful, and fewer resources to save cost while statistical efficiency is low. On the other hand, Pollux knows that the statistical efficiency near the beginning of the job is low for large batch sizes, and gradually increases the number of nodes as the effectiveness of larger batch sizes improves over time). Claim 22 Qiao discloses all the elements as stated in claim 21 and further discloses wherein each of the hyperparameter tuning jobs include an ML model to be trained (Qiao, Section 1 with FIG. 1 of Pages 1-2: DL jobs are resource-intensive and long-running, demanding distributed execution using expensive hardware devices (e.g. GPUs or TPUs) in order to complete within reasonable amounts of time. However, existing schedulers require users submitting jobs to also specify training parameters (i.e., hyper-parameters) that, if set incorrectly, can greatly degrade job performance and resource efficiency. Of these training parameters (i.e., hyper-parameters), the batch size and learning rate of a DL job are strongly dependent on its allocation of resources, making them particularly difficult to decide in advance in shared-resource environments. Furthermore, an allocation of resources that can be efficiently utilized by a DL job not only depends on the structure of the model being trained, but also on the batch size and learning rate. This co-dependence between the resources, batch size, and learning rate creates a complex web of considerations a user must make in order to configure their job for efficient execution and resource utilization. Fundamentally, an efficiently configured DL job strikes a balance between two often opposing desires: (1) system throughput, the number of training examples processed per wall-clock time, and (2) statistical efficiency, the amount of progress made per training example processed. A larger batch size enables higher utilization of more resources (e.g. larger number of GPUs). However, when the batch size is increased, the learning rate must be re-tuned. Otherwise, statistical efficiency will decrease so that the total training time will not be any shorter, wasting the additionally allocated GPUs. Even with an optimally-tuned learning rate, increasing the batch size results in faster-decreasing statistical efficiency. For every distinct allocation of GPUs, there is potentially a different batch size that best balances increasing system throughput with decreasing statistical efficiency, as illustrated in Fig. 1b. Thus, the best choice of batch size and learning rate depends on the resource allocation, which in turn depends on competition from other jobs sharing the cluster. In turn, the best choice of resource allocation depends on the chosen batch size. The batch size, learning rate, and therefore the best resource allocation, all depend on the current training progress of the job. Therefore, we argue that the choice of resource allocations, batch sizes, and learning rates are best made collectively and dynamically by a knowledgeable cluster scheduler. This paper presents Pollux, a hybrid resource scheduler that co-adaptively allocates resources while tuning the batch size and learning rate for every DL job in a shared cluster. We design and implement a scheduling architecture that locally tunes the batch size and learning rate for each DL job, and globally optimizes cluster-wide resource allocations using a genetic algorithm. Both components actively cooperate with each other, and operate based on a common goal of goodput maximization. In cloud environments, Pollux can provision the right amount of resources at the right time, based on job training progress, to maximize statistical efficiency across the entire lifetime of a large DL job. Section 3 in Page 4: As the job runs, Pollux profiles its execution to learn and refine predictive models for both THROUGHPUT and EFFICIENCY. Using these predictive models, Pollux periodically re-tunes a and m for each job, according to cluster-wide resource availability and performance. The learning rate η is re-tuned using AdaScale. Section 4 with FIG. 4 in Pages 5-6: At a job-level granularity, Pollux dynamically tunes the batch size and learning rate for best utilization of the allocated resources. A PolluxAgent runs together with each job. It measures the gradient noise scale and system throughput for that job, and tunes its batch size and learning rate for efficient utilization of its current allocated resources. PolluxAgent and PolluxSched co-adapt to each other. While PolluxAgent adapts each training job to make efficient use of its allocated resources, PolluxSched dynamically re-allocates each job’s resources, taking into account the PolluxAgent’s ability to tune its job. Section 4.1 in Pages 6: With fully specified the DL job’s GOODPUT function at its current training progress, PolluxAgent determines the most efficient batch size, m * = argm ax m ⁡ GOODPUT a ,   m , where a is the job’s current resource allocation. Once a new batch size is found, the job will use it for its subsequent training iterations, using AdaScale to adapt its learning rate appropriately. As the job’s statistical efficiency changes over time, PolluxAgent will periodically re-evaluate the most efficient batch size and learning rate). Claim 23 Qiao discloses all the elements as stated in claim 21 and further discloses wherein the hyperparameter tuning jobs are generated by a hyperparameter tuning algorithm (Qiao, Section 1 with FIG. 1 of Pages 1-2: DL jobs are resource-intensive and long-running, demanding distributed execution using expensive hardware devices (e.g. GPUs or TPUs) in order to complete within reasonable amounts of time. However, existing schedulers require users submitting jobs to also specify training parameters (i.e., hyper-parameters) that, if set incorrectly, can greatly degrade job performance and resource efficiency. Of these training parameters (i.e., hyper-parameters), the batch size and learning rate of a DL job are strongly dependent on its allocation of resources, making them particularly difficult to decide in advance in shared-resource environments. Furthermore, an allocation of resources that can be efficiently utilized by a DL job not only depends on the structure of the model being trained, but also on the batch size and learning rate. This co-dependence between the resources, batch size, and learning rate creates a complex web of considerations a user must make in order to configure their job for efficient execution and resource utilization. Fundamentally, an efficiently configured DL job strikes a balance between two often opposing desires: (1) system throughput, the number of training examples processed per wall-clock time, and (2) statistical efficiency, the amount of progress made per training example processed. A larger batch size enables higher utilization of more resources (e.g. larger number of GPUs). However, when the batch size is increased, the learning rate must be re-tuned. Otherwise, statistical efficiency will decrease so that the total training time will not be any shorter, wasting the additionally allocated GPUs. Even with an optimally-tuned learning rate, increasing the batch size results in faster-decreasing statistical efficiency. For every distinct allocation of GPUs, there is potentially a different batch size that best balances increasing system throughput with decreasing statistical efficiency, as illustrated in Fig. 1b. Thus, the best choice of batch size and learning rate depends on the resource allocation, which in turn depends on competition from other jobs sharing the cluster. In turn, the best choice of resource allocation depends on the chosen batch size. The batch size, learning rate, and therefore the best resource allocation, all depend on the current training progress of the job. Therefore, we argue that the choice of resource allocations, batch sizes, and learning rates are best made collectively and dynamically by a knowledgeable cluster scheduler. This paper presents Pollux, a hybrid resource scheduler that co-adaptively allocates resources while tuning the batch size and learning rate for every DL job in a shared cluster. We design and implement a scheduling architecture that locally tunes the batch size and learning rate for each DL job, and globally optimizes cluster-wide resource allocations using a genetic algorithm. Both components actively cooperate with each other, and operate based on a common goal of goodput maximization. In cloud environments, Pollux can provision the right amount of resources at the right time, based on job training progress, to maximize statistical efficiency across the entire lifetime of a large DL job. Section 3 in Page 4: As the job runs, Pollux profiles its execution to learn and refine predictive models for both THROUGHPUT and EFFICIENCY. Using these predictive models, Pollux periodically re-tunes a and m for each job, according to cluster-wide resource availability and performance. The learning rate η is re-tuned using AdaScale. Section 4 with FIG. 4 in Pages 5-6: At a job-level granularity, Pollux dynamically tunes the batch size and learning rate for best utilization of the allocated resources. A PolluxAgent runs together with each job. It measures the gradient noise scale and system throughput for that job, and tunes its batch size and learning rate for efficient utilization of its current allocated resources. PolluxAgent and PolluxSched co-adapt to each other. While PolluxAgent adapts each training job to make efficient use of its allocated resources, PolluxSched dynamically re-allocates each job’s resources, taking into account the PolluxAgent’s ability to tune its job. Section 4.1 in Pages 6: With fully specified the DL job’s GOODPUT function at its current training progress, PolluxAgent determines the most efficient batch size, m * = argm ax m ⁡ GOODPUT a ,   m , where a is the job’s current resource allocation. Once a new batch size is found, the job will use it for its subsequent training iterations, using AdaScale to adapt its learning rate appropriately. As the job’s statistical efficiency changes over time, PolluxAgent will periodically re-evaluate the most efficient batch size and learning rate. Section 6 in Page 12: Recent work on DL training algorithms have explored dynamically adapting batch sizes for better efficiency and parallelization. AdaBatch [10] increases the batch size at pre-determined iterations during training, while linearly scaling the learning rate. Smith et al. [44] suggest that instead of decaying the learning rate during training, the batch size should be increased instead. CABS [4] adaptively tunes the batch size and learning rate during training using similar gradient statistics as Pollux). Claim 24 Qiao discloses all the elements as stated in claim 21 and further discloses wherein the evaluation and selection of the one or more optimal DT configurations occurs in a DT loop comprising all possible or reasonable resource configurations comprising the tuples (Qiao, Section 2 in Pages 2-4: Training a deep learning model typically involves minimizing a loss function of the form in Eqn. 1, where w [Symbol font/0xCE] R d are the model parameters to be optimized, X is the training dataset, xi are the individual samples in the training data, and l is the loss evaluated at a single sample. The loss function can be minimized using stochastic gradient descent (SGD), which repeatedly applies the update until the loss converges to a stable value. When a DL job is distributed across several nodes, its system throughput is determined by several factors, including (1) the allocation and placement of resources assigned to the job, (2) the method of distributed execution and synchronization, and (3) the batch size used by the SGD algorithm. Data-parallelism is a popular method of distributed execution for DL training. The model parameters w(t) are replicated across a set of distributed GPUs 1, …, K, and each mini-batch M(t) is divided into equal-sized partitions per node, M 1 ( t ) , … , M K ( t ) . Each GPU k computes a local gradient estimate g ^ k ( t ) using its own partition, as shown in Eqn. 4. These local gradient estimates are then averaged across all replicas to obtain the desired g ^ ( t ) , as defined by Eqn. 3. Finally, each node applies the same update using g ^ ( t ) to obtain the new model parameters w(t+1), as defined by Eqn. 2. The run-time of each training iteration is determined by two main components. First, the time spent computing the local gradient estimates g ^ k ( t ) , which we denote by Tgrad. Second, the time spent averaging the local gradients and synchronizing the model parameters across all job replicas, which we denote by Tsync. Tsync can be influenced by the placement of replicas, and is typically smaller when the replicas are co-located within the same physical node or rack, rather than spread across different nodes or racks. Using a larger batch size enables higher system throughput when scaling to more data-parallel replicas. Section 3 in Pages 4-5: We present how goodput can be defined, measured, and predicted for DL jobs, taking into account both system throughput and statistical efficiency. Goodput predictions are leveraged by Pollux to jointly optimize cluster-wide resource allocations and batch sizes. Definition 3.1. (Goodput) The goodput of a DL job at iteration t is the product between its system throughput and its statistical efficiency at iteration t. GOODPUTt (a, m)=THROUGHPUT(a, m)´EFFICIENCYt (m). a ∈ R N is an allocation vector, where a n is the number of GPUs allocated from node n, and m is the batch size. An initial batch size m0 and learning rate η0 are selected by the user when submitting their job. As the job runs, Pollux profiles its execution to learn and refine predictive models for both THROUGHPUT and EFFICIENCY. Using these predictive models, Pollux periodically re-tunes a and m for each job, according to cluster-wide resource availability and performance. The learning rate η is re-tuned using AdaScale. Goodput and statistical efficiency are measured relative to the initial batch size m0 and learning rate η0, and Pollux only considers batch sizes which are at least the initial batch size, i.e.. m ≥ m0. The statistical efficiency of a DL job using batch size m ≥ m0 is the amount of progress made per training example using m, relative to using m0. This quantity can be framed in terms of the gradient noise scale φ t . We can write a concrete measure of statistical efficiency as EFFICIENCYt (m) = ( φ t + m 0 . )/( φ t + m ). Here, we can write the gradient noise scale as φ t = m 0 σ t 2 / μ t 2 , where σ t 2 = Var g ^ ( t ) is the variance and μ t 2 = E g ^ ( t ) 2 the squared norm of the gradient at iteration t using batch size m0. EFFICIENCYt (m) reflects the lifetime-dependent trends exhibited by the true statistical efficiency. The standard way of estimating σ t and μ t involves calculating the sample variance of the local gradient estimates g ^ k ( t ) at each iteration. This can be done efficiently when there are multiple data-parallel processes, by using the different values of g ^ k ( t ) already available on each process. To model and predict the system throughput for data-parallel DL, we aim to predict the time spent per training iteration, Titer, given an allocation vector a and batch size m, and then calculate the throughput as THROUGHPUT a , m = m / T i t e r a , m . We start by separately modeling Tgrad, the time in each iteration spent computing local gradient estimates, and Tsync, the time in each iteration spent averaging gradient estimates and synchronizing model parameters across all GPUs. We model Tgrad as Tgrad ( a , m)=αgrad+βgrad·m/K, where m is the overall batch size, K = ∑ n a n is the number of allocated GPUs, and αgrad and βgrad are learnable parameters. We model Tsync as Eqn. 10, where N is the number of physical nodes occupied by at least one replica. α s y n c l o c a l and β s y n c l o c a l are the constant and retrogression parameters for when all processes are co-located onto the same node. α s y n c n o d e and β s y n c n o d e are the analogous parameters for when at least two process are located on different nodes. Modern DL frameworks can partially overlap Tgrad and Tsync by overlapping gradient computation with network communication. The degree of this overlap depends on structures in the specific DL model being trained, like the ordering and sizes of its layers. To capture the overlap between Tgrad and Tsync, we model Titer as Eqn. 11, where γ≥1 is a learnable parameter. Eqn. 11 has the property that Titer=Tgrad+Tsync when γ =1, and smoothly transitions towards Titer =max(Tgrad, Tsync) as γ→∞. Section 4 with FIG. 4 in Pages 5-6: Pollux performs adaptation at two distinct granularities. First, at a job-level granularity, Pollux dynamically tunes the batch size and learning rate for best utilization of the allocated resources. Second, a the cluster-wide granularity, Pollux dynamically re-allocates resources, driven by the goodput of all jobs sharing the cluster. To achieve this co-adaptivity in a scalable way, Pollux’s design consists of two primary components. First, a PolluxAgent runs together with each job. It measures the gradient noise scale and system throughput for that job, and tunes its batch size and learning rate for efficient utilization of its current allocated resources. PolluxAgent periodically reports the goodput function of its job to the PolluxSched. Second, the PolluxSched periodically optimizes the resource allocations for all jobs in the cluster, taking into account the current statistical efficiency for each job. It uses the goodput function to predict a job’s training performance when allocated different resources. PolluxAgent and PolluxSched co-adapt to each other. While PolluxAgent adapts each training job to make efficient use of its allocated resources, PolluxSched dynamically re-allocates each job’s resources, taking into account the PolluxAgent’s ability to tune its job. Section 4.1 in Page 6: During training, it continually measures the job’s gradient noise scale and system throughput, and reports them to PolluxSched at a fixed interval. It also uses this information to determine the most efficient batch size for its job given its current resource allocations, and adapts its job’s learning rate to this batch size using AdaScale. PolluxAgent measures the time taken per iteration, Titer, and records the triple ( a ,   m , Titer) for all combinations of resource allocations a and batch size m encountered during its lifetime. Periodically, PolluxAgent fits the parameters θsys to all of the throughput data collected so far. Specifically, we minimize the root mean squared logarithmic error (RMSLE) between Eqn. 11 and the collected data triples, using L-BFGS-B. PolluxAgent then reports the updated values of θsys and φ t to PolluxSched. To ensure that Pollux finds efficient resource allocations through systematic exploration, we impose several priors which bias θsys towards the belief that throughput scales perfectly with more resources, until such resource configurations are explored. Each job starts with a single GPU, and is initially assumed to scale perfectly to more GPUs. PolluxSched is then encouraged to allocate more GPUs and/or nodes to the job, naturally as part of its resource optimization, until the PolluxAgent can estimate θsys more accurately. Finally, to prevent a job from being immediately scaled out to arbitrarily many GPUs, we restrict the maximum number of GPUs which can be allocated to at most twice the maximum number of GPUs the job has been allocated in its lifetime. With θsys, mgns, and m0, which fully specify the DL job’s GOODPUT function at its current training progress, PolluxAgent determines the most efficient batch size m * = argm ax m ⁡ GOODPUT a ,   m , where a is the job’s current resource allocation. Once a new batch size is found, the job will use it for its subsequent training iterations, using AdaScale to adapt its learning rate appropriately. As the job’s statistical efficiency changes over time, PolluxAgent will periodically re-evaluate the most efficient batch size and learning rate. Section 4.2 in Pages 6-7: The PolluxSched periodically allocates (or re-allocates) resources for every job in the cluster. To determine a set of efficient cluster-wide resource allocations, it uses a genetic algorithm to maximize a fitness function which is defined as a weighted mean across speedups for each job in Eqn. 14. A is an allocation matrix with each row Aj being the placement vector for a job j, thus Ajn is the number of GPUs on node n allocated to job j. Intuitively, maximizing FITNESS causes PolluxSched to allocate more resources to jobs that achieve a high SPEEDUP when provided with many GPUs (i.e. jobs that scale well). We define the speedup of each job as the factor of goodput improvement using the given placement vector over using a single process with the optimal batch size in Eqn. 15, where GOODPUTj is the goodput of job j at its current training iteration. Similar to the PolluxAgent, PolluxSched performs the maximizations in the numerator and denominator using golden-section search. Although PolluxSched can work well with all weights wj set to 1, we provide the ability to re-weight jobs as an additional tuning knob for cluster operators. We define the weight wj of job j as Eqn. 16. The weight of a job is 1 if its current total GPU-time is at most GPUTIME_THRES, and decays gradually thereafter. Section 4.2.1 with FIG. 5 in Page 7: Our genetic algorithm operates on a population of distinct allocation matrices (see Fig. 5). During each generation of the algorithm, existing allocation matrices are first randomly mutated, then crossed over to produce offspring allocation matrices, and finally modified to satisfy node resource constraints. A constant population size of distinct allocation matrices is maintained after each generation by discarding the allocation matrices with the lowest objective values (according to Eqn. 14). After several generations of the genetic algorithm, the allocation matrix with the highest fitness score is applied to the jobs running in the cluster. Each time a job is re-allocated to a different set of GPUs, it will need to save a checkpoint of its current model parameters, and restart using its new allocation of GPUs. To prevent an excessive number of re-allocations, when PolluxSched evaluates the fitness function for a given allocation matrix, it applies a penalty for every job that needs to restart. When multiple distributed DL jobs share a single node, their network usage while synchronizing gradients and model parameters may interfere with each other, causing both jobs to slow down. PolluxSched mitigates this issue by disallowing different distributed jobs (each using GPUs across multiple nodes) from sharing the same node. Section 4.2.2 in Pages 7-8: In cloud computing environments, GPU nodes can be dynamically provisioned during periods of high demand to decrease job completion time, as well as released during periods of low demand to decrease cost of cluster resources. Beyond adapting to the number and sizes of DL jobs, co-adaptive scheduling presents a unique opportunity for cluster auto-scaling. Since many distributed DL jobs tend to increase in statistical efficiency as training progresses, it may be more cost-effective to provision more cloud resources during the later iterations of a large training job, rather than earlier on. In order to decide when to request or release nodes in the cloud, we first define a measure of cluster resource utility for a given allocation matrix as Eqn. 17. A cluster operator defines a LOW_UTIL_THRES and a HIGH_UTIL_THRES. PolluxSched will request additional nodes (when UTILITY(A) is high) or release existing nodes (when UTILITY(A) is low) until UTILITY(A) falls within this range, where A is the current allocations applied in the cluster. The cluster operator also defines a MIN_NODES and MAX_NODES, which limit the size of the cluster. PolluxSched performs a binary search for the desired number of nodes, with the assumption that UTILITY decreases with increasing numbers of nodes. At each step of the binary search, PolluxSched runs its genetic algorithm to evaluate the UTILITY of the cluster size being queried. In the end, the cluster size with a UTILITY value closest to (LOW_UTIL_THRES+HIGH_UTIL_THRES)=2 is selected. Section 4.3 in Page 8: PolluxAgent inserts performance profiling code which measures the time taken for each iteration of training, as well as calculating the gradient noise scale. At a fixed time interval, PolluxAgent fits the system throughput model (Eqn. 11) to the profiled metrics collected so far, and reports the fitted system throughput parameters, along with the latest gradient statistics, to PolluxSched. After reporting to PolluxSched, PolluxAgent updates the job’s batch size, by optimizing its now up-to-date goodput function (Eqn. 6) with its currently allocated resources. PolluxSched is implemented as a service in Kubernetes. At a fixed time interval, PolluxSched runs a fixed number of generations of its genetic algorithm. It then applies the best allocation matrix to the cluster, by creating and terminating Kubernetes Pods which run the job replicas. Although allocation matrix with the highest fitness score is applied to the cluster, the entire population is saved and used to bootstrap the genetic algorithm in the next scheduling interval). Claims 33-37 are rejected under 35 U.S.C. 102(a)(1) as being anticipated by Rocha et al., ("PipeTune: Pipeline Parallelism of Hyper and System Parameters Tuning for Deep Learning Clusters", arXiv:2010.00501, Oct. 2, 2020, pp. 1-16), hereinafter Rocha. Independent Claim 33 Rocha discloses a computer-implemented method, comprising: generating model hyperparameter values for evaluating a machine learning (ML) model; executing training jobs to determine which of the model hyperparameter values result in optimal ML model performance, wherein execution of the training jobs occurs in accordance with a plurality of resource configurations specifying characteristics of a distributed training environment; selecting an optimal resource configuration from the plurality of resource configurations based on a shortest time to execute the training jobs; and selecting the model hyperparameter values, the use of which in the execution of the training jobs results in a desired level of quality of the ML model (Rocha, Abstract in Page 1: DNN learning jobs are common in today’s clusters due to the advances in AI driven services. The most critical phase of these jobs for model performance and learning cost is the tuning of hyperparameters. Existing approaches make use of techniques such as early stopping criteria to reduce the tuning impact on learning cost. However, these strategies do not consider the impact that certain hyperparameters and systems parameters have on training time. This paper presents PipeTune, a framework for DNN learning jobs that addresses the trade-offs between these two types of parameters. PipeTune takes advantage of the high parallelism and recurring characteristics of such jobs to minimize the learning cost via a pipelined simultaneous tuning of both hyper and system parameters. Section 1 with FIGS. 1-2 and Table 1 in Pages 1-3 with FIG. 3 in Page 4: Several public cloud providers offer native support to deploy, configure and run them, providing tools to automatically or semi-automatically drive the DNN processing pipeline. One important factor is the choice of the DNN hyperparameters (e.g., number of hidden layers, learning rate, dropout rate, momentum, batch size, weight-decay, epochs, pooling size, type of activation function, etc.). DNNs require careful tuning of the hyperparameters, and if done correctly, it can achieve impressive boosts in performance. However, misconfigurations can easily lead to wrong models and hence bad predictions. A naive approach to hyperparameter tuning is to perform a full exploration of the possible configuration variations. Such a tuning approach becomes quickly unpractical, costly and slow, as the number of variations grows exponentially. As a result of proper hyperparameters tuning, one should achieve fast convergence and high accuracy. We observe that some hyperparameters (e.g., number of epochs, batch size, dropout) can drastically reduce training time. Importantly, training a DNN by using different system resources (e.g., number of CPU cores, allocated memory, number of GPUs) lead to different results, as we also demonstrate later in Figure 3 for varying number of cores. The majority of the existing tuning solutions restrict themselves to the sole hyperparameter tuning using a variety of techniques, including grid search, random search, hyperband, Bayesian optimization, evolutionary algorithms, population-based training (PBT), etc. While a possible yet naive approach to treat system parameters is to consider them as possible hyperparameters, this leads to longer training periods (see Table 2). PipeTune strives to optimize both accuracy and training time of DNNs, while simultaneously tuning hyper and system parameters. The key observation of PipeTune is that the backbone of popular training algorithms for DNN is stochastic gradient decent, an iterative algorithm. PipeTune exploits such repetitive patterns as a unique opportunity to improve and achieve fast system parameter tuning. As an example, Figure 2 illustrates the typical repetitive behavior of a training process. Building on this observation, we design, implement and evaluate PipeTune, a middleware solution coordinating between the DNN training applications and systems. In a nutshell, PipeTune relies on low level metrics to profile the training trials on the epoch level and make quick decisions regarding the system parameters. We show that by taking into account the system parameters, the overall tuning runtime can be greatly reduced while at the same time improving the model performance. We show that it is possible to include system parameters in the tuning process and ask the algorithm to optimize the ratio of accuracy to performance. Section II in Pages 3-4: As the process of tuning hyperparameters is, in most cases, crucial to find the best model performance of a given application, there are many proposed approaches and tools addressing this problem. HyperDrive is a package part of Azure Machine Learning which supports hyperparameter tuning. It follows POP’s scheduling algorithm which combines probabilistic model-based classification with dynamic scheduling and early stop techniques. Amazon SageMaker is a fully managed machine learning service. It supports automatic model tuning component that finds the best version of a model by running many training trials on the dataset using the algorithm and ranges of hyperparameters specified by the user. As our approach is an extension of pure hyperparameter tuning, the above mentioned systems and all others which focus on hyperparameter auto-tuning could profit from PipeTune. ByteScheduler is a Bayesian optimization approach. It specifically focuses on auto-tune tensor credit and partition size for different training models under various networking conditions. ByteScheduler uses auto-tune algorithms to find the optimal system related configurations. Instead, PipeTune allows the user to perform hyperparameter auto-tuning and finds the best system configurations independently of this process. AutoKeras [25] supports a form of system parameter tuning, by means of an adaptive search strategy for different GPU memory limits. However, instead of adapting the system parameters to the workload, as we do in PipeTune, AutoKeras limits the size of the neural networks according to the GPU memory. To the best of our knowledge, PipeTune is the first solution that efficiently combines hyper and system parameters in a holistic manner. Section 3 with FIGS. 3-4 in Pages 4-6: A hyperparameter is a configuration external to the model. Choosing the right hyperparameters during the tuning phase is key, as the output accuracy of the trained models can vary significantly. As a result, hyperparameter optimization outputs a tuple of hyperparameters that yields an optimal model which minimizes a predefined loss function on given independent data. Typically the selection criterion considered is model accuracy. However, the hyperparameters values will impact model accuracy, its training time and the energy footprint. The former is typically related to the utility of the trained model, the latter two to its costs. We define system parameters the configurable resources of the underlying computing infrastructure where the training will execute (e.g., memory, CPU cores, CPU frequency). Typically, the hyperparameter optimization fixes the same system parameters for each trial, although they might benefit from different configurations. Regarding the energy observations, we estimate the overall energy consumption of the cluster by calculating the trapezoidal integral of the power values collected every second during training. In summary, these preliminary results in FIG. 3 show the delicate trade-offs between hyper and system parameters. One needs to balance them all towards optimal values, such that the underlying system achieves the best training performance without compromising the model accuracy. A workload is a tuple pairing a model and dataset. In this work, we only consider the training phase of DNN workloads. Moreover, we assume that this training phase includes parameters tuning on top of learning the weights of the model. Hence, tuning a single workload consists of multiple training trials, each divided into epochs. Each epoch involves one forward and one backward pass of the entire input dataset. It is a common practice to train the same model with different datasets, as well as different models using the same dataset. Figure 4 depicts this practice. Our approach leverages the similarity existent among such jobs to improve the tuning performance. Section 4 with FIG. 5 and Table 2 in Pages 5-6: Consider a state-of-the-art hyperparameter auto-tuning system, Tune [35], an open-source library implemented in Python supporting an extensive list of hyperparameters optimization algorithms. First, we consider two versions of Tune. In V1, it is used out-of-the-box to perform hyperparameters tuning with the objective of maximizing accuracy, without taking the system parameters into account. In this version all trials run with the same default system parameters. Then, in V2, the system parameters are included in the list of parameters to be tuned. This second version requires the resources used by each trial to be manually controlled. Also, the objective function must be adapted to maximize the ratio accuracy to duration (i.e., maximize the accuracy and minimize the duration), rather than restricting it to accuracy only. Figure 5 shows the results of Tune’s performance characterization under various system conditions (i.e., the number of cores assigned to the tuning job and the number of jobs assigned to the same logical cores). Tuning under different system conditions significantly impacts the performance of the model being trained. There are only a few system configurations that yield improvements over the baseline for error and training time. Some system configurations caused the tuning to trade better accuracy for faster training. Hyperparameter tuning without system conditions can produce less efficient models. The results in Table 2 show us the following. First, arbitrary values, if not correctly chosen lead to both worse accuracy and training time. Second, if the user’s focus is accuracy only, then PipeTune’s accuracy results are comparable to Tune V1 however achieve in a lower tuning time. Third, if the user’s focus is both accuracy and training time, then PipeTune’s training time results are comparable to Tune V2 but with better accuracy and lower tuning time as well. Section 5 with FIGS. 6-8 and Algorithm 1 in Pages 6-9: One of the first challenges of applying deep learning algorithms in practice is to find the appropriated hyperparameter values for a given workload. In the following we refer to these types of jobs as HPT Jobs (i.e., Hyperparameters Tuning Jobs). A given HPT Job takes as input a given workload, a set of parameters, its respective set of range values, an objective function and the metric of interest (e.g., accuracy, performance, energy). This job spawns a collection of training trials based on the possible values of the parameters, following a given search algorithm (e.g., GridSearch, HyperBand). Each training trial takes as input the workload and a set of fixed values for the parameters of interest, where these values belong to their respective given ranges. These trials can run either sequentially or in parallel depending on the setup. They produce a trained model and a score for the given parameters values. Scores correspond to the metric of interest defined by the user. The optimal set of parameters values is chosen by applying the objective function to the scores. Figure 6 illustrates this process. We consider a deep learning cluster consisting of 𝑁 nodes, each containing 𝐶 cores and 𝑀 GB of memory. HPT Jobs are scheduled in a FIFO manner. We categorize these jobs in the following two main types: Type-I: tuning the same model for different datasets (e.g., recommendation engines). and Type-II: tuning different models for the same dataset (e.g., computer vision). Both types of tuning jobs can still be divided into two sub-types: (a) same set of hyperparameters and ranges, and (b) same set of hyperparameters but different ranges. A key observation is that these jobs could benefit from previously computed results for other jobs in the same category to converge faster. Moreover, training trials spawned by the same HPT Job run all with the same system parameters even though they might require different resources configuration. In summary, our problem’s input consists of an HPT Job with the objective of achieving either maximum accuracy, or maximum accuracy with minimum training time. The former must output the best possible hyperparameters leading to the highest accuracy, independent of training time. For the latter, a combination of optimal hyper and system parameters is expected which leads to the highest accuracy and lowest training time. Figure 7 depicts the architecture components of PipeTune design and the main workflow. While training hyperparameters, a trial is a single training run with a fixed initial hyperparameter configuration. In order to find the best values for a given set of hyperparameters, the system executes a collection of trials, supervised by a given tuning library (e.g., Vizier, Tune) and using one of the supported trial scheduling algorithms (e.g., GridSearch, HyperBand). PipeTune enhances the tuning of system parameters following a pipelined parallelism approach. That is, within each trial, a collection of sub-trials is executed, with the goal of defining the best system configurations for a given optimization function and metric of interest. This sub-trial consists of varying the system configuration on the epoch level and monitoring the system itself as well as the metrics of interest. The execution of sub-trials is controlled by PipeTune, which may also rely on different underlying scheduling algorithms. Algorithm 1 details the pipelined approach. Function train (lines 1-5) is executed during a trial for a given workload (i.e., model and dataset). After initiating the model training using the hyperparameter configuration given for that trial, tuneSystem (line 3) is invoked asynchronously. The profiling phase (lines 7) is initiated for this given trial with the objective of characterizing the workload properties and its systems requirements. We rely on kernel performance counters (e.g., CPU cycles memory stores, instructions) to gather hardware events corresponding to low-level metrics of the underlying system. Once this profiling phase is over, its outcome is used as input to a ground truth phase. This process consists of applying a similarity function (line 8) on the job’s profile. This is done to reuse optimal configurations known by the system for other jobs with similar characteristics. If the score of this similarity function is within a specific confidence level (line 9), then the optimal known configurations are applied (line 10) and no further system metric trials are required. However, if the score does not cross the threshold, a new probing phase starts, searching the optimal system configurations for that trial. The probing requires each system configuration to be applied for a different epoch, following a given scheduling algorithm. We collect several meaningful metrics (e.g., runtime, energy) plus low-level metrics (e.g., hardware events). Then the optimization function is applied over these metrics (line 16) to identify the overall best system configuration. This process consists of iterating over the collected values for each tuple of system parameters, looking for the one which best fits the optimization function (e.g., shortest runtime, lowest energy consumption). Finally, the configuration identified as optimal is applied for the remaining iterations (line 17) and saved for further improving of the ground truth phase. The profiling component leverages hardware performance counters to collect low-level events of the system during the applications execution time. As the number of events collected per time unit is limited by the number of actual hardware counters of the CPU, we filter out highly correlated as well as unsupported events. During Ground Truth phase, new incoming HPT Jobs exploit the ground truth results from historical data collected during the previously completed jobs with similar system characteristics, to accelerate their system-parameter tuning phases. The implementation of ground truth is done as a separate module which is used by PipeTune. Our currently implementation relies on the scikit-learn machine learning library for Python [44] which already supports several clustering algorithms (e.g., affinity propagation, mean-shift, DBSCAN, OPTICS, Birch). The probing phase profiles a given set of workloads in different system conditions, in order to collect sufficient data for a warm start of the ground truth component. In practice, the ground truth model is refined as the similarity of the incoming jobs with the historical data of the system starts to decrease. When this happens, we launch a grid search on the system-parameters at the epoch granularity, yet other search strategies are possible. In this case, the tuning of system parameters for the current job is performed directly on the analytical data collected. Moreover, this collected data is saved to be taken into account once re-clustering is applied. We decide upon the necessity to launch a new probing or not for a given workload based on the similarity score outputted from the ground truth phase. When using k-means, the threshold matches the distance from the new set of data points to their current cluster’s centroid. The distance is compared against the models’ inertia, to measure the reliability of the prediction, or else if a re-clustering is needed. Section 6 in Page 9: Tune is a Python library for hyperparameter search, optimized for deep learning and deep reinforcement learning. Tune provides several trial schedulers based on different optimization algorithms. While we select HyperBand for the remainder of this work, Tune allows to switch among the available ones, as well as to implement new ones. As a consequence, PipeTune indirectly supports all its hyperparameter optimization algorithms. The training applications are executed by BigDL, a distributed deep learning framework on top of Apache Spark. BigDL supports TensorFlow and Keras, hence PipeTune supports models defined using such frameworks. The Ground Truth module is based on a battle-tested k-means implementation openly available in the scikit-learn machine learning library for Python. Section 7 with FIGS. 9-14 in Pages 9-13: Our main findings are: 1. PipeTune achieves significant tuning speedups without affecting model performance (i.e., accuracy); 2. By speeding up the tuning process, we also have a more energy efficient approach, not only due to the runtime reduction but also because of the more efficient utilization of system resources; 3. The proposed approach is sensitive to varying system loads as this is also reflected on the events used to profile and our system adapts on a fine granularity (i.e., epochs level). There are several potential hyperparameters to tune. For practical reasons, in our evaluation we select the 5 described below. 1. Batch size. Number of samples to work through before updating the internal model parameters. Range: [32 - 1024]. 2. Dropout rate. Dropout randomly selects neurons to be ignored during training. The dropout rate value defines the fraction of input to drop to prevent overfitting [41]. Range: [0.0 – 0.5]. 3. Embedding dimensions. Word embeddings provide a mean of transfer learning. Range: [50 – 300]. 4. Learning rate. Rate at which the neural network weights change between iterations. Range: [0.001 - 0.1]. Number of epochs Number times that the learning algorithm will work through the entire training dataset. Range: [10 - 100]. For the purpose of this evaluation, we restrict the list of parameters to number of cores and memory. However, the same mechanisms can be applied to any other parameter of interest (e.g., CPU frequency, CPU voltage). In order to build our initial similarity model we rely on profiling data of the workloads described in Table 3. For each workload, we vary the system configurations as follows. Memory allocation can be 4GB, 8GB, 16GB, and 32GB. The total number of cores that could be allocated were 4, 8, or 16. Finally, batch size could take the values 32, 64, 512, or 1024. In total, this sums up to 48 different configurations for each workload. We consider a single-tenancy scenario, and assume each HPT Job runs in a dedicated cluster, where the required resources demanded by the system parameters are available and exclusive for a given tenant. This prevents interference caused by other jobs co-located on the same cluster. However, as a given HPT Job spawns several training trials asynchronously, the cluster still remains shared among these sub-jobs. We evaluate how PipeTune performs in such stable setting, comparing it against Tune V1 and Tune V2, for all the workloads. Profiling is a fundamental part of our system design and essential for the decision making process. During the profiling of a given epoch, the extra computation introduce additional load, depending on the system configuration. However, as this profiling overhead only occurs in the epoch granularity and does not apply for all the epochs, the performance benefits resulting from tuning the system-parameters overtake the measured overhead. Next, we evaluate PipeTune in a multi-tenancy scenario (i.e., a shared cluster handling multiple HPT Jobs). In this case, we show the average response time of jobs as an indicator of performance. We consider that jobs arrive randomly with the interarrival times being exponentially distributed. For the case where two workload types are considered together, each of them corresponds to 50% of the overall jobs (i.e., equally balanced). In all cases, within a given workload type, the workloads are chosen following a round-robin strategy. The portion of overall unseen jobs corresponds to 20%). Claim 34 Rocha discloses all elements as stated in Claim 33 and further discloses wherein the characteristics of the distributed training environment comprise a combination of distributed training (DT) configuration and computing resource budget (Rocha, Section 5.1 in Pages 6-7: We consider a deep learning cluster consisting of 𝑁 nodes, each containing 𝐶 cores and 𝑀 GB of memory). Claim 35 Rocha discloses all elements as stated in Claim 34 and further discloses wherein the DT configuration comprises at least one of a number of parameter servers, and a number of worker nodes, and wherein the computing resource budget comprises a central processing unit (CPU) allocation, a memory allocation, and disk space (Rocha, Section 5.1 in Pages 6-7: We consider a deep learning cluster consisting of 𝑁 nodes, each containing 𝐶 cores and 𝑀 GB of memory). Claim 36 Rocha discloses all elements as stated in Claim 33 and further disclose wherein the selection of the optimal resource configuration comprises application of Bayesian optimization techniques to determine which of the plurality of resource configurations are to be tested (Rocha, Section 1 of Pages 1-3: The majority of the existing tuning solutions restrict themselves to the sole hyperparameter tuning using a variety of techniques, including grid search [21], random search [9], hyperband [32], Bayesian optimization [50, 52], evolutionary algorithms [55, 62], population-based training (PBT) [24], etc. PipeTune strives to optimize both accuracy and training time of DNNs, while simultaneously tuning hyper and system parameters. Section 2 of Pages 3-4: AutoKeras [25] enables Bayesian optimization to guide the network morphism for efficient neural network architecture search. ByteScheduler is a Bayesian optimization approach. It specifically focuses on auto-tune tensor credit and partition size for different training models under various networking conditions. ByteScheduler uses auto-tune algorithms to find the optimal system related configurations. Instead, PipeTune allows the user to perform hyperparameter auto-tuning and finds the best system configurations independently of this process. FIG. 7 in Page 7 with Section 5.2 in Pages 7-8: Hyperparameter tuning includes Grid search, Generic optimization, Bayesian, Random search, Gradient optimization, and Hyperband; Figure 7 depicts the architecture components of PipeTune design and the main workflow. While training hyperparameters, a trial is a single training run with a fixed initial hyperparameter configuration. In order to find the best values for a given set of hyperparameters, the system executes a collection of trials, supervised by a given tuning library (e.g., Vizier, Tune) and using one of the supported trial scheduling algorithms (e.g., GridSearch, HyperBand). PipeTune enhances the tuning of system parameters following a pipelined parallelism approach. That is, within each trial, a collection of sub-trials is executed, with the goal of defining the best system configurations for a given optimization function and metric of interest. This sub-trial consists of varying the system configuration on the epoch level and monitoring the system itself as well as the metrics of interest. The execution of sub-trials is controlled by PipeTune, which may also rely on different underlying scheduling algorithms). Claim 37 Rocha discloses all elements as stated in Claim 33 and further disclose monitoring resource usage metrics during the execution of the training jobs (Rocha, Section 5.2 with FIG. 7 and Algorithm 1 in Pages 7-8: PipeTune enhances the tuning of system parameters following a pipelined parallelism approach. That is, within each trial, a collection of sub-trials is executed, with the goal of defining the best system configurations for a given optimization function and metric of interest. This sub-trial consists of varying the system configuration on the epoch level and monitoring the system itself as well as the metrics of interest. The execution of sub-trials is controlled by PipeTune, which may also rely on different underlying scheduling algorithms. The profiling phase (lines 7) is initiated for this given trial with the objective of characterizing the workload properties and its systems requirements. This process is done at the granularity of epochs for the currently running trial. We rely on kernel performance counters (e.g., CPU cycles memory stores, instructions) to gather hardware events corresponding to low-level metrics of the underlying system. Once this profiling phase is over, its outcome is used as input to a ground truth phase. This process consists of applying a similarity function (line 8) on the job’s profile. This is done to reuse optimal configurations known by the system for other jobs with similar characteristics. If the score of this similarity function is within a specific confidence level (line 9), then the optimal known configurations are applied (line 10) and no further system metric trials are required. However, if the score does not cross the threshold, a new probing phase starts, searching the optimal system configurations for that trial. The probing requires each system configuration to be applied for a different epoch, following a given scheduling algorithm. We collect several meaningful metrics (e.g., runtime, energy) plus low-level metrics (e.g., hardware events). Then the optimization function is applied over these metrics (line 16) to identify the overall best system configuration. This process consists of iterating over the collected values for each tuple of system parameters, looking for the one which best fits the optimization function (e.g., shortest runtime, lowest energy consumption). Finally, the configuration identified as optimal is applied for the remaining iterations (line 17) and saved for further improving of the ground truth phase. Section 5.3 in Page 8: The profiling component leverages hardware performance counters to collect low-level events of the system during the applications execution time. To mitigate the potential profiling errors, we store the average of results during each epoch’s time window. Section 6 in Page 9: Finally, as storage backend, we leverage InfluxDB (v1.7.4), an open-source time series database. It offers a convenient InfluxDB-Python client for interacting with InfluxDB which we use to query information regarding the collected system metrics. Section 7.3 with FIGS.11-12 in Pages 11-12 with Table 3 in Page 9: Figure 11 presents the results of model accuracy, training and tuning runtime, and overall cluster energy consumption of offline HPT Jobs for the different workloads described in Table 3. Figure 11 (d) reports the energy results. The overall energy consumption of the cluster is directly affected both by the performance decays and gains. Figure 12 compares Tune V1, Tune V2 and PipeTune on a single node. The Type-III workloads used in these experiments have shorter epochs and each a different CNN model). Claim Rejections - 35 USC § 103 The following is a quotation of 35 U.S.C. 103 which forms the basis for all obviousness rejections set forth in this Office action: A patent for a claimed invention may not be obtained, notwithstanding that the claimed invention is not identically disclosed as set forth in section 102, if the differences between the claimed invention and the prior art are such that the claimed invention as a whole would have been obvious before the effective filing date of the claimed invention to a person having ordinary skill in the art to which the claimed invention pertains. Patentability shall not be negated by the manner in which the invention was made. Claims 25-32 are rejected under 35 U.S.C. 103 as being unpatentable over Qiao in view of Fischer et al. ("Machines Tuning Machines: Configuring Distributed Stream Processors with Bayesian Optimization", 2015 IEEE International Conference on Cluster Computing, Sep 8-11, 2015, pp. 22-31), hereinafter Fischer. Claim 25 Qiao discloses all the elements as stated in claim 24 and further discloses wherein the DT loop further comprises a (Qiao, Section 4.1 in Page 6: An instance of PolluxAgent is started with each training job. During training, it continually measures the job’s gradient noise scale and system throughput, and reports them to PolluxSched at a fixed interval. It also uses this information to determine the most efficient batch size for its job given its current resource allocations, and adapts its job’s learning rate to this batch size using AdaScale. PolluxAgent measures the time taken per iteration, Titer, and records the triple ( a ,   m , Titer) for all combinations of resource allocations a and batch size m encountered during its lifetime. Periodically, PolluxAgent fits the parameters θsys to all of the throughput data collected so far. Specifically, we minimize the root mean squared logarithmic error (RMSLE) between Eqn. 11 and the collected data triples, using L-BFGS-B. PolluxAgent then reports the updated values of θsys and φ t to PolluxSched. To ensure that Pollux finds efficient resource allocations through systematic exploration, we impose several priors which bias θsys towards the belief that throughput scales perfectly with more resources, until such resource configurations are explored. Each job starts with a single GPU, and is initially assumed to scale perfectly to more GPUs. PolluxSched is then encouraged to allocate more GPUs and/or nodes to the job, naturally as part of its resource optimization, until the PolluxAgent can estimate θsys more accurately. Finally, to prevent a job from being immediately scaled out to arbitrarily many GPUs, we restrict the maximum number of GPUs which can be allocated to at most twice the maximum number of GPUs the job has been allocated in its lifetime. With θsys, mgns, and m0, which fully specify the DL job’s GOODPUT function at its current training progress, PolluxAgent determines the most efficient batch size m * = argm ax m ⁡ GOODPUT a ,   m , where a is the job’s current resource allocation. Once a new batch size is found, the job will use it for its subsequent training iterations, using AdaScale to adapt its learning rate appropriately. As the job’s statistical efficiency changes over time, PolluxAgent will periodically re-evaluate the most efficient batch size and learning rate. Section 4.2 in Pages 6-7: The PolluxSched periodically allocates (or re-allocates) resources for every job in the cluster. To determine a set of efficient cluster-wide resource allocations, it uses a genetic algorithm to maximize a fitness function which is defined as a weighted mean across speedups for each job in Eqn. 14. A is an allocation matrix with each row Aj being the placement vector for a job j, thus Ajn is the number of GPUs on node n allocated to job j. Intuitively, maximizing FITNESS causes PolluxSched to allocate more resources to jobs that achieve a high SPEEDUP when provided with many GPUs (i.e. jobs that scale well). We define the speedup of each job as the factor of goodput improvement using the given placement vector over using a single process with the optimal batch size in Eqn. 15, where GOODPUTj is the goodput of job j at its current training iteration. Similar to the PolluxAgent, PolluxSched performs the maximizations in the numerator and denominator using golden-section search. Although PolluxSched can work well with all weights wj set to 1, we provide the ability to re-weight jobs as an additional tuning knob for cluster operators. We define the weight wj of job j as Eqn. 16. The weight of a job is 1 if its current total GPU-time is at most GPUTIME_THRES, and decays gradually thereafter. Section 4.2.1 with FIG. 5 in Page 7: Our genetic algorithm operates on a population of distinct allocation matrices (see Fig. 5). During each generation of the algorithm, existing allocation matrices are first randomly mutated, then crossed over to produce offspring allocation matrices, and finally modified to satisfy node resource constraints. A constant population size of distinct allocation matrices is maintained after each generation by discarding the allocation matrices with the lowest objective values (according to Eqn. 14). After several generations of the genetic algorithm, the allocation matrix with the highest fitness score is applied to the jobs running in the cluster. Mutation: Each element Ajn is mutated with probability 1/N, where N is the total number of nodes. Thus, each job suffers on average one mutation in each generation. Crossover: When two allocation matrices are crossed over, their rows are randomly mixed. In other words, the offspring allocation matrix consists of job allocations which are randomly selected between its two parent allocation matrices. Repair: After the mutation and crossover operations, the resultant allocation matrices may no longer satisfy resource constraints, and try to request more GPUs than are available on a node. To address this issue, random elements are decremented within columns of the allocation matrix that correspond to overcapacity nodes, until the GPU resource constraints are satisfied. Penalty: Each time a job is re-allocated to a different set of GPUs, it will need to save a checkpoint of its current model parameters, and restart using its new allocation of GPUs. To prevent an excessive number of re-allocations, when PolluxSched evaluates the fitness function for a given allocation matrix, it applies a penalty for every job that needs to restart. Interference avoidance: When multiple distributed DL jobs share a single node, their network usage while synchronizing gradients and model parameters may interfere with each other, causing both jobs to slow down. PolluxSched mitigates this issue by disallowing different distributed jobs (each using GPUs across multiple nodes) from sharing the same node. Section 4.2.2 in Pages 7-8: In cloud computing environments, GPU nodes can be dynamically provisioned during periods of high demand to decrease job completion time, as well as released during periods of low demand to decrease cost of cluster resources. Beyond adapting to the number and sizes of DL jobs, co-adaptive scheduling presents a unique opportunity for cluster auto-scaling. Since many distributed DL jobs tend to increase in statistical efficiency as training progresses, it may be more cost-effective to provision more cloud resources during the later iterations of a large training job, rather than earlier on. In order to decide when to request or release nodes in the cloud, we first define a measure of cluster resource utility for a given allocation matrix as Eqn. 17. A cluster operator defines a LOW_UTIL_THRES and a HIGH_UTIL_THRES. PolluxSched will request additional nodes (when UTILITY(A) is high) or release existing nodes (when UTILITY(A) is low) until UTILITY(A) falls within this range, where A is the current allocations applied in the cluster. The cluster operator also defines a MIN_NODES and MAX_NODES, which limit the size of the cluster. PolluxSched performs a binary search for the desired number of nodes, with the assumption that UTILITY decreases with increasing numbers of nodes. At each step of the binary search, PolluxSched runs its genetic algorithm to evaluate the UTILITY of the cluster size being queried. In the end, the cluster size with a UTILITY value closest to (LOW_UTIL_THRES+HIGH_UTIL_THRES)=2 is selected). Qiao fails to explicitly disclose a Bayesian optimization configuration generator. Fischer teaches a system and a method related to configuring a distributed computing framework (Fischer, Abstract in Page 22), wherein a Bayesian optimization configuration generator (Fischer, Abstract and Section I in Page 22: Modern distributed computing frameworks such as Apache Hadoop, Spark, or Storm distribute the workload of applications across a large number of machines. Whilst they abstract the details of distribution they do require the programmer to set a number of configuration parameters before deployment. These parameter settings (usually) have a substantial impact on execution efficiency. Finding the right values for these parameters is considered a difficult task and requires domain, application, and framework expertise. In this paper, we propose a machine learning approach to the problem of configuring a distributed computing framework. Specifically, we propose using Bayesian Optimization to find good parameter settings. In an extensive empirical evaluation, we show that Bayesian Optimization can effectively find good parameter settings for four different stream processing topologies implemented in Apache Storm resulting in significant gains over a parallel linear approach. To address the tedious manual parameter experimentation, this paper proposes an automated process based on Bayesian Optimization for finding optimal parameter configurations. Specifically, we present empirical results from a series of experiments in which we evaluated the suitability of Bayesian Optimization for the configuration of distributed stream processing systems built using Apache Storm. Our contributions are: • We present an auto-configuration approach for distributed stream processing systems (SPS) using Bayesian Optimization. • We provide an extensive empirical evaluation showing the effectiveness of our approach on a cluster of 80 machines (320 cores) running Storm topologies (applications) of varying sizes and characteristics. • We introduce a reusable benchmark consisting of a set of operator graphs as well as generation approach. Section II.B in Page 23: Bayesian Optimization is a probabilistic technique to optimize systems with unknown cost functions. It has successfully been applied in cases where the performance of systems is strongly dependent on configuration parameters, and no mathematical closed-form cost model is known such as finding good hyperparameter settings in machine learning problems (e.g., classification or feature selection). There are several Bayesian Optimization frameworks (e.g., Spearmint3 [17], SMAC4[18], HyperOpt5, or BayesOpt6) available for research. We are not aware of any previous work that has investigated the applicability of Bayesian Optimization for the configuration of distributed systems or for distributed stream processing systems in particular; Section III with FIG. 1 in Pages 23-25: This section describes how we employ Bayesian Optimization to configure a distributed stream processing system based on the Storm distributed real-time computation framework. In contrast to batch-based distributed systems such as Apache MapReduce, Storm ingests data continuously. As in MapReduce, a Storm application allows the user to partition the data and to distribute parts of the processing across a compute cluster. A Storm application—a topology—is a directed graph consisting of spout and bolt nodes as depicted in Figure 1 on the left. Spouts emit data to downstream nodes. Bolts consume data from upstream nodes and emit data to downstream nodes. Spout nodes are typically used to connect a Storm topology to external data sources such as queues, web-services, or file systems. For each spout and bolt, the programmer defines how many instances of this node should be created in the physical instantiation of the topology – the task instances. This results in a physical topology depicted on the right of Figure 1, which is different from the logical representation. The parameter used to define the degree of parallelism of a node is called a parallelism hint, as Storm may change these hints for consistency purposes. The task instances, or tasks, are distributed across all machines of the compute cluster to which a topology has been assigned. Each edge in the topology graph defines a grouping strategy according to which messages that pass between the nodes—the tuples—are sent to downstream nodes. Storm offers a number of configuration parameters that allow the programmer, as well as the system administrators, to configure various aspects of the system. Table I lists the parameters that we used in our evaluations: parameters that are most commonly tuned are the already mentioned parallelism hints, the batch size, and the batch parallelism of a topology. Note that the “parallelism hints” parameter is not one single value, but a list of values that contains one number for each node of the topology. While any single one of these configuration options impacts the runtime behavior, overall performance is a result of the combination of all of these parameters working together. To tackle the problem of choosing good configuration parameters, we investigate the possibility of having a computer program choose these parameters. To this end, we employ the technique of Bayesian Optimization. Bayesian Optimization has first been proposed by Jonas Mockus as an optimization strategy for situations in which the objective function is a non-convex blackbox function (i.e., a function for which no closed-form solution or derivative is known). The function is assumed to be Lipschitz-continuous (i.e., smooth and does not change dramatically). Also, sampling the function is assumed to be costly, either in terms of time or money. Thus, it can pay off to invest computational resources into computing the point in the parameter space where to sample next. For our domain, we assume the function to be the actual system performance of our distributed stream processor, given all the configuration parameters chosen. Obviously, given the black-box nature of the system, no mathematical representation exists, and determining the value of the function given certain parameter settings is achieved by running the system on a cluster with these settings and, hence, is costly. The process of choosing the next set of parameters is conducted using a Bayesian approach, which combines our prior assumptions about the function with the observed performance from previous runs. Thus, we reason about the likelihood of observing the results of an evaluation run, given our prior beliefs about how the system would change in response to parameter modifications. Using the results of each evaluation run, a posterior distribution P(M|E) is computed and integrated into the model. The decision of where to sample next is made by maximizing an acquisition function. Often, Gaussian Processes are used to model the noise within the acquisition function. The purpose of the acquisition function is to balance the tradeoff between exploration and exploitation. The goal is to sample the next measurement in a region where either the uncertainty of the expected performance is high, the expected performance is high, or both. Bayesian Optimization is an iterative process in which we sample an objective function repeatedly. Our prior beliefs about f can be expressed as a prior distribution P(f). We then collect observations (measured samples) and add them to the set D1:t = {x1:t, y1:t} of all evidence to date. In each step, we update our posterior belief with the newly collected evidence: P(f|D1:t) ∝ P(D1:t|f)P(f)). The new evidence is used to fit a Gaussian Process (GP) that describes our prior beliefs of how f is distributed. f(x) ∼ GP(m(x), k(x, x')) where m is the mean function at position x and k is the covariance function depending on x as well as on the closest previously sampled point at x'. The result is a function estimating the expected performance of any parameter value combination given some confidence interval. An acquisition function u(x|D) is built using these two parameters (expected performance and confidence intervals) that are derived from the data D. The goal of the acquisition function is to create a tradeoff between exploration (try points with high uncertainty/variance) and exploitation (try points with a high expected performance). Hence, the next sample point x is determined by maximizing u(x) (i.e. the x where the tradeoff between exploration and exploitation is optimal): xt+1 = argmaxxu(x|D1:t). There are several different ways of defining the acquisition function such as Probability of Improvement (PI), Expected Improvement (EI), or GP Upper Confidence Bound to name the most common ones. In this paper, we use Expected Improvement, as it provides a good tradeoff between exploration and exploitation and it is the method implemented in Spearmint, the toolkit we use in our experiments. So the next x would be chosen at the position, where the expected improvement between the new sample point (ft+1(x)) and the current best sample point (fmax) is maximized. In this project, we leverage Spearmint for the following reasons: first, it showed good performance in comparison with other main-stream Bayesian Optimization frameworks. Second, it is well documented and its source code is openly available. Last, it supports pausing and resuming the optimization process, a feature that turned out to be important in our evaluation setup) Qiao and Fischer are analogous art because they are from the same field of endeavor, a system and a method related to configuring a distributed computing framework. Therefore, it would have been obvious to one of ordinary skill in the art before the effective filling date of the claimed invention to apply the teaching of Fischer to Qiao. Motivation for doing so would shorten the number of evaluation runs necessary to achieve convergence (Fischer, Section V.B in Pages 28-29) . Claim 26 Qiao in view of Fischer discloses all elements as stated in Claim 25 and further discloses wherein the running of the one or more ML jobs occurs in a resource allocation loop, and wherein while the one or more ML jobs are running, monitoring performance of the one or more ML jobs (Qiao, Section 3 in Pages 4-5: We present how goodput can be defined, measured, and predicted for DL jobs, taking into account both system throughput and statistical efficiency. Goodput predictions are leveraged by Pollux to jointly optimize cluster-wide resource allocations and batch sizes. Definition 3.1. (Goodput) The goodput of a DL job at iteration t is the product between its system throughput and its statistical efficiency at iteration t. GOODPUTt (a, m)=THROUGHPUT(a, m)´EFFICIENCYt (m). a ∈ R N is an allocation vector, where a n is the number of GPUs allocated from node n, and m is the batch size. An initial batch size m0 and learning rate η0 are selected by the user when submitting their job. As the job runs, Pollux profiles its execution to learn and refine predictive models for both THROUGHPUT and EFFICIENCY. Using these predictive models, Pollux periodically re-tunes a and m for each job, according to cluster-wide resource availability and performance. The learning rate η is re-tuned using AdaScale. Goodput and statistical efficiency are measured relative to the initial batch size m0 and learning rate η0, and Pollux only considers batch sizes which are at least the initial batch size, i.e.. m ≥ m0. The statistical efficiency of a DL job using batch size m ≥ m0 is the amount of progress made per training example using m, relative to using m0. This quantity can be framed in terms of the gradient noise scale φ t . We can write a concrete measure of statistical efficiency as EFFICIENCYt (m) = ( φ t + m 0 . )/( φ t + m ). Here, we can write the gradient noise scale as φ t = m 0 σ t 2 / μ t 2 , where σ t 2 = Var g ^ ( t ) is the variance and μ t 2 = E g ^ ( t ) 2 the squared norm of the gradient at iteration t using batch size m0. EFFICIENCYt (m) reflects the lifetime-dependent trends exhibited by the true statistical efficiency. The standard way of estimating σ t and μ t involves calculating the sample variance of the local gradient estimates g ^ k ( t ) at each iteration. This can be done efficiently when there are multiple data-parallel processes, by using the different values of g ^ k ( t ) already available on each process. To model and predict the system throughput for data-parallel DL, we aim to predict the time spent per training iteration, Titer, given an allocation vector a and batch size m, and then calculate the throughput as THROUGHPUT a , m = m / T i t e r a , m . We start by separately modeling Tgrad, the time in each iteration spent computing local gradient estimates, and Tsync, the time in each iteration spent averaging gradient estimates and synchronizing model parameters across all GPUs. We model Tgrad as Tgrad ( a , m)=αgrad+βgrad·m/K, where m is the overall batch size, K = ∑ n a n is the number of allocated GPUs, and αgrad and βgrad are learnable parameters. We model Tsync as Eqn. 10, where N is the number of physical nodes occupied by at least one replica. α s y n c l o c a l and β s y n c l o c a l are the constant and retrogression parameters for when all processes are co-located onto the same node. α s y n c n o d e and β s y n c n o d e are the analogous parameters for when at least two process are located on different nodes. Modern DL frameworks can partially overlap Tgrad and Tsync by overlapping gradient computation with network communication. The degree of this overlap depends on structures in the specific DL model being trained, like the ordering and sizes of its layers. To capture the overlap between Tgrad and Tsync, we model Titer as Eqn. 11, where γ≥1 is a learnable parameter. Eqn. 11 has the property that Titer=Tgrad+Tsync when γ =1, and smoothly transitions towards Titer =max(Tgrad, Tsync) as γ→∞. Section 4 with FIG. 4 in Pages 5-6: Pollux performs adaptation at two distinct granularities. First, at a job-level granularity, Pollux dynamically tunes the batch size and learning rate for best utilization of the allocated resources. Second, a the cluster-wide granularity, Pollux dynamically re-allocates resources, driven by the goodput of all jobs sharing the cluster. To achieve this co-adaptivity in a scalable way, Pollux’s design consists of two primary components. First, a PolluxAgent runs together with each job. It measures the gradient noise scale and system throughput for that job, and tunes its batch size and learning rate for efficient utilization of its current allocated resources. PolluxAgent periodically reports the goodput function of its job to the PolluxSched. Second, the PolluxSched periodically optimizes the resource allocations for all jobs in the cluster, taking into account the current statistical efficiency for each job. It uses the goodput function to predict a job’s training performance when allocated different resources. PolluxAgent and PolluxSched co-adapt to each other. While PolluxAgent adapts each training job to make efficient use of its allocated resources, PolluxSched dynamically re-allocates each job’s resources, taking into account the PolluxAgent’s ability to tune its job. Section 4.1 in Page 6: During training, it continually measures the job’s gradient noise scale and system throughput, and reports them to PolluxSched at a fixed interval. It also uses this information to determine the most efficient batch size for its job given its current resource allocations, and adapts its job’s learning rate to this batch size using AdaScale. PolluxAgent measures the time taken per iteration, Titer, and records the triple ( a ,   m , Titer) for all combinations of resource allocations a and batch size m encountered during its lifetime. Periodically, PolluxAgent fits the parameters θsys to all of the throughput data collected so far. Specifically, we minimize the root mean squared logarithmic error (RMSLE) between Eqn. 11 and the collected data triples, using L-BFGS-B. PolluxAgent then reports the updated values of θsys and φ t to PolluxSched. To ensure that Pollux finds efficient resource allocations through systematic exploration, we impose several priors which bias θsys towards the belief that throughput scales perfectly with more resources, until such resource configurations are explored. Each job starts with a single GPU, and is initially assumed to scale perfectly to more GPUs. PolluxSched is then encouraged to allocate more GPUs and/or nodes to the job, naturally as part of its resource optimization, until the PolluxAgent can estimate θsys more accurately. Finally, to prevent a job from being immediately scaled out to arbitrarily many GPUs, we restrict the maximum number of GPUs which can be allocated to at most twice the maximum number of GPUs the job has been allocated in its lifetime. With θsys, mgns, and m0, which fully specify the DL job’s GOODPUT function at its current training progress, PolluxAgent determines the most efficient batch size m * = argm ax m ⁡ GOODPUT a ,   m , where a is the job’s current resource allocation. Once a new batch size is found, the job will use it for its subsequent training iterations, using AdaScale to adapt its learning rate appropriately. As the job’s statistical efficiency changes over time, PolluxAgent will periodically re-evaluate the most efficient batch size and learning rate. Section 4.2 in Pages 6-7: The PolluxSched periodically allocates (or re-allocates) resources for every job in the cluster. To determine a set of efficient cluster-wide resource allocations, it uses a genetic algorithm to maximize a fitness function which is defined as a weighted mean across speedups for each job in Eqn. 14. A is an allocation matrix with each row Aj being the placement vector for a job j, thus Ajn is the number of GPUs on node n allocated to job j. Intuitively, maximizing FITNESS causes PolluxSched to allocate more resources to jobs that achieve a high SPEEDUP when provided with many GPUs (i.e. jobs that scale well). We define the speedup of each job as the factor of goodput improvement using the given placement vector over using a single process with the optimal batch size in Eqn. 15, where GOODPUTj is the goodput of job j at its current training iteration. Similar to the PolluxAgent, PolluxSched performs the maximizations in the numerator and denominator using golden-section search. Although PolluxSched can work well with all weights wj set to 1, we provide the ability to re-weight jobs as an additional tuning knob for cluster operators. We define the weight wj of job j as Eqn. 16. The weight of a job is 1 if its current total GPU-time is at most GPUTIME_THRES, and decays gradually thereafter. Section 4.2.1 with FIG. 5 in Page 7: Our genetic algorithm operates on a population of distinct allocation matrices (see Fig. 5). During each generation of the algorithm, existing allocation matrices are first randomly mutated, then crossed over to produce offspring allocation matrices, and finally modified to satisfy node resource constraints. A constant population size of distinct allocation matrices is maintained after each generation by discarding the allocation matrices with the lowest objective values (according to Eqn. 14). After several generations of the genetic algorithm, the allocation matrix with the highest fitness score is applied to the jobs running in the cluster. Each time a job is re-allocated to a different set of GPUs, it will need to save a checkpoint of its current model parameters, and restart using its new allocation of GPUs. To prevent an excessive number of re-allocations, when PolluxSched evaluates the fitness function for a given allocation matrix, it applies a penalty for every job that needs to restart. When multiple distributed DL jobs share a single node, their network usage while synchronizing gradients and model parameters may interfere with each other, causing both jobs to slow down. PolluxSched mitigates this issue by disallowing different distributed jobs (each using GPUs across multiple nodes) from sharing the same node. Section 4.2.2 in Pages 7-8: In cloud computing environments, GPU nodes can be dynamically provisioned during periods of high demand to decrease job completion time, as well as released during periods of low demand to decrease cost of cluster resources. Beyond adapting to the number and sizes of DL jobs, co-adaptive scheduling presents a unique opportunity for cluster auto-scaling. Since many distributed DL jobs tend to increase in statistical efficiency as training progresses, it may be more cost-effective to provision more cloud resources during the later iterations of a large training job, rather than earlier on. In order to decide when to request or release nodes in the cloud, we first define a measure of cluster resource utility for a given allocation matrix as Eqn. 17. A cluster operator defines a LOW_UTIL_THRES and a HIGH_UTIL_THRES. PolluxSched will request additional nodes (when UTILITY(A) is high) or release existing nodes (when UTILITY(A) is low) until UTILITY(A) falls within this range, where A is the current allocations applied in the cluster. The cluster operator also defines a MIN_NODES and MAX_NODES, which limit the size of the cluster. PolluxSched performs a binary search for the desired number of nodes, with the assumption that UTILITY decreases with increasing numbers of nodes. At each step of the binary search, PolluxSched runs its genetic algorithm to evaluate the UTILITY of the cluster size being queried. In the end, the cluster size with a UTILITY value closest to (LOW_UTIL_THRES+HIGH_UTIL_THRES)=2 is selected. Section 4.3 in Page 8: PolluxAgent inserts performance profiling code which measures the time taken for each iteration of training, as well as calculating the gradient noise scale. At a fixed time interval, PolluxAgent fits the system throughput model (Eqn. 11) to the profiled metrics collected so far, and reports the fitted system throughput parameters, along with the latest gradient statistics, to PolluxSched. After reporting to PolluxSched, PolluxAgent updates the job’s batch size, by optimizing its now up-to-date goodput function (Eqn. 6) with its currently allocated resources. PolluxSched is implemented as a service in Kubernetes. At a fixed time interval, PolluxSched runs a fixed number of generations of its genetic algorithm. It then applies the best allocation matrix to the cluster, by creating and terminating Kubernetes Pods which run the job replicas. Although allocation matrix with the highest fitness score is applied to the cluster, the entire population is saved and used to bootstrap the genetic algorithm in the next scheduling interval). Claim 27 Qiao in view of Fischer discloses all elements as stated in Claim 26 and further discloses reporting the monitored performance of the one or more ML jobs upon completion thereof to the DT loop, in response to which, the DT loop updates the current performance context with the reported, monitored performance (Qiao, Section 2 in Pages 2-4: Training a deep learning model typically involves minimizing a loss function of the form in Eqn. 1, where w [Symbol font/0xCE] R d are the model parameters to be optimized, X is the training dataset, xi are the individual samples in the training data, and l is the loss evaluated at a single sample. The loss function can be minimized using stochastic gradient descent (SGD), which repeatedly applies the update until the loss converges to a stable value. When a DL job is distributed across several nodes, its system throughput is determined by several factors, including (1) the allocation and placement of resources assigned to the job, (2) the method of distributed execution and synchronization, and (3) the batch size used by the SGD algorithm. Data-parallelism is a popular method of distributed execution for DL training. The model parameters w(t) are replicated across a set of distributed GPUs 1, …, K, and each mini-batch M(t) is divided into equal-sized partitions per node, M 1 ( t ) , … , M K ( t ) . Each GPU k computes a local gradient estimate g ^ k ( t ) using its own partition, as shown in Eqn. 4. These local gradient estimates are then averaged across all replicas to obtain the desired g ^ ( t ) , as defined by Eqn. 3. Finally, each node applies the same update using g ^ ( t ) to obtain the new model parameters w(t+1), as defined by Eqn. 2. The run-time of each training iteration is determined by two main components. First, the time spent computing the local gradient estimates g ^ k ( t ) , which we denote by Tgrad. Second, the time spent averaging the local gradients and synchronizing the model parameters across all job replicas, which we denote by Tsync. Tsync can be influenced by the placement of replicas, and is typically smaller when the replicas are co-located within the same physical node or rack, rather than spread across different nodes or racks. Using a larger batch size enables higher system throughput when scaling to more data-parallel replicas. Section 3 in Pages 4-5: We present how goodput can be defined, measured, and predicted for DL jobs, taking into account both system throughput and statistical efficiency. Goodput predictions are leveraged by Pollux to jointly optimize cluster-wide resource allocations and batch sizes. Definition 3.1. (Goodput) The goodput of a DL job at iteration t is the product between its system throughput and its statistical efficiency at iteration t. GOODPUTt (a, m)=THROUGHPUT(a, m)´EFFICIENCYt (m). a ∈ R N is an allocation vector, where a n is the number of GPUs allocated from node n, and m is the batch size. An initial batch size m0 and learning rate η0 are selected by the user when submitting their job. As the job runs, Pollux profiles its execution to learn and refine predictive models for both THROUGHPUT and EFFICIENCY. Using these predictive models, Pollux periodically re-tunes a and m for each job, according to cluster-wide resource availability and performance. The learning rate η is re-tuned using AdaScale. Goodput and statistical efficiency are measured relative to the initial batch size m0 and learning rate η0, and Pollux only considers batch sizes which are at least the initial batch size, i.e.. m ≥ m0. The statistical efficiency of a DL job using batch size m ≥ m0 is the amount of progress made per training example using m, relative to using m0. This quantity can be framed in terms of the gradient noise scale φ t . We can write a concrete measure of statistical efficiency as EFFICIENCYt (m) = ( φ t + m 0 . )/( φ t + m ). Here, we can write the gradient noise scale as φ t = m 0 σ t 2 / μ t 2 , where σ t 2 = Var g ^ ( t ) is the variance and μ t 2 = E g ^ ( t ) 2 the squared norm of the gradient at iteration t using batch size m0. EFFICIENCYt (m) reflects the lifetime-dependent trends exhibited by the true statistical efficiency. The standard way of estimating σ t and μ t involves calculating the sample variance of the local gradient estimates g ^ k ( t ) at each iteration. This can be done efficiently when there are multiple data-parallel processes, by using the different values of g ^ k ( t ) already available on each process. To model and predict the system throughput for data-parallel DL, we aim to predict the time spent per training iteration, Titer, given an allocation vector a and batch size m, and then calculate the throughput as THROUGHPUT a , m = m / T i t e r a , m . We start by separately modeling Tgrad, the time in each iteration spent computing local gradient estimates, and Tsync, the time in each iteration spent averaging gradient estimates and synchronizing model parameters across all GPUs. We model Tgrad as Tgrad ( a , m)=αgrad+βgrad·m/K, where m is the overall batch size, K = ∑ n a n is the number of allocated GPUs, and αgrad and βgrad are learnable parameters. We model Tsync as Eqn. 10, where N is the number of physical nodes occupied by at least one replica. α s y n c l o c a l and β s y n c l o c a l are the constant and retrogression parameters for when all processes are co-located onto the same node. α s y n c n o d e and β s y n c n o d e are the analogous parameters for when at least two process are located on different nodes. Modern DL frameworks can partially overlap Tgrad and Tsync by overlapping gradient computation with network communication. The degree of this overlap depends on structures in the specific DL model being trained, like the ordering and sizes of its layers. To capture the overlap between Tgrad and Tsync, we model Titer as Eqn. 11, where γ≥1 is a learnable parameter. Eqn. 11 has the property that Titer=Tgrad+Tsync when γ =1, and smoothly transitions towards Titer =max(Tgrad, Tsync) as γ→∞. Section 4 with FIG. 4 in Pages 5-6: Pollux performs adaptation at two distinct granularities. First, at a job-level granularity, Pollux dynamically tunes the batch size and learning rate for best utilization of the allocated resources. Second, a the cluster-wide granularity, Pollux dynamically re-allocates resources, driven by the goodput of all jobs sharing the cluster. To achieve this co-adaptivity in a scalable way, Pollux’s design consists of two primary components. First, a PolluxAgent runs together with each job. It measures the gradient noise scale and system throughput for that job, and tunes its batch size and learning rate for efficient utilization of its current allocated resources. PolluxAgent periodically reports the goodput function of its job to the PolluxSched. Second, the PolluxSched periodically optimizes the resource allocations for all jobs in the cluster, taking into account the current statistical efficiency for each job. It uses the goodput function to predict a job’s training performance when allocated different resources. PolluxAgent and PolluxSched co-adapt to each other. While PolluxAgent adapts each training job to make efficient use of its allocated resources, PolluxSched dynamically re-allocates each job’s resources, taking into account the PolluxAgent’s ability to tune its job. Section 4.1 in Page 6: An instance of PolluxAgent is started with each training job. During training, it continually measures the job’s gradient noise scale and system throughput, and reports them to PolluxSched at a fixed interval. It also uses this information to determine the most efficient batch size for its job given its current resource allocations, and adapts its job’s learning rate to this batch size using AdaScale. PolluxAgent measures the time taken per iteration, Titer, and records the triple ( a ,   m , Titer) for all combinations of resource allocations a and batch size m encountered during its lifetime. Periodically, PolluxAgent fits the parameters θsys to all of the throughput data collected so far. Specifically, we minimize the root mean squared logarithmic error (RMSLE) between Eqn. 11 and the collected data triples, using L-BFGS-B. PolluxAgent then reports the updated values of θsys and φ t to PolluxSched. To ensure that Pollux finds efficient resource allocations through systematic exploration, we impose several priors which bias θsys towards the belief that throughput scales perfectly with more resources, until such resource configurations are explored. Each job starts with a single GPU, and is initially assumed to scale perfectly to more GPUs. PolluxSched is then encouraged to allocate more GPUs and/or nodes to the job, naturally as part of its resource optimization, until the PolluxAgent can estimate θsys more accurately. Finally, to prevent a job from being immediately scaled out to arbitrarily many GPUs, we restrict the maximum number of GPUs which can be allocated to at most twice the maximum number of GPUs the job has been allocated in its lifetime. With θsys, mgns, and m0, which fully specify the DL job’s GOODPUT function at its current training progress, PolluxAgent determines the most efficient batch size m * = argm ax m ⁡ GOODPUT a ,   m , where a is the job’s current resource allocation. Once a new batch size is found, the job will use it for its subsequent training iterations, using AdaScale to adapt its learning rate appropriately. As the job’s statistical efficiency changes over time, PolluxAgent will periodically re-evaluate the most efficient batch size and learning rate. Section 4.2 in Pages 6-7: The PolluxSched periodically allocates (or re-allocates) resources for every job in the cluster. To determine a set of efficient cluster-wide resource allocations, it uses a genetic algorithm to maximize a fitness function which is defined as a weighted mean across speedups for each job in Eqn. 14. A is an allocation matrix with each row Aj being the placement vector for a job j, thus Ajn is the number of GPUs on node n allocated to job j. Intuitively, maximizing FITNESS causes PolluxSched to allocate more resources to jobs that achieve a high SPEEDUP when provided with many GPUs (i.e. jobs that scale well). We define the speedup of each job as the factor of goodput improvement using the given placement vector over using a single process with the optimal batch size in Eqn. 15, where GOODPUTj is the goodput of job j at its current training iteration. Similar to the PolluxAgent, PolluxSched performs the maximizations in the numerator and denominator using golden-section search. Although PolluxSched can work well with all weights wj set to 1, we provide the ability to re-weight jobs as an additional tuning knob for cluster operators. We define the weight wj of job j as Eqn. 16. The weight of a job is 1 if its current total GPU-time is at most GPUTIME_THRES, and decays gradually thereafter. Section 4.2.1 with FIG. 5 in Page 7: Our genetic algorithm operates on a population of distinct allocation matrices (see Fig. 5). During each generation of the algorithm, existing allocation matrices are first randomly mutated, then crossed over to produce offspring allocation matrices, and finally modified to satisfy node resource constraints. A constant population size of distinct allocation matrices is maintained after each generation by discarding the allocation matrices with the lowest objective values (according to Eqn. 14). After several generations of the genetic algorithm, the allocation matrix with the highest fitness score is applied to the jobs running in the cluster. Each time a job is re-allocated to a different set of GPUs, it will need to save a checkpoint of its current model parameters, and restart using its new allocation of GPUs. To prevent an excessive number of re-allocations, when PolluxSched evaluates the fitness function for a given allocation matrix, it applies a penalty for every job that needs to restart. When multiple distributed DL jobs share a single node, their network usage while synchronizing gradients and model parameters may interfere with each other, causing both jobs to slow down. PolluxSched mitigates this issue by disallowing different distributed jobs (each using GPUs across multiple nodes) from sharing the same node. Section 4.2.2 in Pages 7-8: In cloud computing environments, GPU nodes can be dynamically provisioned during periods of high demand to decrease job completion time, as well as released during periods of low demand to decrease cost of cluster resources. Beyond adapting to the number and sizes of DL jobs, co-adaptive scheduling presents a unique opportunity for cluster auto-scaling. Since many distributed DL jobs tend to increase in statistical efficiency as training progresses, it may be more cost-effective to provision more cloud resources during the later iterations of a large training job, rather than earlier on. In order to decide when to request or release nodes in the cloud, we first define a measure of cluster resource utility for a given allocation matrix as Eqn. 17. A cluster operator defines a LOW_UTIL_THRES and a HIGH_UTIL_THRES. PolluxSched will request additional nodes (when UTILITY(A) is high) or release existing nodes (when UTILITY(A) is low) until UTILITY(A) falls within this range, where A is the current allocations applied in the cluster. The cluster operator also defines a MIN_NODES and MAX_NODES, which limit the size of the cluster. PolluxSched performs a binary search for the desired number of nodes, with the assumption that UTILITY decreases with increasing numbers of nodes. At each step of the binary search, PolluxSched runs its genetic algorithm to evaluate the UTILITY of the cluster size being queried. In the end, the cluster size with a UTILITY value closest to (LOW_UTIL_THRES+HIGH_UTIL_THRES)=2 is selected. Section 4.3 in Page 8: PolluxAgent inserts performance profiling code which measures the time taken for each iteration of training, as well as calculating the gradient noise scale. At a fixed time interval, PolluxAgent fits the system throughput model (Eqn. 11) to the profiled metrics collected so far, and reports the fitted system throughput parameters, along with the latest gradient statistics, to PolluxSched. After reporting to PolluxSched, PolluxAgent updates the job’s batch size, by optimizing its now up-to-date goodput function (Eqn. 6) with its currently allocated resources. PolluxSched is implemented as a service in Kubernetes. At a fixed time interval, PolluxSched runs a fixed number of generations of its genetic algorithm. It then applies the best allocation matrix to the cluster, by creating and terminating Kubernetes Pods which run the job replicas. Although allocation matrix with the highest fitness score is applied to the cluster, the entire population is saved and used to bootstrap the genetic algorithm in the next scheduling interval). Claim 28 Qiao in view of Fischer discloses all elements as stated in Claim 27 and further discloses adaptively re-allocating at least one of the parameter servers and the worker nodes in real-time prior to completion of the one or more ML jobs when at least one of the parameter servers and the worker nodes are operating at or above a utilization threshold (Qiao, Abstract in Page 1: Pollux improves scheduling performance in deep learning (DL) clusters by adaptively co-optimizing inter-dependent factors both at the per-job level and at the cluster-wide level. By observing each job during training, Pollux models how their goodput (system throughput combined with statistical efficiency) would change by adding or removing resources. Leveraging these models, Pollux dynamically (re-)assigns resources to maximize cluster-wide goodput, while continually optimizing each DL job to better utilize those resources. Section 3 in Page 4: Using these predictive models, Pollux periodically re-tunes a and m for each job, according to cluster-wide resource availability and performance. Section 4 with FIG. 4 in Pages 5-6: At the cluster-wide granularity, Pollux dynamically re-allocates resources, driven by the goodput of all jobs sharing the cluster. PolluxAgent and PolluxSched co-adapt to each other. While PolluxAgent adapts each training job to make efficient use of its allocated resources, PolluxSched dynamically re-allocates each job’s resources, taking into account the PolluxAgent’s ability to tune its job. Section 4.2 in Pages 6-7: The PolluxSched periodically allocates (or re-allocates) resources for every job in the cluster. Section 4.2.2 in Pages 7-8: In cloud computing environments, GPU nodes can be dynamically provisioned during periods of high demand to decrease job completion time, as well as released during periods of low demand to decrease cost of cluster resources. Beyond adapting to the number and sizes of DL jobs, co-adaptive scheduling presents a unique opportunity for cluster auto-scaling. Since many distributed DL jobs tend to increase in statistical efficiency as training progresses, it may be more cost-effective to provision more cloud resources during the later iterations of a large training job, rather than earlier on. In order to decide when to request or release nodes in the cloud, we first define a measure of cluster resource utility for a given allocation matrix as Eqn. 17. A cluster operator defines a LOW_UTIL_THRES and a HIGH_UTIL_THRES. PolluxSched will request additional nodes (when UTILITY(A) is high) or release existing nodes (when UTILITY(A) is low) until UTILITY(A) falls within this range, where A is the current allocations applied in the cluster. The cluster operator also defines a MIN_NODES and MAX_NODES, which limit the size of the cluster. PolluxSched performs a binary search for the desired number of nodes, with the assumption that UTILITY decreases with increasing numbers of nodes. At each step of the binary search, PolluxSched runs its genetic algorithm to evaluate the UTILITY of the cluster size being queried. In the end, the cluster size with a UTILITY value closest to (LOW_UTIL_THRES+HIGH_UTIL_THRES)=2 is selected). Claim 29 Qiao in view of Fischer discloses all elements as stated in Claim 28 and further discloses saving progress of the one or more ML jobs to identify a checkpoint in response to triggering of the re- allocation of the at least one of the parameter servers and worker nodes, and resuming the one or more ML jobs from the checkpoint (Qiao, Section 4.2.1 in Page 7: Each time a job is re-allocated to a different set of GPUs, it will need to save a checkpoint of its current model parameters, and restart using its new allocation of GPUs. This process can typically take 30s-60s depending on the size of the model being trained. To prevent an excessive number of re-allocations, when PolluxSched evaluates the fitness function for a given allocation matrix, it applies a penalty for every job that needs to restart; i.e., SPEEDUPj(Aj)←SPEEDUPj(Aj) – RESTART_PENALTY). Claim 30 Qiao in view of Fischer discloses all elements as stated in Claim 28 and further discloses wherein the adaptive re-allocation of the at least one of the parameter servers and the worker nodes comprises re-allocating the at least one of the parameter servers and the worker nodes to one or more other nodes with at least one of higher resource demand or utilization than that of a current node to which the parameter servers and the worker nodes are allocated (Qiao, Section 4.2.2 in Pages 7-8: In cloud computing environments, GPU nodes can be dynamically provisioned during periods of high demand to decrease job completion time, as well as released during periods of low demand to decrease cost of cluster resources. Beyond adapting to the number and sizes of DL jobs, co-adaptive scheduling presents a unique opportunity for cluster auto-scaling. Since many distributed DL jobs tend to increase in statistical efficiency as training progresses, it may be more cost-effective to provision more cloud resources during the later iterations of a large training job, rather than earlier on. In order to decide when to request or release nodes in the cloud, we first define a measure of cluster resource utility for a given allocation matrix as Eqn. 17. A cluster operator defines a LOW_UTIL_THRES and a HIGH_UTIL_THRES. PolluxSched will request additional nodes (when UTILITY(A) is high) or release existing nodes (when UTILITY(A) is low) until UTILITY(A) falls within this range, where A is the current allocations applied in the cluster. The cluster operator also defines a MIN_NODES and MAX_NODES, which limit the size of the cluster. PolluxSched performs a binary search for the desired number of nodes, with the assumption that UTILITY decreases with increasing numbers of nodes. At each step of the binary search, PolluxSched runs its genetic algorithm to evaluate the UTILITY of the cluster size being queried. In the end, the cluster size with a UTILITY value closest to (LOW_UTIL_THRES+HIGH_UTIL_THRES)=2 is selected). Claim 31 Qiao in view of Fischer discloses all elements as stated in Claim 27 and further discloses determining whether a stopping criterion has been reached, the stopping criterion comprising a level of improvement achieved in the performance of the one or more ML jobs that is at or below a defined threshold performance difference level (Fischer, Section V.A in Pages 27-29: We set the maximum number of evaluation runs to be 60. To prevent unnecessary evaluation runs for the pla strategies, we stopped the optimizer after measuring zero performance in three consecutive runs.). Claim 32 Qiao in view of Fischer discloses all elements as stated in Claim 27 and further discloses determining whether a stopping criterion has been reached, the stopping criterion comprising a number of executions of the one or more ML jobs that have been performed in accordance with a defined maximum number of DT configurations (Fischer, Section V.A in Pages 27-29: We set the maximum number of evaluation runs to be 60. To prevent unnecessary evaluation runs for the pla strategies, we stopped the optimizer after measuring zero performance in three consecutive runs) Claims 38-40 are rejected under 35 U.S.C. 103 as being unpatentable over Rocha in view of Qiao. Claim 38 Rocha discloses all elements as stated in Claim 37 (see also 112(b) Rejections to Claim 38) and except failing to explicitly disclose evaluating the monitored resource usage metrics against one or more threshold levels of utilization or idleness of resources of the plurality of resource configurations. Qiao teaches a system and a method relating to assign resources for DL jobs (Qiao, Abstract in Page 1), wherein evaluating the monitored resource usage metrics against one or more threshold levels of utilization or idleness of resources of the plurality of resource configurations (Qiao, Section 4.2.2 in Pages 7-8: In cloud computing environments, GPU nodes can be dynamically provisioned during periods of high demand to decrease job completion time, as well as released during periods of low demand to decrease cost of cluster resources. Beyond adapting to the number and sizes of DL jobs, co-adaptive scheduling presents a unique opportunity for cluster auto-scaling. Since many distributed DL jobs tend to increase in statistical efficiency as training progresses, it may be more cost-effective to provision more cloud resources during the later iterations of a large training job, rather than earlier on. In order to decide when to request or release nodes in the cloud, we first define a measure of cluster resource utility for a given allocation matrix as Eqn. 17. A cluster operator defines a LOW_UTIL_THRES and a HIGH_UTIL_THRES. PolluxSched will request additional nodes (when UTILITY(A) is high) or release existing nodes (when UTILITY(A) is low) until UTILITY(A) falls within this range, where A is the current allocations applied in the cluster. The cluster operator also defines a MIN_NODES and MAX_NODES, which limit the size of the cluster. PolluxSched performs a binary search for the desired number of nodes, with the assumption that UTILITY decreases with increasing numbers of nodes. At each step of the binary search, PolluxSched runs its genetic algorithm to evaluate the UTILITY of the cluster size being queried. In the end, the cluster size with a UTILITY value closest to (LOW_UTIL_THRES+HIGH_UTIL_THRES)=2 is selected). Rocha and Qiao are analogous art because they are from the same field of endeavor, a system and a method relating to assign resources for DL jobs. Therefore, it would have been obvious to one of ordinary skill in the art before the effective filling date of the claimed invention to apply the teaching of Qiao to Rocha. Motivation for doing so would allow jobs to use resources more efficiently (Qiao, Abstract and Section 1 in Pages 1-2; Section 3.1 in Pages 4-5; Section 4 in Pages 5-6; and Section 4.2 in Pages 6-8). Claim 39 Rocha discloses all elements as stated in Claim 37 and except failing to explicitly disclose re-allocating one or more resources of the plurality of resource configurations. Qiao teaches a system and a method relating to assign resources for DL jobs (Qiao, Abstract in Page 1), wherein re-allocating one or more resources of the plurality of resource configurations (Qiao, Abstract in Page 1: Pollux improves scheduling performance in deep learning (DL) clusters by adaptively co-optimizing inter-dependent factors both at the per-job level and at the cluster-wide level. By observing each job during training, Pollux models how their goodput (system throughput combined with statistical efficiency) would change by adding or removing resources. Leveraging these models, Pollux dynamically (re-)assigns resources to maximize cluster-wide goodput, while continually optimizing each DL job to better utilize those resources. Section 3 in Page 4: Using these predictive models, Pollux periodically re-tunes a and m for each job, according to cluster-wide resource availability and performance. Section 4 with FIG. 4 in Pages 5-6: At the cluster-wide granularity, Pollux dynamically re-allocates resources, driven by the goodput of all jobs sharing the cluster. PolluxAgent and PolluxSched co-adapt to each other. While PolluxAgent adapts each training job to make efficient use of its allocated resources, PolluxSched dynamically re-allocates each job’s resources, taking into account the PolluxAgent’s ability to tune its job. Section 4.2 in Pages 6-7: The PolluxSched periodically allocates (or re-allocates) resources for every job in the cluster. Section 4.2.2 in Pages 7-8: In cloud computing environments, GPU nodes can be dynamically provisioned during periods of high demand to decrease job completion time, as well as released during periods of low demand to decrease cost of cluster resources. Beyond adapting to the number and sizes of DL jobs, co-adaptive scheduling presents a unique opportunity for cluster auto-scaling. Since many distributed DL jobs tend to increase in statistical efficiency as training progresses, it may be more cost-effective to provision more cloud resources during the later iterations of a large training job, rather than earlier on. In order to decide when to request or release nodes in the cloud, we first define a measure of cluster resource utility for a given allocation matrix as Eqn. 17. A cluster operator defines a LOW_UTIL_THRES and a HIGH_UTIL_THRES. PolluxSched will request additional nodes (when UTILITY(A) is high) or release existing nodes (when UTILITY(A) is low) until UTILITY(A) falls within this range, where A is the current allocations applied in the cluster. The cluster operator also defines a MIN_NODES and MAX_NODES, which limit the size of the cluster. PolluxSched performs a binary search for the desired number of nodes, with the assumption that UTILITY decreases with increasing numbers of nodes. At each step of the binary search, PolluxSched runs its genetic algorithm to evaluate the UTILITY of the cluster size being queried. In the end, the cluster size with a UTILITY value closest to (LOW_UTIL_THRES+HIGH_UTIL_THRES)=2 is selected). Rocha and Qiao are analogous art because they are from the same field of endeavor, a system and a method relating to assign resources for DL jobs. Therefore, it would have been obvious to one of ordinary skill in the art before the effective filling date of the claimed invention to apply the teaching of Qiao to Rocha. Motivation for doing so would allow jobs to use resources more efficiently (Qiao, Abstract and Section 1 in Pages 1-2; Section 3.1 in Pages 4-5; Section 4 in Pages 5-6; and Section 4.2 in Pages 6-8). Claim 40 Rocha in view of Qiao discloses all elements as stated in Claim 39 (see also 112(b) Rejections to Claim 39) and further disclose executing the training jobs in accordance with an updated resource configuration after the re-allocation of the one or more resources (Qiao, Abstract in Page 1: Pollux improves scheduling performance in deep learning (DL) clusters by adaptively co-optimizing inter-dependent factors both at the per-job level and at the cluster-wide level. By observing each job during training, Pollux models how their goodput (system throughput combined with statistical efficiency) would change by adding or removing resources. Leveraging these models, Pollux dynamically (re-)assigns resources to maximize cluster-wide goodput, while continually optimizing each DL job to better utilize those resources. Section 3 in Page 4: Using these predictive models, Pollux periodically re-tunes a and m for each job, according to cluster-wide resource availability and performance. Section 4 with FIG. 4 in Pages 5-6: Pollux performs adaptation at two distinct granularities. First, at a job-level granularity, Pollux dynamically tunes the batch size and learning rate for best utilization of the allocated resources. Second, a the cluster-wide granularity, Pollux dynamically re-allocates resources, driven by the goodput of all jobs sharing the cluster. To achieve this co-adaptivity in a scalable way, Pollux’s design consists of two primary components. First, a PolluxAgent runs together with each job. It measures the gradient noise scale and system throughput for that job, and tunes its batch size and learning rate for efficient utilization of its current allocated resources. PolluxAgent periodically reports the goodput function of its job to the PolluxSched. Second, the PolluxSched periodically optimizes the resource allocations for all jobs in the cluster, taking into account the current statistical efficiency for each job. It uses the goodput function to predict a job’s training performance when allocated different resources. PolluxAgent and PolluxSched co-adapt to each other. While PolluxAgent adapts each training job to make efficient use of its allocated resources, PolluxSched dynamically re-allocates each job’s resources, taking into account the PolluxAgent’s ability to tune its job. Section 4.1 in Page 6: During training, it continually measures the job’s gradient noise scale and system throughput, and reports them to PolluxSched at a fixed interval. It also uses this information to determine the most efficient batch size for its job given its current resource allocations, and adapts its job’s learning rate to this batch size using AdaScale. PolluxAgent measures the time taken per iteration, Titer, and records the triple ( a ,   m , Titer) for all combinations of resource allocations a and batch size m encountered during its lifetime. PolluxAgent determines the most efficient batch size, m * = argm ax m ⁡ GOODPUT a ,   m , where a is the job’s current resource allocation. Once a new batch size is found, the job will use it for its subsequent training iterations, using AdaScale to adapt its learning rate appropriately. Section 4.2 in Pages 6-7: The PolluxSched periodically allocates (or re-allocates) resources for every job in the cluster. Section 4.2.2 in Pages 7-8: In cloud computing environments, GPU nodes can be dynamically provisioned during periods of high demand to decrease job completion time, as well as released during periods of low demand to decrease cost of cluster resources. Beyond adapting to the number and sizes of DL jobs, co-adaptive scheduling presents a unique opportunity for cluster auto-scaling. Since many distributed DL jobs tend to increase in statistical efficiency as training progresses, it may be more cost-effective to provision more cloud resources during the later iterations of a large training job, rather than earlier on. In order to decide when to request or release nodes in the cloud, we first define a measure of cluster resource utility for a given allocation matrix as Eqn. 17. A cluster operator defines a LOW_UTIL_THRES and a HIGH_UTIL_THRES. PolluxSched will request additional nodes (when UTILITY(A) is high) or release existing nodes (when UTILITY(A) is low) until UTILITY(A) falls within this range, where A is the current allocations applied in the cluster. The cluster operator also defines a MIN_NODES and MAX_NODES, which limit the size of the cluster. PolluxSched performs a binary search for the desired number of nodes, with the assumption that UTILITY decreases with increasing numbers of nodes. At each step of the binary search, PolluxSched runs its genetic algorithm to evaluate the UTILITY of the cluster size being queried. In the end, the cluster size with a UTILITY value closest to (LOW_UTIL_THRES+HIGH_UTIL_THRES)=2 is selected). Conclusion The prior art made of record and not relied upon is considered pertinent to applicant's disclosure. Liaw et al. ("HyperSched: Dynamic Resource Reallocation for Model Development on a Deadline", arXiv:2001.02338, Jan 8, 2020, pp. 1-13) discloses in Abstract of Page 1 that Prior research in resource scheduling for machine learning training workloads has largely focused on minimizing job completion times. Commonly, these model training workloads collectively search over a large number of parameter values that control the learning process in a hyperparameter search. It is preferable to identify and maximally provision the best-performing hyperparameter configuration (trial) to achieve the highest accuracy result as soon as possible. To optimally trade-off evaluating multiple configurations and training the most promising ones by a fixed deadline, we design and build HyperSched—a dynamic application-level resource scheduler to track, identify, and preferentially allocate resources to the best performing trials to maximize accuracy by the deadline. HyperSched leverages three properties of a hyperparameter search workload overlooked in prior work – trial disposability, progressively identifiable rankings among different configurations, and space-time constraints – to outperform standard hyperparameter search algorithms across a variety of benchmarks. Liaw further discloses in Section 3 in Pages 3-4 that HyperSched aims to address one objective: to provide the best-trained model by a given deadline. A trial corresponds to an evaluation of a specific hyperparameter configuration. An experiment corresponds to a hyperparameter search composed of multiple trials. In this work, we focus on scheduling the trials for a single experiment over a fixed set of resources and with a known deadline. At the deadline, the system must return a model with a good hyperparameter configuration that has also been adequately trained so that it can be immediately tested or deployed with maximum accuracy. Parallelizable training: Even though HyperSched is able to support a wide range of training workloads, it is targeted towards models that can effectively utilize different quantities of allocated resources and parallelize the training procedure. We assume that trials will return intermediate training progress. Further, model training functions need to support checkpoint-restore functionality, which is commonly offered in many distributed deep learning training frameworks such as TensorFlow and RLlib. To both identify the most promising configuration and optimize/train it as much as possible by the given deadline, an algorithm needs to balance both exploration and exploitation. Liaw also discloses in Section 4 with FIG. 4 in Pages 4-5 that At a high level, HyperSched extends ASHA for constrained resource space-time setting by using an exploration policy that is deadline-aware and by dynamically allocating more resources to fewer promising training trials. HyperSched requires the user to specify a logical unit of resource allocation, or atom, along with a total available count of atoms N. HyperSched reallocates resources in integer quantities of atoms, up to N. The inputs to HyperSched are similar to ASHA but require three more parameters from the user - the time deadline, total resource atoms available, and a scaling function s(a). Given finite resource and time constraints, it is necessary to be deadline-aware for a more effective trade-off of exploration and exploitation. For this, HyperSched utilizes speculative evaluation and a deadline-aware entrance policy. If at any time the trial score at the rung is not in the top 1/η of all seen trials at the rung, the trial will be paused. HyperSched Entrance Policy: HyperSched will stop running new configurations after a certain threshold. Unlike SHA, which aims to address the problem of best arm identification, HyperSched aims to optimize the highest accuracy at deadline, or "best configuration exploitation". This is done via an exploration/exploitation policy that incorporates dynamic resource allocation as illustrated in Fig. 4. Li et al. ("Job Placement Strategy with Opportunistic Resource Sharing for Distributed Deep Learning Clusters", 2020 IEEE 22nd International Conference on High Performance Computing and Communications (HPCC), Dec 14-16, 2020, pp. 620-627) discloses in Abstract of Page 620 that Distributed deep learning frameworks train large deep leaning workload with multiple training jobs on shared distributed GPU servers. There are new challenges when scheduling resources for these systems. Modern deep learning training jobs tend to consume large amount of GPU memory. A training job has an iterative nature that causes the memory usage fluctuate overtime. Jobs sharing a host may suffer from significant performance degradation caused by memory overload in runtime. Moreover, even without memory overloads, deep learning training jobs still experience different levels of performance interference when sharing a GPU device. This paper studies these two issues. We introduced an opportunistic memory sharing model to allocate resources for training jobs with time-varying memory requirements. Based on this model, we introduced an Opportunistic Job Placement Problem (OJPP) for shared GPU clusters that seeks job placement configurations using minimum number of GPU devices and guarantees user-defined performance requirements. We proposed a greedy algorithm and a heuristic algorithm with computational complexities of O(n log n) and O(n2 log n), respectively, to solve the problem. Any inquiry concerning this communication or earlier communications from the examiner should be directed to HWEI-MIN LU whose telephone number is (313)446-4913. The examiner can normally be reached Mon - Fri: 9:00 AM - 6:00 PM EST. Examiner interviews are available via telephone, in-person, and video conferencing using a USPTO supplied web-based collaboration tool. To schedule an interview, applicant is encouraged to use the USPTO Automated Interview Request (AIR) at http://www.uspto.gov/interviewpractice. If attempts to reach the examiner by telephone are unsuccessful, the examiner’s supervisor, Mariela D. Reyes can be reached at (571) 270-1006. The fax phone number for the organization where this application or proceeding is assigned is 571-273-8300. Information regarding the status of published or unpublished applications may be obtained from Patent Center. Unpublished application information in Patent Center is available to registered users. To file and manage patent submissions in Patent Center, visit: https://patentcenter.uspto.gov. Visit https://www.uspto.gov/patents/apply/patent-center for more information about Patent Center and https://www.uspto.gov/patents/docx for information about filing in DOCX format. For additional questions, contact the Electronic Business Center (EBC) at 866-217-9197 (toll-free). If you would like assistance from a USPTO Customer Service Representative, call 800-786-9199 (IN USA OR CANADA) or 571-272-1000. /HWEI-MIN LU/Primary Examiner, Art Unit 2142
Read full office action

Prosecution Timeline

May 03, 2024
Application Filed
Jun 21, 2024
Response after Non-Final Action
Sep 23, 2026
Non-Final Rejection mailed — §102, §103, §112 (current)

Precedent Cases

Applications granted by this same examiner with similar technology

Patent 12749017
MACHINE LEARNING EVALUATION FOR DETECTING FEATURE BIAS
3y 9m to grant Granted Sep 29, 2026
Patent 12737096
DRAWER PAGE OVERLAY FOR MULTITASKING
2y 3m to grant Granted Sep 15, 2026
Patent 12718096
ENHANCED DISCRIMINATE FEATURE LEARNING DEEP RESIDUAL CNN FOR MULTI-TASK ROTATING MACHINERY FAULT DIAGNOSIS WITH INFORMATION FUSION
3y 5m to grant Granted Aug 25, 2026
Patent 12705533
SYSTEMS AND METHODS FOR IMPROVING PREDICTION PROCESS USING AUTOMATED RULE LEARNING FRAMEWORK
3y 8m to grant Granted Aug 11, 2026
Patent 12700003
SYSTEMS AND METHODS FOR FREQUENT MACHINE LEARNING MODEL RETRAINING AND RULE OPTIMIZATION
4y 2m to grant Granted Aug 04, 2026
Study what changed to get past this examiner. Based on 5 most recent grants.

Strategy Recommendation AI-generated — please review before filing

Get a prosecution strategy drawn from examiner precedents, rejection analysis, and claim mapping.
Typically takes 5-10 seconds — AI-generated, attorney review required before filing

Prosecution Projections

1-2
Expected OA Rounds
63%
Grant Probability
99%
With Interview (+40.2%)
2y 11m (~6m remaining)
Median Time to Grant
Low
PTA Risk
Based on 240 resolved cases by this examiner. Grant probability derived from career allowance rate.

Sign in with your work email

Enter your email to receive a magic link. No password needed.

Personal email addresses (Gmail, Yahoo, etc.) are not accepted.

Free tier: 3 strategy analyses per month