{"id":21421,"date":"2024-07-10T13:30:26","date_gmt":"2024-07-10T20:30:26","guid":{"rendered":"https:\/\/engineering.fb.com\/?p=21421"},"modified":"2024-07-10T13:03:30","modified_gmt":"2024-07-10T20:03:30","slug":"tail-utilization-ads-inference-meta","status":"publish","type":"post","link":"https:\/\/engineering.fb.com\/2024\/07\/10\/production-engineering\/tail-utilization-ads-inference-meta\/","title":{"rendered":"Taming the tail utilization of ads inference at Meta scale"},"content":{"rendered":"<ul>\n<li style=\"font-weight: 400;\" aria-level=\"1\"><span style=\"font-weight: 400;\">Tail utilization is a significant system issue and a major factor in overload-related failures and low compute utilization.<\/span><\/li>\n<li style=\"font-weight: 400;\" aria-level=\"1\"><span style=\"font-weight: 400;\">The tail utilization optimizations at Meta have had a profound impact on model serving capacity footprint and reliability.\u00a0<\/span><\/li>\n<li style=\"font-weight: 400;\" aria-level=\"1\"><span style=\"font-weight: 400;\">Failure rates, which are mostly timeout errors, were reduced by two-thirds; the compute footprint delivered 35% more work for the same amount of resources; and p99 latency was cut in half.<\/span><\/li>\n<\/ul>\n<p><span style=\"font-weight: 400;\">The inference platforms that serve the sophisticated machine learning models used by Meta\u2019s ads delivery system require significant infrastructure capacity across CPUs, GPUs, storage, networking, and databases. Improving tail utilization \u2013 the utilization level of the top 5% of the servers when ranked by utilization\u2013 within our infrastructure is imperative to operate our fleet efficiently and sustainably.<\/span><\/p>\n<p><span style=\"font-weight: 400;\">With the growing complexity and computational intensity of these models, as well as the strict latency and throughput requirements to deliver ads, we\u2019ve implemented system optimizations and best practices to address tail utilization. The solutions we\u2019ve implemented for our ads inference service have positively impacted compute utilization in our ads fleet in several ways, including increasing work output by 35 percent without additional resources, decreasing timeout error rates by two-thirds, and reducing tail latency at p99 by half.<\/span><\/p>\n<h2><span style=\"font-weight: 400;\">How Meta\u2019s ads model inference service works\u00a0<\/span><\/h2>\n<p><span style=\"font-weight: 400;\">When placing an ad, client requests are routed to the inference service to get predictions. A single request from a client typically results in multiple model inferences being requested, depending on experiment setup, page type, and ad attributes. This is shown below in figure 1 as a request from the ads core services to the model inference service. The actual request flow is more complex but for the purpose of this post, the below schematic model should serve well.<\/span><\/p>\n<p><span style=\"font-weight: 400;\">The inference service leverages Meta infrastructure capabilities such as <\/span><a href=\"https:\/\/atscaleconference.com\/servicerouter-hyperscale-service-mesh-at-meta\/\"><span style=\"font-weight: 400;\">ServiceRouter<\/span><\/a><span style=\"font-weight: 400;\"> for service discovery, load balancing, and other reliability features. The service is set up as a <\/span><span style=\"font-weight: 400;\">sharded service where each model is a shard and multiple models are hosted in a single host of a job that spans multiple hosts.<\/span><\/p>\n<p><span style=\"font-weight: 400;\">This is supported by Meta\u2019s sharding service, <\/span><a href=\"https:\/\/engineering.fb.com\/2020\/08\/24\/production-engineering\/scaling-services-with-shard-manager\/\"><span style=\"font-weight: 400;\">Shard Manager<\/span><\/a><span style=\"font-weight: 400;\">, a universal infrastructure solution that facilitates efficient development and operation of reliable sharded applications. Meta\u2019s advertising team leverages Shard Manager\u2019sload balancing and shard scaling capabilities to effectively handle shards across heterogeneous hardware.<\/span><\/p>\n<figure id=\"attachment_21422\" aria-describedby=\"caption-attachment-21422\" style=\"width: 600px\" class=\"wp-caption aligncenter\"><img loading=\"lazy\" decoding=\"async\" class=\"wp-image-21422\" src=\"https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-1.png?w=916\" alt=\"\" width=\"600\" height=\"466\" srcset=\"https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-1.png 1496w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-1.png?resize=916,711 916w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-1.png?resize=768,597 768w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-1.png?resize=1024,795 1024w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-1.png?resize=96,75 96w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-1.png?resize=192,149 192w\" sizes=\"auto, (max-width: 992px) 100vw, 62vw\" \/><figcaption id=\"caption-attachment-21422\" class=\"wp-caption-text\">Figure 1: The ads inference architecture.<\/figcaption><\/figure>\n<h2><span style=\"font-weight: 400;\">Challenges of load balancing<\/span><\/h2>\n<p><span style=\"font-weight: 400;\">There are two approaches to load balancing:\u00a0<\/span><\/p>\n<ul>\n<li style=\"font-weight: 400;\" aria-level=\"1\"><span style=\"font-weight: 400;\">Routing load balancing \u2013 load balancing across replicas of a single model. We use ServiceRouter to enable routing based load balancing.\u00a0<\/span><\/li>\n<li style=\"font-weight: 400;\" aria-level=\"1\"><span style=\"font-weight: 400;\">Placement load balancing \u2013 balancing load on hosts by moving replicas of a model across hosts.<\/span><\/li>\n<\/ul>\n<p><span style=\"font-weight: 400;\">Fundamental concepts like replica estimation, snapshot transition and multi-service deployments are key aspects of model productionisation that make load balancing in this environment a complex problem.<\/span><\/p>\n<h3><span style=\"font-weight: 400;\">Replica estimation<\/span><\/h3>\n<p><span style=\"font-weight: 400;\">When a new version of the model enters the system, the number of replicas needed for the new model version is estimated based on historical data of the replica usage of the model.<\/span><\/p>\n<h3><span style=\"font-weight: 400;\">Snapshot transition<\/span><\/h3>\n<p><span style=\"font-weight: 400;\">Ads models are continuously updated to improve their performance. The ads inference system then transitions traffic from the older model to the new version. Updated and refreshed models get a new snapshot ID. Snapshot transition is the mechanism by which the refreshed model replaces the current model serving production traffic.<\/span><\/p>\n<h3><span style=\"font-weight: 400;\">Multi-service deployment<\/span><\/h3>\n<p><span style=\"font-weight: 400;\">Models are deployed to multiple service tiers to take advantage of hardware heterogeneity and <\/span><a href=\"https:\/\/engineering.fb.com\/2020\/09\/14\/networking-traffic\/throughput-autoscaling\/\"><span style=\"font-weight: 400;\">elastic capacity<\/span><\/a><span style=\"font-weight: 400;\">.<\/span><\/p>\n<h2><span style=\"font-weight: 400;\">Why is tail utilization a problem?<\/span><\/h2>\n<p><span style=\"font-weight: 400;\">Tail utilization is a problem because as the number of requests increases, servers that contribute to high tail utilization become overloaded and fail, ultimately affecting our service level agreements (SLAs). Consequently, the extra headroom or buffer needed to handle increased traffic is directly determined by the tail utilization.\u00a0<\/span><\/p>\n<p><span style=\"font-weight: 400;\">This is challenging because it leads to overallocation of capacity for the service. If demand increases, capacity headroom is necessary in constrained servers to maintain service levels when accommodating new demand. Since capacity is uniformly added to all servers in a cluster, generating headroom in constrained servers involves adding significantly more capacity than required for headroom.<\/span><\/p>\n<p><span style=\"font-weight: 400;\">In addition, tail utilization for most constrained servers grows faster than lower percentile utilization due to the non linear relationship between traffic increase and utilization. This is the reason why more capacity is needed even while the system is under utilized on average.<\/span><\/p>\n<p><span style=\"font-weight: 400;\">Making the utilization distribution tighter across the fleet unlocks capacity within servers running at low utilization, i.e. the fleet can support more requests and model launches while maintaining SLAs.<\/span><\/p>\n<figure id=\"attachment_21423\" aria-describedby=\"caption-attachment-21423\" style=\"width: 600px\" class=\"wp-caption aligncenter\"><img loading=\"lazy\" decoding=\"async\" class=\"wp-image-21423\" src=\"https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-2.png?w=916\" alt=\"\" width=\"600\" height=\"342\" srcset=\"https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-2.png 1508w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-2.png?resize=916,522 916w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-2.png?resize=768,438 768w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-2.png?resize=1024,584 1024w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-2.png?resize=96,55 96w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-2.png?resize=192,109 192w\" sizes=\"auto, (max-width: 992px) 100vw, 62vw\" \/><figcaption id=\"caption-attachment-21423\" class=\"wp-caption-text\">Figure 2: Divergence in the tail utilization distribution across percentile ranges.<\/figcaption><\/figure>\n<h2><span style=\"font-weight: 400;\">How we optimized tail utilization\u00a0<\/span><\/h2>\n<p><span style=\"font-weight: 400;\">The implemented solution comprises a class of technical optimizations that attempt to balance the objectives of improving utilization and reducing error rate and latency.<\/span><\/p>\n<p><span style=\"font-weight: 400;\">The improvements made the utilization distribution tighter. This created the ability to move work from crunched servers to low utilization servers and absorb increased demand. As a result, the system has been able to absorb up to 35% load increase with no additional capacity.<\/span><\/p>\n<figure id=\"attachment_21424\" aria-describedby=\"caption-attachment-21424\" style=\"width: 600px\" class=\"wp-caption aligncenter\"><img loading=\"lazy\" decoding=\"async\" class=\"wp-image-21424\" src=\"https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-3.png?w=916\" alt=\"\" width=\"600\" height=\"325\" srcset=\"https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-3.png 1760w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-3.png?resize=916,495 916w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-3.png?resize=768,415 768w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-3.png?resize=1024,554 1024w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-3.png?resize=1536,831 1536w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-3.png?resize=96,52 96w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-3.png?resize=192,104 192w\" sizes=\"auto, (max-width: 992px) 100vw, 62vw\" \/><figcaption id=\"caption-attachment-21424\" class=\"wp-caption-text\">Figure 3: Convergence of tail utilization distribution across percentiles.<\/figcaption><\/figure>\n<p><span style=\"font-weight: 400;\">The reliability also improved, reducing the timeout error rate by two-thirds and cutting latency by half.<\/span><\/p>\n<figure id=\"attachment_21425\" aria-describedby=\"caption-attachment-21425\" style=\"width: 600px\" class=\"wp-caption aligncenter\"><img loading=\"lazy\" decoding=\"async\" class=\"wp-image-21425\" src=\"https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-4.png?w=916\" alt=\"\" width=\"600\" height=\"347\" srcset=\"https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-4.png 1714w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-4.png?resize=916,529 916w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-4.png?resize=768,444 768w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-4.png?resize=1024,591 1024w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-4.png?resize=1536,887 1536w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-4.png?resize=96,55 96w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-4.png?resize=192,111 192w\" sizes=\"auto, (max-width: 992px) 100vw, 62vw\" \/><figcaption id=\"caption-attachment-21425\" class=\"wp-caption-text\">Figure 4: System reliability over time.<\/figcaption><\/figure>\n<p><span style=\"font-weight: 400;\">The solution involved two approaches:\u00a0<\/span><\/p>\n<ol>\n<li><span style=\"font-weight: 400;\">Tuning load balancing mechanisms<\/span><\/li>\n<li><span style=\"font-size: 1rem;\">Making system level changes in model productionisation.\u00a0<\/span><\/li>\n<\/ol>\n<p><span style=\"font-weight: 400;\">The first approach is well understood in the industry. The second one required significant trial, testing, and nuanced execution.<\/span><\/p>\n<h2><span style=\"font-weight: 400;\">Tuning load balancing mechanisms<\/span><\/h2>\n<h3><span style=\"font-weight: 400;\">The power of two choices<\/span><\/h3>\n<p><span style=\"font-weight: 400;\">The service mesh, ServiceRouter, provides detailed instrumentation that allows a better understanding of the load balancing characteristics. Specifically relevant to tail utilization is suboptimal load balancing because of load staleness. To address this we leveraged <\/span><a href=\"https:\/\/www.eecs.harvard.edu\/~michaelm\/postscripts\/mythesis.pdf\"><span style=\"font-weight: 400;\">the power of two choices in a randomized load balancing mechanism<\/span><\/a><span style=\"font-weight: 400;\">. This algorithm requires load data from the servers. This telemetry is collected either by polling \u2013 query server load before request dispatch; or by load-header \u2013\u00a0 piggyback on response.\u00a0<\/span><\/p>\n<p><span style=\"font-weight: 400;\">Polling provides fresh load, while it adds an additional hop, but on the other side, load-header results in reading stale load. Load staleness is a significant issue for large services with substantial clients. Any error here due to staleness would result in random load balancing. For polling, given the inference request is computationally expensive, the overhead was found to be negligible. Using polling improved tail utilization noticeably because heavily loaded hosts were actively avoided. This approach worked very well specifically for inference requests greater than 10s of milliseconds.<\/span><\/p>\n<p><span style=\"font-weight: 400;\">ServiceRouter provides various tuning load-balancing capabilities. We tested many of these techniques, including the number of choices for server selection (i.e., power of k instead of 2), backup request configuration, and <\/span><a href=\"https:\/\/www.usenix.org\/system\/files\/osdi23-saokar.pdf\"><span style=\"font-weight: 400;\">hardware-specific routing weights<\/span><\/a><span style=\"font-weight: 400;\">.\u00a0<\/span><\/p>\n<p><span style=\"font-weight: 400;\">These changes offered marginal improvements. CPU utilization as load-counter was especially insightful. While it is intuitive to balance based on CPU utilization, it turned out to be not useful because: CPU utilization is aggregated over some period of time versus the need for instant load information in this case; and outstanding active tasks waiting on I\/O were not taken into account correctly.<\/span><\/p>\n<h3><span style=\"font-weight: 400;\">Placement load balancing<\/span><\/h3>\n<p><span style=\"font-weight: 400;\">Placement load balancing helped a lot. Given the diversity in model resource demand characteristics and machine resource supply, there is significant variance in server utilization. There is an opportunity to make the utilization distribution tighter by tuning the Shard Manager load balancing configurations, such as load bands, thresholds, and balancing frequency. The basic tuning above helped and provided big gains. It also exposed a deeper problem like spiky tail utilization, which was hidden behind the high tail utilization and was fixed once identified .<\/span><\/p>\n<h2><span style=\"font-weight: 400;\">System level changes<\/span><\/h2>\n<p><span style=\"font-weight: 400;\">There wasn&#8217;t a single significant cause for the utilization variance and several intriguing issues emerged among them that offered valuable insights into the system characteristics.\u00a0<\/span><\/p>\n<h3><span style=\"font-weight: 400;\">Memory bandwidth<\/span><\/h3>\n<p><span style=\"font-weight: 400;\">CPU spikes were observed when new replicas, placed on hosts already hosting other models, began serving traffic. Ideally, this should not happen because Shard Manager should only place a replica when the resource requirements are met. Upon examining the spike pattern, the team discovered that the stall cycles were increasing significantly. Using <\/span><a href=\"https:\/\/developers.facebook.com\/blog\/post\/2022\/11\/16\/dynolog-open-source-system-observability\/\"><span style=\"font-weight: 400;\">dynolog perf instrumentations<\/span><\/a><span style=\"font-weight: 400;\">, we determined that memory latency was increasing as well, which aligned with memory latency benchmarks.<\/span><\/p>\n<p><span style=\"font-weight: 400;\">Memory latency starts to increase exponentially at around 65-70% utilization. It appears to be an increase in CPU utilization, but the actual issue was that the CPU was stalling. The solution involved considering memory bandwidth as a resource during replica placement in Shard Manager.\u00a0<\/span><\/p>\n<h3><span style=\"font-weight: 400;\">ServiceRouter and Shard Manager expectation mismatch<\/span><\/h3>\n<p><span style=\"font-weight: 400;\">There is a service control plane component called ReplicaEstimator that performs replica count estimation for a model. When ReplicaEstimator performs this estimation, the expectation is that each replica roughly receives the same amount of traffic. Shard Manager also works under this assumption that replicas of the same model will roughly be equal in their resource usage on a host. Shard Manager load balancing also assumes this property. There are also cases where Shard Manager uses load information from other replicas if load fetch fails. So ReplicaEstimator and Shard Manager share the same expectation that each replica will end up doing roughly the same amount of work.<\/span><\/p>\n<p><span style=\"font-weight: 400;\">ServiceRouter employs the default load counter, which encompasses both active and queued outstanding requests on a host. In general, this works fine when there is only one replica per host and they are expected to receive the same amount of load. However, this assumption is broken due to multi-tenancy, resulting in each host potentially having different models and outstanding requests on a host cannot be used to compare load as it can vary greatly. For example, two hosts serving the same model could have completely different load metrics leading to significant CPU imbalance issues.<\/span><\/p>\n<p><span style=\"font-weight: 400;\">The imbalance of replica load created because of the host level consolidated load counter violates Shard Manager and ReplicaEstimator expectations. A simple and elegant solution to this problem is <\/span><i><span style=\"font-weight: 400;\">a per-model load counter<\/span><\/i><span style=\"font-weight: 400;\">. If each model were to expose a load counter based on its own load on the server, ServiceRouter will end up balancing load across model replicas, and Shard Manger will end up more accurately balancing hosts. Replica estimation also ends up being more accurate. All expectations are aligned.\u00a0<\/span><\/p>\n<p><span style=\"font-weight: 400;\">Support for this was added to the prediction client by explicitly setting the load counter per model client and exposing appropriate per model load metric on the server side. The model replica load distribution as expected became much tighter with a per-model load counter and helps with the problems discussed above.<\/span><\/p>\n<p><span style=\"font-weight: 400;\">But this also presented some challenges. Enabling per-model load counter changes the load distribution instantaneously, causing spikes until Shard Manager catches up and rebalances. The team built a mechanism to make the transition smooth by gradually rolling out the load counter change to the client. Then there are models with low load that end up having per-model load counter values of \u20180\u2019, making it essentially random. In the default load counter configuration, such models end up using the host level load as a good proxy to decide which server to send the request to.<\/span><\/p>\n<p><span style=\"font-weight: 400;\">&#8220;Outstanding examples CPU\u201d was the most promising load counter among many that were tested. It is the estimated total CPU time spent on active requests, and better represents the cost of outstanding work. The counter is normalized by the number of cores to account for machine heterogeneity.<\/span><\/p>\n<figure id=\"attachment_21426\" aria-describedby=\"caption-attachment-21426\" style=\"width: 600px\" class=\"wp-caption aligncenter\"><img loading=\"lazy\" decoding=\"async\" class=\"wp-image-21426\" src=\"https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-5.png?w=1024\" alt=\"\" width=\"600\" height=\"288\" srcset=\"https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-5.png 1712w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-5.png?resize=916,440 916w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-5.png?resize=768,369 768w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-5.png?resize=1024,492 1024w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-5.png?resize=1536,737 1536w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-5.png?resize=96,46 96w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-5.png?resize=192,92 192w\" sizes=\"auto, (max-width: 992px) 100vw, 62vw\" \/><figcaption id=\"caption-attachment-21426\" class=\"wp-caption-text\">Figure 5: Throughput as measured by requests per second across hosts in a tier.<\/figcaption><\/figure>\n<h3><span style=\"font-weight: 400;\">Snapshot transition<\/span><\/h3>\n<p><span style=\"font-weight: 400;\">Some ads models are retrained more frequently than others. Discounting real-time updated models, the majority of the models involve transitioning traffic from a previous model snapshot to the new model snapshot. Snapshot transition is a major disruption to a balanced system, especially when the transitioning models have a large number of replicas.<\/span><\/p>\n<p><span style=\"font-weight: 400;\">During peak traffic, snapshot transition can have a significant impact on utilization. Figure 6 below illustrates the issue. The snapshot transition of large models during a crunched time causes utilization to be very unbalanced until Shard Manager is able to bring it back in balance. This takes a few load balancing runs because the placement of the new model during peak traffic ends up violating CPU soft thresholds. The problem of load counters, as discussed earlier, further complicates Shard Manager&#8217;s ability to resolve issues.<\/span><\/p>\n<figure id=\"attachment_21427\" aria-describedby=\"caption-attachment-21427\" style=\"width: 600px\" class=\"wp-caption aligncenter\"><img loading=\"lazy\" decoding=\"async\" class=\"wp-image-21427\" src=\"https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-6.png?w=1024\" alt=\"\" width=\"600\" height=\"353\" srcset=\"https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-6.png 1756w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-6.png?resize=916,539 916w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-6.png?resize=768,452 768w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-6.png?resize=1024,603 1024w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-6.png?resize=1536,904 1536w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-6.png?resize=96,57 96w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-6.png?resize=192,113 192w\" sizes=\"auto, (max-width: 992px) 100vw, 62vw\" \/><figcaption id=\"caption-attachment-21427\" class=\"wp-caption-text\">Figure 6: A utilization spike due to the snapshot transition.<\/figcaption><\/figure>\n<p><span style=\"font-weight: 400;\">To mitigate this issue, the team added the snapshot transition budget capability. This allows for snapshot transitions to occur only when resource utilization is below a configured threshold. The trade-off here is between snapshot staleness and failure rate. Fast scale down of old snapshots helped minimize the overhead of snapshot staleness while maintaining lower failure rates.<\/span><\/p>\n<h3><span style=\"font-weight: 400;\">Cross-service load balancing<\/span><\/h3>\n<p><span style=\"font-weight: 400;\">After optimizing load balancing within a single service, the next step was to extend this to multiple services. Each regional model inference service is made up of multiple sub-services depending on hardware type and capacity pools \u2013 <\/span><a href=\"https:\/\/engineering.fb.com\/2020\/09\/14\/networking-traffic\/throughput-autoscaling\/\"><span style=\"font-weight: 400;\">guaranteed and elastic pools<\/span><\/a><span style=\"font-weight: 400;\">. We changed the calculation to the compute capacity of the hosts instead of the host number. This helped with a more balanced load across tiers.\u00a0<\/span><\/p>\n<p><span style=\"font-weight: 400;\">Certain hardware types are more loaded than others. Given that clients maintain separate connections to these tiers, ServiceRouter load balancing, which performs balancing within tiers, did not help. Given the production setup, it was non-trivial to put all these tiers behind a single parent tier. Therefore, the team added a small utilization balancing feedback controller to adjust traffic routing percentages and achieve balance between these tiers. Figure 7 shows\u00a0 an example of this being rolled out.<\/span><\/p>\n<figure id=\"attachment_21428\" aria-describedby=\"caption-attachment-21428\" style=\"width: 600px\" class=\"wp-caption aligncenter\"><img loading=\"lazy\" decoding=\"async\" class=\"wp-image-21428\" src=\"https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-7.png?w=1024\" alt=\"\" width=\"600\" height=\"329\" srcset=\"https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-7.png 1736w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-7.png?resize=916,502 916w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-7.png?resize=768,421 768w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-7.png?resize=1024,562 1024w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-7.png?resize=1536,842 1536w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-7.png?resize=96,53 96w, https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/tail-utilization_figure-7.png?resize=192,105 192w\" sizes=\"auto, (max-width: 992px) 100vw, 62vw\" \/><figcaption id=\"caption-attachment-21428\" class=\"wp-caption-text\">Figure 7: Request per service.<\/figcaption><\/figure>\n<h3><span style=\"font-weight: 400;\">Replica estimation and predictive scaling\u00a0<\/span><\/h3>\n<p><span style=\"font-weight: 400;\">Shard Manager employs a reactive approach to load by scaling up replicas in response to a load increase. This meant increased error rates during the time replicas were scaled up and became ready. <\/span><span style=\"font-weight: 400;\">This is exacerbated by the fact that replicas with higher utilization are more prone to utilization spikes given the non-linear relationship between queries per second (QPS) and utilization. <\/span><span style=\"font-weight: 400;\">To add to this, when auto-scaling kicks in, it responds to a much larger CPU requirement and results in over-replication. We designed a simple predictive replica estimation system for the models that predicts future resource usage based on current and past usage patterns up to two hours in advance. This approach yielded significant improvements in failure rate during peak periods.\u00a0<\/span><\/p>\n<h2><span style=\"font-weight: 400;\">Next steps\u00a0<\/span><\/h2>\n<p><span style=\"font-weight: 400;\">The next step in our journey is to adopt our learnings around tail utilization to new system architectures and platforms. For example, we\u2019re actively working to apply the utilizations discussed here to <\/span><a href=\"https:\/\/atscaleconference.com\/ipnext-metas-next-generation-inference-platform\/\"><span style=\"font-weight: 400;\">IPnext<\/span><\/a><span style=\"font-weight: 400;\">, Meta\u2019s next-generation unified platform for managing the entire lifecycle of machine learning model deployments, from publishing to serving. IPnext&#8217;s modular design enables us to support various model architectures (e.g., for ranking or GenAI applications) through a single platform spanning multiple data center regions. Optimizing tail utilization within IPnext thereby delivering these benefits to a broader range of expanding machine learning inference use cases at Meta.<\/span><\/p>\n","protected":false},"excerpt":{"rendered":"<p>Tail utilization is a significant system issue and a major factor in overload-related failures and low compute utilization. The tail utilization optimizations at Meta have had a profound impact on model serving capacity footprint and reliability.\u00a0 Failure rates, which are mostly timeout errors, were reduced by two-thirds; the compute footprint delivered 35% more work for [&#8230;]<\/p>\n<p><a class=\"btn btn-secondary understrap-read-more-link\" href=\"https:\/\/engineering.fb.com\/2024\/07\/10\/production-engineering\/tail-utilization-ads-inference-meta\/\">Read More&#8230;<\/a><\/p>\n","protected":false},"author":51,"featured_media":21449,"comment_status":"closed","ping_status":"closed","sticky":false,"template":"","format":"standard","meta":{"_jetpack_newsletter_access":"","_jetpack_dont_email_post_to_subs":false,"_jetpack_newsletter_tier_id":0,"_jetpack_memberships_contains_paywalled_content":false,"_jetpack_feature_clip_id":0,"_jetpack_memberships_contains_paid_content":false,"footnotes":"","jetpack_post_was_ever_published":false},"categories":[72,70,67],"tags":[],"coauthors":[2100,2099,2101],"class_list":["post-21421","post","type-post","status-publish","format-standard","has-post-thumbnail","hentry","category-ml-applications","category-networking-traffic","category-production-engineering","fb_content_type-article"],"yoast_head":"<!-- This site is optimized with the Yoast SEO Premium plugin v19.3 (Yoast SEO v27.3) - https:\/\/yoast.com\/product\/yoast-seo-premium-wordpress\/ -->\n<title>Taming the tail utilization of ads inference at Meta scale - Engineering at Meta<\/title>\n<meta name=\"robots\" content=\"index, follow, max-snippet:-1, max-image-preview:large, max-video-preview:-1\" \/>\n<link rel=\"canonical\" href=\"https:\/\/engineering.fb.com\/2024\/07\/10\/production-engineering\/tail-utilization-ads-inference-meta\/\" \/>\n<meta name=\"twitter:label1\" content=\"Written by\" \/>\n\t<meta name=\"twitter:data1\" content=\"Rohith Menon, Bikash Sharma, Deepak Tiwari\" \/>\n\t<meta name=\"twitter:label2\" content=\"Est. reading time\" \/>\n\t<meta name=\"twitter:data2\" content=\"13 minutes\" \/>\n<script type=\"application\/ld+json\" class=\"yoast-schema-graph\">{\"@context\":\"https:\\\/\\\/schema.org\",\"@graph\":[{\"@type\":\"Article\",\"@id\":\"https:\\\/\\\/engineering.fb.com\\\/2024\\\/07\\\/10\\\/production-engineering\\\/tail-utilization-ads-inference-meta\\\/#article\",\"isPartOf\":{\"@id\":\"https:\\\/\\\/engineering.fb.com\\\/2024\\\/07\\\/10\\\/production-engineering\\\/tail-utilization-ads-inference-meta\\\/\"},\"author\":{\"@id\":\"https:\\\/\\\/engineering.fb.com\\\/2024\\\/07\\\/10\\\/production-engineering\\\/tail-utilization-ads-inference-meta\\\/#author\",\"name\":\"\"},\"headline\":\"Taming the tail utilization of ads inference at Meta scale\",\"datePublished\":\"2024-07-10T20:30:26+00:00\",\"mainEntityOfPage\":{\"@id\":\"https:\\\/\\\/engineering.fb.com\\\/2024\\\/07\\\/10\\\/production-engineering\\\/tail-utilization-ads-inference-meta\\\/\"},\"wordCount\":2520,\"publisher\":{\"@id\":\"https:\\\/\\\/engineering.fb.com\\\/#organization\"},\"image\":{\"@id\":\"https:\\\/\\\/engineering.fb.com\\\/2024\\\/07\\\/10\\\/production-engineering\\\/tail-utilization-ads-inference-meta\\\/#primaryimage\"},\"thumbnailUrl\":\"https:\\\/\\\/engineering.fb.com\\\/wp-content\\\/uploads\\\/2024\\\/07\\\/Meta-Data-Center-Cold-Storage.jpg\",\"articleSection\":[\"ML Applications\",\"Networking &amp; Traffic\",\"Production Engineering\"],\"inLanguage\":\"en-US\"},{\"@type\":\"WebPage\",\"@id\":\"https:\\\/\\\/engineering.fb.com\\\/2024\\\/07\\\/10\\\/production-engineering\\\/tail-utilization-ads-inference-meta\\\/\",\"url\":\"https:\\\/\\\/engineering.fb.com\\\/2024\\\/07\\\/10\\\/production-engineering\\\/tail-utilization-ads-inference-meta\\\/\",\"name\":\"Taming the tail utilization of ads inference at Meta scale - Engineering at Meta\",\"isPartOf\":{\"@id\":\"https:\\\/\\\/engineering.fb.com\\\/#website\"},\"primaryImageOfPage\":{\"@id\":\"https:\\\/\\\/engineering.fb.com\\\/2024\\\/07\\\/10\\\/production-engineering\\\/tail-utilization-ads-inference-meta\\\/#primaryimage\"},\"image\":{\"@id\":\"https:\\\/\\\/engineering.fb.com\\\/2024\\\/07\\\/10\\\/production-engineering\\\/tail-utilization-ads-inference-meta\\\/#primaryimage\"},\"thumbnailUrl\":\"https:\\\/\\\/engineering.fb.com\\\/wp-content\\\/uploads\\\/2024\\\/07\\\/Meta-Data-Center-Cold-Storage.jpg\",\"datePublished\":\"2024-07-10T20:30:26+00:00\",\"breadcrumb\":{\"@id\":\"https:\\\/\\\/engineering.fb.com\\\/2024\\\/07\\\/10\\\/production-engineering\\\/tail-utilization-ads-inference-meta\\\/#breadcrumb\"},\"inLanguage\":\"en-US\",\"potentialAction\":[{\"@type\":\"ReadAction\",\"target\":[\"https:\\\/\\\/engineering.fb.com\\\/2024\\\/07\\\/10\\\/production-engineering\\\/tail-utilization-ads-inference-meta\\\/\"]}]},{\"@type\":\"ImageObject\",\"inLanguage\":\"en-US\",\"@id\":\"https:\\\/\\\/engineering.fb.com\\\/2024\\\/07\\\/10\\\/production-engineering\\\/tail-utilization-ads-inference-meta\\\/#primaryimage\",\"url\":\"https:\\\/\\\/engineering.fb.com\\\/wp-content\\\/uploads\\\/2024\\\/07\\\/Meta-Data-Center-Cold-Storage.jpg\",\"contentUrl\":\"https:\\\/\\\/engineering.fb.com\\\/wp-content\\\/uploads\\\/2024\\\/07\\\/Meta-Data-Center-Cold-Storage.jpg\",\"width\":2576,\"height\":1719},{\"@type\":\"BreadcrumbList\",\"@id\":\"https:\\\/\\\/engineering.fb.com\\\/2024\\\/07\\\/10\\\/production-engineering\\\/tail-utilization-ads-inference-meta\\\/#breadcrumb\",\"itemListElement\":[{\"@type\":\"ListItem\",\"position\":1,\"name\":\"Home\",\"item\":\"https:\\\/\\\/engineering.fb.com\\\/\"},{\"@type\":\"ListItem\",\"position\":2,\"name\":\"Taming the tail utilization of ads inference at Meta scale\"}]},{\"@type\":\"WebSite\",\"@id\":\"https:\\\/\\\/engineering.fb.com\\\/#website\",\"url\":\"https:\\\/\\\/engineering.fb.com\\\/\",\"name\":\"Engineering at Meta\",\"description\":\"Engineering at Meta Blog\",\"publisher\":{\"@id\":\"https:\\\/\\\/engineering.fb.com\\\/#organization\"},\"potentialAction\":[{\"@type\":\"SearchAction\",\"target\":{\"@type\":\"EntryPoint\",\"urlTemplate\":\"https:\\\/\\\/engineering.fb.com\\\/?s={search_term_string}\"},\"query-input\":{\"@type\":\"PropertyValueSpecification\",\"valueRequired\":true,\"valueName\":\"search_term_string\"}}],\"inLanguage\":\"en-US\"},{\"@type\":\"Organization\",\"@id\":\"https:\\\/\\\/engineering.fb.com\\\/#organization\",\"name\":\"Meta\",\"url\":\"https:\\\/\\\/engineering.fb.com\\\/\",\"logo\":{\"@type\":\"ImageObject\",\"inLanguage\":\"en-US\",\"@id\":\"https:\\\/\\\/engineering.fb.com\\\/#\\\/schema\\\/logo\\\/image\\\/\",\"url\":\"https:\\\/\\\/engineering.fb.com\\\/wp-content\\\/uploads\\\/2023\\\/08\\\/Meta_lockup_positive-primary_RGB.jpg\",\"contentUrl\":\"https:\\\/\\\/engineering.fb.com\\\/wp-content\\\/uploads\\\/2023\\\/08\\\/Meta_lockup_positive-primary_RGB.jpg\",\"width\":29011,\"height\":12501,\"caption\":\"Meta\"},\"image\":{\"@id\":\"https:\\\/\\\/engineering.fb.com\\\/#\\\/schema\\\/logo\\\/image\\\/\"},\"sameAs\":[\"https:\\\/\\\/www.facebook.com\\\/Engineering\\\/\",\"https:\\\/\\\/x.com\\\/fb_engineering\"]},[]]}<\/script>\n<!-- \/ Yoast SEO Premium plugin. -->","yoast_head_json":{"title":"Taming the tail utilization of ads inference at Meta scale - Engineering at Meta","robots":{"index":"index","follow":"follow","max-snippet":"max-snippet:-1","max-image-preview":"max-image-preview:large","max-video-preview":"max-video-preview:-1"},"canonical":"https:\/\/engineering.fb.com\/2024\/07\/10\/production-engineering\/tail-utilization-ads-inference-meta\/","twitter_misc":{"Written by":"Rohith Menon, Bikash Sharma, Deepak Tiwari","Est. reading time":"13 minutes"},"schema":{"@context":"https:\/\/schema.org","@graph":[{"@type":"Article","@id":"https:\/\/engineering.fb.com\/2024\/07\/10\/production-engineering\/tail-utilization-ads-inference-meta\/#article","isPartOf":{"@id":"https:\/\/engineering.fb.com\/2024\/07\/10\/production-engineering\/tail-utilization-ads-inference-meta\/"},"author":{"@id":"https:\/\/engineering.fb.com\/2024\/07\/10\/production-engineering\/tail-utilization-ads-inference-meta\/#author","name":""},"headline":"Taming the tail utilization of ads inference at Meta scale","datePublished":"2024-07-10T20:30:26+00:00","mainEntityOfPage":{"@id":"https:\/\/engineering.fb.com\/2024\/07\/10\/production-engineering\/tail-utilization-ads-inference-meta\/"},"wordCount":2520,"publisher":{"@id":"https:\/\/engineering.fb.com\/#organization"},"image":{"@id":"https:\/\/engineering.fb.com\/2024\/07\/10\/production-engineering\/tail-utilization-ads-inference-meta\/#primaryimage"},"thumbnailUrl":"https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/Meta-Data-Center-Cold-Storage.jpg","articleSection":["ML Applications","Networking &amp; Traffic","Production Engineering"],"inLanguage":"en-US"},{"@type":"WebPage","@id":"https:\/\/engineering.fb.com\/2024\/07\/10\/production-engineering\/tail-utilization-ads-inference-meta\/","url":"https:\/\/engineering.fb.com\/2024\/07\/10\/production-engineering\/tail-utilization-ads-inference-meta\/","name":"Taming the tail utilization of ads inference at Meta scale - Engineering at Meta","isPartOf":{"@id":"https:\/\/engineering.fb.com\/#website"},"primaryImageOfPage":{"@id":"https:\/\/engineering.fb.com\/2024\/07\/10\/production-engineering\/tail-utilization-ads-inference-meta\/#primaryimage"},"image":{"@id":"https:\/\/engineering.fb.com\/2024\/07\/10\/production-engineering\/tail-utilization-ads-inference-meta\/#primaryimage"},"thumbnailUrl":"https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/Meta-Data-Center-Cold-Storage.jpg","datePublished":"2024-07-10T20:30:26+00:00","breadcrumb":{"@id":"https:\/\/engineering.fb.com\/2024\/07\/10\/production-engineering\/tail-utilization-ads-inference-meta\/#breadcrumb"},"inLanguage":"en-US","potentialAction":[{"@type":"ReadAction","target":["https:\/\/engineering.fb.com\/2024\/07\/10\/production-engineering\/tail-utilization-ads-inference-meta\/"]}]},{"@type":"ImageObject","inLanguage":"en-US","@id":"https:\/\/engineering.fb.com\/2024\/07\/10\/production-engineering\/tail-utilization-ads-inference-meta\/#primaryimage","url":"https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/Meta-Data-Center-Cold-Storage.jpg","contentUrl":"https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/Meta-Data-Center-Cold-Storage.jpg","width":2576,"height":1719},{"@type":"BreadcrumbList","@id":"https:\/\/engineering.fb.com\/2024\/07\/10\/production-engineering\/tail-utilization-ads-inference-meta\/#breadcrumb","itemListElement":[{"@type":"ListItem","position":1,"name":"Home","item":"https:\/\/engineering.fb.com\/"},{"@type":"ListItem","position":2,"name":"Taming the tail utilization of ads inference at Meta scale"}]},{"@type":"WebSite","@id":"https:\/\/engineering.fb.com\/#website","url":"https:\/\/engineering.fb.com\/","name":"Engineering at Meta","description":"Engineering at Meta Blog","publisher":{"@id":"https:\/\/engineering.fb.com\/#organization"},"potentialAction":[{"@type":"SearchAction","target":{"@type":"EntryPoint","urlTemplate":"https:\/\/engineering.fb.com\/?s={search_term_string}"},"query-input":{"@type":"PropertyValueSpecification","valueRequired":true,"valueName":"search_term_string"}}],"inLanguage":"en-US"},{"@type":"Organization","@id":"https:\/\/engineering.fb.com\/#organization","name":"Meta","url":"https:\/\/engineering.fb.com\/","logo":{"@type":"ImageObject","inLanguage":"en-US","@id":"https:\/\/engineering.fb.com\/#\/schema\/logo\/image\/","url":"https:\/\/engineering.fb.com\/wp-content\/uploads\/2023\/08\/Meta_lockup_positive-primary_RGB.jpg","contentUrl":"https:\/\/engineering.fb.com\/wp-content\/uploads\/2023\/08\/Meta_lockup_positive-primary_RGB.jpg","width":29011,"height":12501,"caption":"Meta"},"image":{"@id":"https:\/\/engineering.fb.com\/#\/schema\/logo\/image\/"},"sameAs":["https:\/\/www.facebook.com\/Engineering\/","https:\/\/x.com\/fb_engineering"]},[]]}},"jetpack_featured_media_url":"https:\/\/engineering.fb.com\/wp-content\/uploads\/2024\/07\/Meta-Data-Center-Cold-Storage.jpg","jetpack_shortlink":"https:\/\/wp.me\/pa0Lhq-5zv","jetpack_sharing_enabled":true,"_links":{"self":[{"href":"https:\/\/engineering.fb.com\/wp-json\/wp\/v2\/posts\/21421","targetHints":{"allow":["GET"]}}],"collection":[{"href":"https:\/\/engineering.fb.com\/wp-json\/wp\/v2\/posts"}],"about":[{"href":"https:\/\/engineering.fb.com\/wp-json\/wp\/v2\/types\/post"}],"author":[{"embeddable":true,"href":"https:\/\/engineering.fb.com\/wp-json\/wp\/v2\/users\/51"}],"replies":[{"embeddable":true,"href":"https:\/\/engineering.fb.com\/wp-json\/wp\/v2\/comments?post=21421"}],"version-history":[{"count":16,"href":"https:\/\/engineering.fb.com\/wp-json\/wp\/v2\/posts\/21421\/revisions"}],"predecessor-version":[{"id":21473,"href":"https:\/\/engineering.fb.com\/wp-json\/wp\/v2\/posts\/21421\/revisions\/21473"}],"wp:featuredmedia":[{"embeddable":true,"href":"https:\/\/engineering.fb.com\/wp-json\/wp\/v2\/media\/21449"}],"wp:attachment":[{"href":"https:\/\/engineering.fb.com\/wp-json\/wp\/v2\/media?parent=21421"}],"wp:term":[{"taxonomy":"category","embeddable":true,"href":"https:\/\/engineering.fb.com\/wp-json\/wp\/v2\/categories?post=21421"},{"taxonomy":"post_tag","embeddable":true,"href":"https:\/\/engineering.fb.com\/wp-json\/wp\/v2\/tags?post=21421"},{"taxonomy":"author","embeddable":true,"href":"https:\/\/engineering.fb.com\/wp-json\/wp\/v2\/coauthors?post=21421"}],"curies":[{"name":"wp","href":"https:\/\/api.w.org\/{rel}","templated":true}]}}