Files
nexus/sreweekly/articles/361/05-consistent-hashing-algorithm.html
2026-09-12 17:23:01 +08:00

562 lines
70 KiB
HTML
Raw Permalink Blame History

This file contains invisible Unicode characters
This file contains invisible Unicode characters that are indistinguishable to humans but may be processed differently by a computer. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
<!DOCTYPE html>
<html lang="en">
<head>
<title>Consistent hashing algorithm - High Scalability -</title>
<meta charset="utf-8">
<meta name="viewport" content="width=device-width, initial-scale=1.0">
<link rel="preload" as="style" href="https://highscalability.com/assets/built/screen.css?v=QpNBjlSbNzyRyl6s">
<link rel="preload" as="script" href="https://highscalability.com/assets/built/source.js?v=uMULaSMoSjM4gY5C">
<link rel="preload" as="font" type="font/woff2" href="https://highscalability.com/assets/fonts/inter-roman.woff2?v=OecsB5TBLy27FKD2" crossorigin="anonymous">
<style>
@font-face {
font-family: "Inter";
font-style: normal;
font-weight: 100 900;
font-display: optional;
src: url(https://highscalability.com/assets/fonts/inter-roman.woff2?v=OecsB5TBLy27FKD2) format("woff2");
unicode-range: U+0000-00FF, U+0131, U+0152-0153, U+02BB-02BC, U+02C6, U+02DA, U+02DC, U+0304, U+0308, U+0329, U+2000-206F, U+2074, U+20AC, U+2122, U+2191, U+2193, U+2212, U+2215, U+FEFF, U+FFFD;
}
</style>
<link rel="stylesheet" type="text/css" href="https://highscalability.com/assets/built/screen.css?v=QpNBjlSbNzyRyl6s">
<style>
:root {
--background-color: #ffffff
}
</style>
<script>
/* The script for calculating the color contrast has been taken from
https://gomakethings.com/dynamically-changing-the-text-color-based-on-background-color-contrast-with-vanilla-js/ */
var accentColor = getComputedStyle(document.documentElement).getPropertyValue('--background-color');
accentColor = accentColor.trim().slice(1);
if (accentColor.length === 3) {
accentColor = accentColor[0] + accentColor[0] + accentColor[1] + accentColor[1] + accentColor[2] + accentColor[2];
}
var r = parseInt(accentColor.substr(0, 2), 16);
var g = parseInt(accentColor.substr(2, 2), 16);
var b = parseInt(accentColor.substr(4, 2), 16);
var yiq = ((r * 299) + (g * 587) + (b * 114)) / 1000;
var textColor = (yiq >= 128) ? 'dark' : 'light';
document.documentElement.className = `has-${textColor}-text`;
</script>
<link rel="canonical" href="https://highscalability.com/consistent-hashing-algorithm/">
<meta name="referrer" content="no-referrer-when-downgrade">
<meta property="og:site_name" content="High Scalability">
<meta property="og:type" content="article">
<meta property="og:title" content="Consistent hashing algorithm - High Scalability -">
<meta property="og:description" content="You can subscribe to the system design newsletter to excel in system design interviews a...">
<meta property="og:url" content="https://highscalability.com/consistent-hashing-algorithm/">
<meta property="og:image" content="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/52705024790_8a89921cfb-1.jpg">
<meta property="article:published_time" content="2023-02-22T16:39:15.000Z">
<meta property="article:modified_time" content="2024-02-02T12:32:33.000Z">
<meta property="article:tag" content="consistent hashing">
<meta property="article:tag" content="sharding">
<meta property="article:tag" content="Load Balancing">
<meta property="article:publisher" content="https://www.facebook.com/ghost">
<meta name="twitter:card" content="summary_large_image">
<meta name="twitter:title" content="Consistent hashing algorithm - High Scalability -">
<meta name="twitter:description" content="You can subscribe to the system design newsletter to excel in system design interviews and software architecture.
You can view the original article Consistent hashing explained on systemdesign.one website.
How does consistent hashing work?
At a high level, consistent hashing performs the following operations:
1. The output of the">
<meta name="twitter:url" content="https://highscalability.com/consistent-hashing-algorithm/">
<meta name="twitter:image" content="https://static.ghost.org/v5.0.0/images/publication-cover.jpg">
<meta name="twitter:label1" content="Written by">
<meta name="twitter:data1" content="NK">
<meta name="twitter:label2" content="Filed under">
<meta name="twitter:data2" content="consistent hashing, sharding, Load Balancing">
<meta name="twitter:site" content="@ghost">
<meta property="og:image:width" content="500">
<meta property="og:image:height" content="333">
<script type="application/ld+json">
{
"@context": "https://schema.org",
"@type": "Article",
"publisher": {
"@type": "Organization",
"name": "High Scalability",
"url": "https://highscalability.com/",
"logo": {
"@type": "ImageObject",
"url": "https://highscalability.com/favicon.ico",
"width": 48,
"height": 48
}
},
"author": {
"@type": "Person",
"name": "NK",
"url": "https://highscalability.com/author/sdxyz/",
"sameAs": []
},
"headline": "Consistent hashing algorithm - High Scalability -",
"url": "https://highscalability.com/consistent-hashing-algorithm/",
"datePublished": "2023-02-22T16:39:15.000Z",
"dateModified": "2024-02-02T12:32:33.000Z",
"keywords": "consistent hashing, sharding, Load Balancing",
"description": "You can subscribe to the system design newsletter to excel in system design interviews and software architecture.\n\nYou can view the original article Consistent hashing explained on systemdesign.one website.\n\n\nHow does consistent hashing work?\n\nAt a high level, consistent hashing performs the following operations:\n\n 1. The output of the hash function is placed on a virtual ring structure (known as the hash ring)\n 2. The hashed IP addresses of the nodes are used to assign a position for the nodes ",
"mainEntityOfPage": "https://highscalability.com/consistent-hashing-algorithm/"
}
</script>
<meta name="generator" content="Ghost 6.64">
<link rel="alternate" type="application/rss+xml" title="High Scalability" href="https://highscalability.com/rss/">
<script defer src="https://cdn.jsdelivr.net/ghost/portal@~2.71/umd/portal.min.js" data-i18n="true" data-ghost="https://highscalability.com/" data-key="e7388272b3fcb33a0abfc2f95c" data-api="https://high-scalability.ghost.io/ghost/api/content/" data-locale="en" crossorigin="anonymous"></script><style id="gh-members-styles">.gh-post-upgrade-cta-content,
.gh-post-upgrade-cta {
display: flex;
flex-direction: column;
align-items: center;
font-family: -apple-system, BlinkMacSystemFont, 'Segoe UI', Roboto, Oxygen, Ubuntu, Cantarell, 'Open Sans', 'Helvetica Neue', sans-serif;
text-align: center;
width: 100%;
color: #ffffff;
font-size: 16px;
}
.gh-post-upgrade-cta-content {
border-radius: 8px;
padding: 40px 4vw;
}
.gh-post-upgrade-cta h2 {
color: #ffffff;
font-size: 28px;
letter-spacing: -0.2px;
margin: 0;
padding: 0;
}
.gh-post-upgrade-cta p {
margin: 20px 0 0;
padding: 0;
}
.gh-post-upgrade-cta small {
font-size: 16px;
letter-spacing: -0.2px;
}
.gh-post-upgrade-cta a {
color: #ffffff;
cursor: pointer;
font-weight: 500;
box-shadow: none;
text-decoration: underline;
}
.gh-post-upgrade-cta a:hover {
color: #ffffff;
opacity: 0.8;
box-shadow: none;
text-decoration: underline;
}
.gh-post-upgrade-cta a.gh-btn {
display: block;
background: #ffffff;
text-decoration: none;
margin: 28px 0 0;
padding: 8px 18px;
border-radius: 4px;
font-size: 16px;
font-weight: 600;
}
.gh-post-upgrade-cta a.gh-btn:hover {
opacity: 0.92;
}</style>
<script defer src="https://cdn.jsdelivr.net/ghost/sodo-search@~1.8/umd/sodo-search.min.js" data-key="e7388272b3fcb33a0abfc2f95c" data-styles="https://cdn.jsdelivr.net/ghost/sodo-search@~1.8/umd/main.css" data-sodo-search="https://high-scalability.ghost.io/" data-locale="en" crossorigin="anonymous"></script>
<link href="https://highscalability.com/webmentions/receive/" rel="webmention">
<script defer src="/public/cards.min.js?v=ShRHxgy4po8zN-Wf"></script>
<link rel="stylesheet" type="text/css" href="/public/cards.min.css?v=WwnU9jw5ancNC8Gc">
<script defer src="/public/comment-counts.min.js?v=oFYkaGLdiMqB8VN9" data-ghost-comments-counts-api="https://highscalability.com/members/api/comments/counts/"></script>
<script defer src="/public/member-attribution.min.js?v=AKG4hWena9j3yX3I"></script>
<script defer src="/public/ghost-stats.min.js?v=vFcCUf6ZQ0Hyhc8h" data-stringify-payload="false" data-datasource="analytics_events" data-storage="localStorage" data-host="https://highscalability.com/.ghost/analytics/api/v1/page_hit" tb_site_uuid="647204cd-7ad2-4539-b98a-3489074b932d" tb_post_uuid="213c4df8-871a-47ea-b29a-f2764009f564" tb_post_type="post" tb_member_uuid="undefined" tb_member_status="undefined" tb_gift_link=""></script><style>:root {--ghost-accent-color: #35cea0;}</style>
<!-- Fathom - beautiful, simple website analytics -->
<script src="https://cdn.usefathom.com/script.js" data-site="XBRJSZNU" defer></script>
<!-- / Fathom -->
<style>
/* Hide feature image on single post page in Source theme */
.post-template .gh-article-image {
display: none;
}
.gh-footer-copyright { display: none; }
</style>
</head>
<body class="post-template tag-consistent-hashing tag-sharding tag-load-balancing tag-hash-sqs has-sans-title has-sans-body">
<div class="gh-viewport">
<header id="gh-navigation" class="gh-navigation is-middle-logo gh-outer">
<div class="gh-navigation-inner gh-inner">
<div class="gh-navigation-brand">
<a class="gh-navigation-logo is-title" href="https://highscalability.com">
High Scalability
</a>
<button class="gh-search gh-icon-button" aria-label="Search this site" data-ghost-search>
<svg xmlns="http://www.w3.org/2000/svg" fill="none" viewBox="0 0 24 24" stroke="currentColor" stroke-width="2" width="20" height="20"><path stroke-linecap="round" stroke-linejoin="round" d="M21 21l-6-6m2-5a7 7 0 11-14 0 7 7 0 0114 0z"></path></svg></button> <button class="gh-burger gh-icon-button" aria-label="Menu">
<svg xmlns="http://www.w3.org/2000/svg" width="24" height="24" fill="currentColor" viewBox="0 0 256 256"><path d="M224,128a8,8,0,0,1-8,8H40a8,8,0,0,1,0-16H216A8,8,0,0,1,224,128ZM40,72H216a8,8,0,0,0,0-16H40a8,8,0,0,0,0,16ZM216,184H40a8,8,0,0,0,0,16H216a8,8,0,0,0,0-16Z"></path></svg> <svg xmlns="http://www.w3.org/2000/svg" width="24" height="24" fill="currentColor" viewBox="0 0 256 256"><path d="M205.66,194.34a8,8,0,0,1-11.32,11.32L128,139.31,61.66,205.66a8,8,0,0,1-11.32-11.32L116.69,128,50.34,61.66A8,8,0,0,1,61.66,50.34L128,116.69l66.34-66.35a8,8,0,0,1,11.32,11.32L139.31,128Z"></path></svg> </button>
</div>
<nav class="gh-navigation-menu">
<ul class="nav">
<li class="nav-home"><a href="https://highscalability.com/">Home</a></li>
<li class="nav-system-design-interview-course"><a href="https://bit.ly/hishscalcourse">System Design Interview Course</a></li>
</ul>
</nav>
<div class="gh-navigation-actions">
<button class="gh-search gh-icon-button" aria-label="Search this site" data-ghost-search>
<svg xmlns="http://www.w3.org/2000/svg" fill="none" viewBox="0 0 24 24" stroke="currentColor" stroke-width="2" width="20" height="20"><path stroke-linecap="round" stroke-linejoin="round" d="M21 21l-6-6m2-5a7 7 0 11-14 0 7 7 0 0114 0z"></path></svg></button> <div class="gh-navigation-members">
<a href="#/portal/signin" data-portal="signin">Sign in</a>
<a class="gh-button" href="#/portal/signup" data-portal="signup">Subscribe</a>
</div>
</div>
</div>
</header>
<main class="gh-main">
<article class="gh-article post tag-consistent-hashing tag-sharding tag-load-balancing tag-hash-sqs no-image">
<header class="gh-article-header gh-canvas">
<a class="gh-article-tag" href="https://highscalability.com/tag/consistent-hashing/">consistent hashing</a>
<h1 class="gh-article-title is-title">Consistent hashing algorithm</h1>
<div class="gh-meta-share">
<div class="gh-article-meta">
<div class="gh-article-author-image instapaper_ignore">
<a href="/author/sdxyz/"><svg viewBox="0 0 24 24" xmlns="http://www.w3.org/2000/svg"><g fill="none" fill-rule="evenodd"><path d="M3.513 18.998C4.749 15.504 8.082 13 12 13s7.251 2.504 8.487 5.998C18.47 21.442 15.417 23 12 23s-6.47-1.558-8.487-4.002zM12 12c2.21 0 4-2.79 4-5s-1.79-4-4-4-4 1.79-4 4 1.79 5 4 5z" fill="#FFF"/></g></svg>
</a>
</div>
<div class="gh-article-meta-wrapper">
<h4 class="gh-article-author-name"><a href="/author/sdxyz/">NK</a></h4>
<div class="gh-article-meta-content">
<time class="gh-article-meta-date" datetime="2023-02-22">22 Feb 2023</time>
<span class="gh-article-meta-length"><span class="bull">—</span> 16 min read</span>
</div>
</div>
</div>
<a href="#/share" class="gh-button gh-button-share">
Share
</a>
</div>
</header>
<section class="gh-content gh-canvas is-body">
<figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/52705024790_8a89921cfb.jpg" class="kg-image" alt loading="lazy"></figure><p>You can <strong>subscribe to the <a href="https://newsletter.systemdesign.one/?ref=highscalability.com">system design newsletter</a> to excel in system design interviews and software architecture</strong>.</p><p>You can view the original article <a href="https://systemdesign.one/consistent-hashing-explained/?ref=highscalability.com">Consistent hashing explained</a> on <a href="https://systemdesign.one/?ref=highscalability.com">systemdesign.one</a> website.</p><h2 id="how-does-consistent-hashing-work">How does consistent hashing work?</h2><p>At a high level, consistent hashing performs the following operations:</p><ol><li>The output of the hash function is placed on a virtual ring structure (known as the hash ring)</li><li>The hashed IP addresses of the nodes are used to assign a position for the nodes on the hash ring</li><li>The key of a data object is hashed using the same hash function to find the position of the key on the hash ring</li><li>The hash ring is traversed in the clockwise direction starting from the position of the key until a node is found</li><li>The data object is stored or retrieved from the node that was found</li></ol><h2 id="terminology">Terminology</h2><p>The following terminology might be useful for you:</p><ul><li>Node: a server that provides functionality to other services</li><li>Hash function: a mathematical function used to map data of arbitrary size to fixed-size values</li><li>Data partitioning: a technique of distributing data across multiple nodes to improve the performance and scalability of the system</li><li>Data replication: a technique of storing multiple copies of the same data on different nodes to improve the availability and durability of the system</li><li>Hotspot: A performance-degraded node in a distributed system due to a large share of data storage and a high volume of retrieval or storage requests</li><li><a href="https://systemdesign.one/gossip-protocol/?ref=highscalability.com" rel="noopener noreffer ">Gossip protocol</a>: peer-to-peer communication technique used by nodes to periodically exchange state information</li></ul><h2 id="requirements">Requirements</h2><h3 id="functional-requirements">Functional Requirements</h3><ul><li>Design an algorithm to horizontally scale the cache servers</li><li>The algorithm must minimize the occurrence of hotspots in the network</li><li>The algorithm must be able to handle internet-scale dynamic load</li><li>The algorithm must reuse existing network protocols such as TCP/IP</li></ul><h3 id="non-functional-requirements">Non-Functional Requirements</h3><ul><li>Scalable</li><li>High availability</li><li>Low latency</li><li>Reliable</li></ul><h2 id="introduction">Introduction</h2><p>A website can become extremely popular in a relatively short time frame. The increased load might swamp and degrade the performance of the website. The cache server is used to improve the latency and reduce the load on the system. The cache servers must scale to meet the dynamic demand as a fixed collection of cache servers will not be able to handle the dynamic load. In addition, the occurrence of multiple cache misses might swamp the origin server.</p><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/yle8ukj.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/yle8ukj.png 600w"></figure><p>Figure 1: Cache replication</p><p>The replication of the cache improves the availability of the system. However, replication of the cache does not solve the dynamic load problem as only a limited data set can be cached [1]. The tradeoffs of the cache replication approach are the following:</p><ul><li>only a limited data set is cached</li><li><a href="https://systemdesign.one/consistency-patterns/?ref=highscalability.com" rel="noopener noreffer ">consistency</a> between cache replicas is expensive to maintain</li></ul><p>The <strong>spread</strong> is the number of cache servers holding the same key-value pair (data object). The <strong>load</strong> is the number of distinct data objects assigned to a cache server. The optimal configuration for the high performance of a cache server is to keep the spread and the load at a minimum [2].</p><figure class="kg-card kg-image-card"><a href="http://localhost:2369/content/images/2024/02/pgcf6hf.png?ref=highscalability.com"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/pgcf6hf-__squarespace_cacheversion-1675879926185.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/pgcf6hf-__squarespace_cacheversion-1675879926185.png 600w"></a></figure><p>Figure 2: Dynamic hashing</p><p>The data set must be partitioned (shard) among multiple cache servers (<strong>nodes</strong>) to horizontally scale. The replication and partitioning of nodes are orthogonal to each other. Multiple data partitions are stored on a single node for improved fault tolerance and increased throughput. The reasons for partitioning are the following [1]:</p><ul><li>a cache server is memory bound</li><li>increased throughput</li></ul><h2 id="partitioning">Partitioning</h2><p>The data set is partitioned among multiple nodes to horizontally scale out. The different techniques for partitioning the cache servers are the following [1]:</p><ul><li>Random assignment</li><li>Single global cache</li><li>Key range partitioning</li><li>Static hash partitioning</li><li>Consistent hashing</li></ul><h3 id="random-assignment">Random assignment</h3><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/lch5ja1-__squarespace_cacheversion-1675880048742.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/lch5ja1-__squarespace_cacheversion-1675880048742.png 600w"></figure><p>Figure 3: Partitioning; Random assignment</p><p>The server distributes the data objects randomly across the cache servers. The random assignment of a large data set results in a relatively uniform distribution of data. However, the client cannot easily identify the node to retrieve the data due to random distribution. In conclusion, the random assignment solution will not scale to handle the dynamic load.</p><h3 id="single-global-cache">Single global cache</h3><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/gvgsnzo-__squarespace_cacheversion-1675880110463.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/gvgsnzo-__squarespace_cacheversion-1675880110463.png 600w"></figure><p>Figure 4: Partitioning; Single global cache</p><p>The server stores the whole data set on a single global cache server. The data objects are easily retrieved by the client at the expense of degraded performance and decreased availability of the system. In conclusion, the single global cache solution will not scale to handle the dynamic load.</p><h3 id="key-range-partitioning">Key range partitioning</h3><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/fbri3nw-__squarespace_cacheversion-1675880154142.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/fbri3nw-__squarespace_cacheversion-1675880154142.png 600w"></figure><p>Figure 5: Key range partitioning</p><p>The cache servers are partitioned using the key range of the data set. The client can easily retrieve the data from cache servers. The data set is not necessarily uniformly distributed among the cache servers as there might be more keys in a certain key range. In conclusion, the key range partitioning solution will not scale to handle the dynamic load.</p><h3 id="static-hash-partitioning">Static hash partitioning</h3><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/8wmqv5c-__squarespace_cacheversion-1675880212408.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/8wmqv5c-__squarespace_cacheversion-1675880212408.png 600w"></figure><p>Figure 6: Static hash partitioning</p><p>The identifiers (internet protocol address or domain name) of the nodes are placed on an array of length N. The modulo hash service computes the hash of the data key and executes a modulo N operation to locate the array index (node identifier) to store or retrieve a key. The time complexity to locate a node identifier (<strong>ID</strong>) in static hash partitioning is constant O(1).</p><p>Static hash partitioning</p><p><strong>node ID = hash(key) mod N</strong></p><p>where N is the array’s length and the key is the key of the data object.</p><p>A collision occurs when multiple nodes are assigned to the same position on the array. The techniques to resolve a collision are <a href="https://en.wikipedia.org/wiki/Open_addressing?ref=highscalability.com" rel="noopener noreffer ">open addressing</a> and <a href="https://en.wikipedia.org/wiki/Hash_table?ref=highscalability.com#Collision_resolution" rel="noopener noreffer ">chaining</a>. The occurrence of a collision degrades the time complexity of the cache nodes.</p><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/q1dlmkl-__squarespace_cacheversion-1675880256016.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/q1dlmkl-__squarespace_cacheversion-1675880256016.png 600w"></figure><p>Figure 7: Static hash partitioning; Node failure</p><p>The static hash partitioning is not horizontally scalable. The removal of a node (due to a server crash) breaks the existing mappings between the keys and nodes. The keys must be rehashed to restore mapping between keys and nodes [3].</p><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/rcktxu9-__squarespace_cacheversion-1675880299125.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/rcktxu9-__squarespace_cacheversion-1675880299125.png 600w"></figure><p>Figure 8: Static hash partitioning; Node added</p><p>New nodes must be provisioned to handle the increasing load. The addition of a node breaks the existing mappings between the keys and nodes. The following are the drawbacks of static hash partitioning:</p><ul><li>nodes will not horizontally scale to handle the dynamic load</li><li>the addition or removal of a node breaks the mapping between keys and nodes</li><li>massive data movement when the number of nodes changes</li></ul><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/i0zcxej-__squarespace_cacheversion-1675880346738.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/i0zcxej-__squarespace_cacheversion-1675880346738.png 600w"></figure><p>Figure 9: Static hash partitioning; Data movement due to node failure</p><p>In conclusion, the data set must be rehashed or moved between nodes when the number of nodes changes. The majority of the requests in the meantime will result in cache misses. The requests are delegated to the origin server on cache misses. The heavy load on the origin server might swamp and degrade the service [3].</p><h2 id="further-system-design-learning-resources">Further System Design Learning Resources</h2><p><strong><a href="https://newsletter.systemdesign.one/?ref=highscalability.com" rel="noopener noreffer ">Subscribe to the system design newsletter</a></strong> and never miss a new blog post again. You will also receive the <strong>ultimate guide</strong> to approaching <strong>system design interviews</strong> on <a href="https://newsletter.systemdesign.one/?ref=highscalability.com" rel="noopener noreffer ">newsletter sign-up</a>.</p><h3 id="consistent-hashing">Consistent hashing</h3><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/pyy2uyb-__squarespace_cacheversion-1675880389034.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/pyy2uyb-__squarespace_cacheversion-1675880389034.png 600w"></figure><p>Figure 10: Consistent hashing</p><p>Consistent hashing is a distributed systems technique that operates by assigning the data objects and nodes a position on a virtual ring structure (hash ring). Consistent hashing minimizes the number of keys to be remapped when the total number of nodes changes [4].</p><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/9azf3jd-__squarespace_cacheversion-1675880438400.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/9azf3jd-__squarespace_cacheversion-1675880438400.png 600w"></figure><p>Figure 11: Hash function mapping</p><p>The basic gist behind the consistent hashing algorithm is to hash both node identifiers and data keys using the same hash function. A uniform and independent hashing function such as message-digest 5 (<a href="https://en.wikipedia.org/wiki/MD5?ref=highscalability.com" rel="noopener noreffer "><strong>MD5</strong></a>) is used to find the position of the nodes and keys (data objects) on the hash ring. The output range of the hash function must be of reasonable size to prevent collisions.</p><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/oibeqcu-__squarespace_cacheversion-1675880480125.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/oibeqcu-__squarespace_cacheversion-1675880480125.png 600w"></figure><p>Figure 12: Consistent hash ring</p><p>The output space of the hash function is treated as a fixed circular space to form the hash ring. The largest hash value wraps around the smallest hash value. The hash ring is considered to have a finite number of positions [5].</p><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/tcf1r4l-__squarespace_cacheversion-1675880540189.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/tcf1r4l-__squarespace_cacheversion-1675880540189.png 600w"></figure><p>Figure 13: Consistent hashing; Positioning the nodes on the hash ring</p><p>The following operations are executed to locate the position of a node on the hash ring [4]:</p><ol><li>Hash the internet protocol (<strong>IP</strong>) address or domain name of the node using a hash function</li><li>The hash code is base converted</li><li>Modulo the hash code with the total number of available positions on the hash ring</li></ol><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/qcd8qf1-__squarespace_cacheversion-1675880841904.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/qcd8qf1-__squarespace_cacheversion-1675880841904.png 600w"></figure><p>Figure 14: Consistent hashing; Node position</p><p>Suppose the hash function produces an output space size of 10 bits (2¹⁰ = 1024), the hash ring formed is a virtual circle with a number range starting from 0 to 1023. The hashed value of the IP address of a node is used to assign a location for the node on the hash ring.</p><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/bbn13ap-__squarespace_cacheversion-1675880910715.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/bbn13ap-__squarespace_cacheversion-1675880910715.png 600w"></figure><p>Figure 15: Consistent hashing; Storing a data object (key)</p><p>The key of the data object is hashed using the same hash function to locate the position of the key on the hash ring. The hash ring is traversed in the clockwise direction starting from the position of the key until a node is found. The data object is stored on the node that was found. In simple words, the first node with a position value greater than the position of the key stores the data object [6].</p><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/bw0tzkx-__squarespace_cacheversion-1675880954242.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/bw0tzkx-__squarespace_cacheversion-1675880954242.png 600w"></figure><p>Figure 16: Consistent hashing; Retrieving a data object (key)</p><p>The key of the data object is hashed using the same hash function to locate the position of the key on the hash ring. The hash ring is traversed in the clockwise direction starting from the position of the key until a node is found. The data object is retrieved from the node that was found. In simple words, the first node with a position value greater than the position of the key must hold the data object.</p><p>Each node is responsible for the region on the ring between the node and its predecessor node on the hash ring. The origin server must be queried on a cache miss. In conclusion, the following operations are performed for consistent hashing [7]:</p><ol><li>The output of the hash function such as MD5 is placed on the hash ring</li><li>The IP address of the nodes is hashed to find the position of the nodes on the hash ring</li><li>The key of the data object is hashed using the same hash function to locate the position of the key on the hash ring</li><li>Traverse the hash ring in the clockwise direction starting from the position of the key until the next node to identify the correct node to store or retrieve the data object</li></ol><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/ezsk5bf-__squarespace_cacheversion-1675880999540.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/ezsk5bf-__squarespace_cacheversion-1675880999540.png 600w"></figure><p>Figure 17: Consistent hashing; Deletion of a node</p><p>The failure (crash) of a node results in the movement of data objects from the failed node to the immediate neighboring node in the clockwise direction. The remaining nodes on the hash ring are unaffected [5].</p><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/9d3fwgy-__squarespace_cacheversion-1675881046368.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/9d3fwgy-__squarespace_cacheversion-1675881046368.png 600w"></figure><p>Figure 18: Consistent hashing; Addition of a node</p><p>When a new node is provisioned and added to the hash ring, the keys (data objects) that fall within the range of the new node are moved out from the immediate neighboring node in the clockwise direction.</p><p>Consistent hashing</p><p><strong>Average number of keys stored on a node = k/N</strong></p><p>where k is the total number of keys (data objects) and N is the number of nodes.</p><p>The deletion or addition of a node results in the movement of an average number of keys stored on a single node. Consistent hashing aid cloud computing by minimizing the movement of data when the total number of nodes changes due to dynamic load [8].</p><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/lbalo0c-__squarespace_cacheversion-1675881100463.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/lbalo0c-__squarespace_cacheversion-1675881100463.png 600w"></figure><p>Figure 19: Consistent hashing; Non-uniform positioning of nodes</p><p>There is a chance that nodes are not uniformly distributed on the consistent hash ring. The nodes that receive a huge amount of traffic become hotspots resulting in <a href="https://en.wikipedia.org/wiki/Cascading_failure?ref=highscalability.com" rel="noopener noreffer ">cascading failure</a> of the nodes.</p><figure class="kg-card kg-image-card"><img src="https://imgur.com/h01BDY14" class="kg-image" alt loading="lazy"></figure><p>Figure 20: Consistent hashing; Virtual nodes</p><p>The nodes are assigned to multiple positions on the hash ring by hashing the node IDs through distinct hash functions to ensure uniform distribution of keys among the nodes. The technique of assigning multiple positions to a node is known as a <strong>virtual node</strong>. The virtual nodes improve the load balancing of the system and prevent hotspots. The number of positions for a node is decided by the heterogeneity of the node. In other words, the nodes with a higher capacity are assigned more positions on the hash ring [5].</p><p>The data objects can be replicated on adjacent nodes to minimize the data movement when a node crashes or when a node is added to the hash ring. In conclusion, consistent hashing resolves the problem of dynamic load.</p><h2 id="consistent-hashing-implementation">Consistent hashing implementation</h2><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/a1ro3ri-__squarespace_cacheversion-1675881211193.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/a1ro3ri-__squarespace_cacheversion-1675881211193.png 600w"></figure><p>Figure 21: Consistent hashing implementation; Binary search tree storing the node positions</p><p>The self-balancing binary search tree (<strong>BST</strong>) data structure is used to store the positions of the nodes on the hash ring. The BST offers logarithmic O(log n) time complexity for search, insert, and delete operations. The keys of the BST contain the positions of the nodes on the hash ring.</p><p>The BST data structure is stored on a centralized highly available service. As an alternative, the BST data structure is stored on each node, and the state information between the nodes is synchronized through the <a href="https://systemdesign.one/gossip-protocol/?ref=highscalability.com" rel="noopener noreffer ">gossip protocol</a> [7].</p><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/gkufpyg-__squarespace_cacheversion-1675881248016.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/gkufpyg-__squarespace_cacheversion-1675881248016.png 600w"></figure><p>Figure 22: Consistent hashing implementation; Insertion of a data object (key)</p><p>In the diagram, suppose the hash of an arbitrary key ‘xyz’ yields the hash code output 5. The successor BST node is 6 and the data object with the key ‘xyz’ is stored on the node that is at position 6. In general, the following operations are executed to insert a key (data object):</p><ol><li>Hash the key of the data object</li><li>Search the BST in logarithmic time to find the BST node immediately greater than the hashed output</li><li>Store the data object in the successor node</li></ol><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/7zs0lss-__squarespace_cacheversion-1675881284630.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/7zs0lss-__squarespace_cacheversion-1675881284630.png 600w"></figure><p>Figure 23: Consistent hashing implementation; Insertion of a node</p><p>The insertion of a new node results in the movement of data objects that fall within the range of the new node from the successor node. Each node might store an internal or an external BST to track the keys allocated in the node. The following operations are executed to insert a node on the hash ring:</p><ol><li>Insert the hash of the node ID in BST in logarithmic time</li><li>Identify the keys that fall within the subrange of the new node from the successor node on BST</li><li>Move the keys to the new node</li></ol><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/pipd4fz-__squarespace_cacheversion-1675881320902.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/pipd4fz-__squarespace_cacheversion-1675881320902.png 600w"></figure><p>Figure 24: Consistent hashing implementation; Deletion of a node</p><p>The deletion of a node results in the movement of data objects that fall within the range of the decommissioned node to the successor node. An additional external BST can be used to track the keys allocated in the node. The following operations are executed to delete a node on the hash ring:</p><ol><li>Delete the hash of the decommissioned node ID in BST in logarithmic time</li><li>Identify the keys that fall within the range of the decommissioned node</li><li>Move the keys to the successor node</li></ol><h2 id="what-is-the-asymptotic-complexity-of-consistent-hashing">What is the asymptotic complexity of consistent hashing?</h2><p>The asymptotic complexity of consistent hashing operations are the following:</p><!--kg-card-begin: html--><table>
<thead>
<tr>
<th>Operation</th><th>Time Complexity</th><th>Description</th>
</tr>
</thead>
<tbody>
<tr>
<td>Add a node</td>
<td>O(k/n + logn)</td>
<td>O(k/n) for redistribution of keys O(logn) for binary search tree traversal</td>
</tr>
<tr>
<td>Remove a node</td>
<td>O(k/n + logn)</td>
<td>O(k/n) for redistribution of keys O(logn) for binary search tree traversal</td>
</tr>
<tr>
<td>Add a key</td>
<td>O(logn)</td>
<td>O(logn) for binary search tree traversal</td>
</tr>
<tr>
<td>Remove a key</td>
<td>O(logn)</td>
<td>O(logn) for binary search tree traversal</td>
</tr>
</tbody>
</table><!--kg-card-end: html--><p>where k = total number of keys, n = total number of nodes [2, 7].</p><h2 id="further-system-design-learning-resources-1">Further System Design Learning Resources</h2><p><strong><a href="https://systemdesignone.substack.com/" rel="noopener noreffer ">Subscribe to the system design newsletter</a></strong> and never miss a new blog post again. You will also receive the <strong>ultimate guide</strong> to approaching <strong>system design interviews</strong> on <a href="https://newsletter.systemdesign.one/?ref=highscalability.com" rel="noopener noreffer ">newsletter sign-up</a>.</p><h2 id="how-to-handle-concurrency-in-consistent-hashing">How to handle concurrency in consistent hashing?</h2><p>The BST that stores the positions of the nodes is a mutable data structure that must be synchronized when multiple nodes are added or removed at the same time on the hash ring. The <a href="https://en.wikipedia.org/wiki/Readers%E2%80%93writer_lock?ref=highscalability.com" rel="noopener noreffer ">readers-writer lock</a> is used to synchronize BST at the expense of a slight increase in latency.</p><h2 id="what-hash-functions-are-used-in-consistent-hashing">What hash functions are used in consistent hashing?</h2><p>An optimal hash function for consistent hashing must be fast and produce uniform output. The cryptographic hash functions such as <a href="https://en.wikipedia.org/wiki/MD5?ref=highscalability.com" rel="noopener noreffer ">MD5</a>, and the secure hash algorithms <a href="https://en.wikipedia.org/wiki/SHA-1?ref=highscalability.com" rel="noopener noreffer ">SHA-1</a> and <a href="https://en.wikipedia.org/wiki/SHA-2?ref=highscalability.com" rel="noopener noreffer ">SHA-256</a> are not relatively fast. <a href="https://en.wikipedia.org/wiki/MurmurHash?ref=highscalability.com" rel="noopener noreffer ">MurmurHash</a> is a relatively cheaper hash function. The non-cryptographic hash functions like <a href="https://github.com/cespare/xxhash?ref=highscalability.com" rel="noopener noreffer ">xxHash</a>, <a href="https://github.com/dgryski/go-metro?ref=highscalability.com" rel="noopener noreffer ">MetroHash</a>, or <a href="https://github.com/dgryski/go-sip13?ref=highscalability.com" rel="noopener noreffer ">SipHash1–3</a> are other potential candidates [6].</p><h2 id="what-are-the-benefits-of-consistent-hashing">What are the benefits of consistent hashing?</h2><p>The following are the advantages of consistent hashing [3]:</p><ul><li>horizontally scalable</li><li>minimized data movement when the number of nodes changes</li><li>quick replication and partitioning of data</li></ul><p>The following are the advantages of virtual nodes [5]:</p><ul><li>load handled by a node is uniformly distributed across the remaining available nodes during downtime</li><li>the newly provisioned node accepts an equivalent amount of load from the available nodes</li><li>fair distribution of load among heterogenous nodes</li></ul><h2 id="what-are-the-drawbacks-of-consistent-hashing">What are the drawbacks of consistent hashing?</h2><p>The following are the disadvantages of consistent hashing [5]:</p><ul><li>cascading failure due to hotspots</li><li>non-uniform distribution of nodes and data</li><li>oblivious to the heterogeneity in the performance of nodes</li></ul><p>The following are the disadvantages of virtual nodes [5, 6, 8]:</p><ul><li>when a specific data object becomes extremely popular, consistent hashing will still send all the requests for the popular data object to the same subset of nodes resulting in a degradation of the service</li><li><a href="https://systemdesign.one/back-of-the-envelope?ref=highscalability.com" rel="noopener noreffer ">capacity planning</a> is trickier with virtual nodes</li><li>memory costs and operational complexity increase due to the maintenance of BST</li><li>replication of data objects is challenging due to the additional logic to identify the distinct physical nodes</li><li>downtime of a virtual node affects multiple nodes on the ring</li></ul><h2 id="what-are-the-consistent-hashing-examples">What are the consistent hashing examples?</h2><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/t47t8lk-__squarespace_cacheversion-1675881374193.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/t47t8lk-__squarespace_cacheversion-1675881374193.png 600w"></figure><p>Figure 25: Consistent hashing example: Discord</p><p>The discord server (discord space or chat room) is hosted on a set of nodes. The client of the discord chat application identifies the set of nodes that hosts a specific discord server using consistent hashing [9].</p><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/efyr0wb-__squarespace_cacheversion-1675881596983.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/efyr0wb-__squarespace_cacheversion-1675881596983.png 600w"></figure><p>Figure 26: Consistent hashing example: Amazon Dynamo</p><p>The distributed NoSQL data stores such as Amazon DynamoDB, Apache Cassandra, and Riak use consistent hashing to dynamically partition the data set across the set of nodes. The data is partitioned for incremental scalability [5].</p><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/0cglhiz-__squarespace_cacheversion-1675881653619.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/0cglhiz-__squarespace_cacheversion-1675881653619.png 600w"></figure><p>Figure 27: Consistent hashing example: Vimeo</p><p>The video storage and streaming service Vimeo uses consistent hashing for load balancing the traffic to stream videos [8].</p><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/hkdzadr-__squarespace_cacheversion-1675881674703.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/hkdzadr-__squarespace_cacheversion-1675881674703.png 600w"></figure><p>Figure 28: Consistent hashing example: Netflix</p><p>The video streaming service Netflix uses consistent hashing to distribute the uploaded video content across the content delivery network (<strong>CDN</strong>) [10].</p><h2 id="consistent-hashing-algorithm-real-world-implementation">Consistent hashing algorithm real-world implementation</h2><p>The clients of Memcached (Ketama), and Amazon Dynamo support consistent hashing out of the box [3, 11]. The HAProxy includes the bounded-load consistent hashing algorithm for load balancing the traffic [8]. As an alternative, the consistent hashing algorithm can be implemented, in the language of choice.</p><h2 id="consistent-hashing-optimization">Consistent hashing optimization</h2><p>Some of the popular variants of consistent hashing are the following:</p><ul><li>Multi-probe consistent hashing</li><li>Consistent hashing with bounded loads</li></ul><h3 id="multi-probe-consistent-hashing">Multi-probe consistent hashing</h3><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/ykzvyl2-__squarespace_cacheversion-1675881740883.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/ykzvyl2-__squarespace_cacheversion-1675881740883.png 600w"></figure><p>Figure 29: Consistent hashing optimization; Multi-probe consistent hashing</p><p>The Multi-probe consistent hashing offers linear O(n) space complexity to store the positions of nodes on the hash ring. There are no virtual nodes but a node is assigned only a single position on the hash ring. The amortized time complexity for the addition and removal of nodes is constant O(1). However, the key (data object) lookups are relatively slower.</p><p>The basic gist of multi-probe consistent hashing is to hash the key (data object) multiple times using distinct hash functions on lookup and the closest node in the clockwise direction returns the data object [12].</p><h3 id="consistent-hashing-with-bounded-loads">Consistent hashing with bounded loads</h3><figure class="kg-card kg-image-card"><img src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/pkv04wr-__squarespace_cacheversion-1675881762243.png" class="kg-image" alt loading="lazy" srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/2024/02/pkv04wr-__squarespace_cacheversion-1675881762243.png 600w"></figure><p>Figure 30: Consistent hashing optimization; Bounded-load consistent hashing</p><p>The consistent hashing with bounded load puts an upper limit on the load received by a node on the hash ring, relative to the average load of the whole hash ring. The distribution of requests is the same as consistent hashing as long as the nodes are not overloaded [13].</p><p>When a specific data object becomes extremely popular, the node hosting the data object receives a significant amount of traffic resulting in the degradation of the service. If a node is overloaded, the incoming request is delegated to a fallback node. The list of fallback nodes will be the same for the same request hash. In simple words, the same node(s) will consistently be the “second choice” for a popular data object. The fallback nodes resolve the popular data object caching problem.</p><p>If a node is overloaded, the list of the fallback nodes will usually be different for different request hashes. In other words, the requests to an overloaded node are distributed among the available nodes instead of a single fallback node.</p><h2 id="summary">Summary</h2><p>Consistent hashing is popular among distributed systems such internet-scale <a href="https://systemdesign.one/url-shortening-system-design/?ref=highscalability.com">URL shortener</a> , and <a href="https://systemdesign.one/system-design-pastebin/?ref=highscalability.com" rel="noopener noreffer ">Pastebin</a>. The most common use cases of consistent hashing are data partitioning and load balancing.</p><h2 id="further-system-design-learning-resources-2">Further System Design Learning Resources</h2><p><strong><a href="https://systemdesignone.substack.com/" rel="noopener noreffer ">Subscribe to the system design newsletter</a></strong> and never miss a new blog post again. You will also receive the <strong>ultimate guide</strong> to approaching <strong>system design interviews</strong> on <a href="https://newsletter.systemdesign.one/?ref=highscalability.com" rel="noopener noreffer ">newsletter sign-up</a>.</p><h2 id="references">References</h2><p>[1] Lindsey Kuper, <a href="https://www.youtube.com/watch?v=uNQGP0yupn0&list=PLNPUF5QyWU8PydLG2cIJrCvnn5I_exhYx&index=18&ref=highscalability.com" rel="noopener noreffer ">UC Santa Cruz CSE138 (Distributed Systems) Lecture 15: introduction to sharding; consistent hashing</a> (2021)</p><p>[2] David Karger, Eric Lehman, Tom Leighton, Rina Panigrahy, Matthew Levine, Daniel Lewin, <a href="https://github.com/papers-we-love/papers-we-love/blob/master/distributed_systems/consistent-hashing-and-random-trees.pdf?ref=highscalability.com" rel="noopener noreffer ">Consistent Hashing and Random Trees: Distributed Caching Protocols for Relieving Hot Spots on the World Wide Web</a> (1997)</p><p>[3] Tom White, <a href="http://tom-e-white.com/2007/11/consistent-hashing.html?ref=highscalability.com" rel="noopener noreffer ">Consistent Hashing</a> (2007), tom-e-white.com</p><p>[4] Srushtika Neelakantam, <a href="https://ably.com/blog/implementing-efficient-consistent-hashing?ref=highscalability.com" rel="noopener noreffer ">Consistent hashing explained</a> (2018), ably.com</p><p>[5] Giuseppe DeCandia, Deniz Hastorun, Madan Jampani, Gunavardhan Kakulapati, Avinash Lakshman, Alex Pilchin, Swaminathan Sivasubramanian, Peter Vosshall and Werner Vogels, <a href="https://github.com/papers-we-love/papers-we-love/blob/master/datastores/dynamo-amazons-highly-available-key-value-store.pdf?ref=highscalability.com" rel="noopener noreffer ">Dynamo: Amazon’s Highly Available Key-value Store</a> (2007)</p><p>[6] Damian Gryski, <a href="https://dgryski.medium.com/consistent-hashing-algorithmic-tradeoffs-ef6b8e2fcae8?ref=highscalability.com" rel="noopener noreffer ">Consistent Hashing: Algorithmic Tradeoffs</a> (2018), medium.com</p><p>[7] <a href="http://people.csail.mit.edu/moitra/854.html?ref=highscalability.com" rel="noopener noreffer ">MIT 6.854 Spring 2016 Lecture 3: Consistent Hashing and Random Trees</a> (2016)</p><p>[8] <a href="https://medium.com/vimeo-engineering-blog/improving-load-balancing-with-a-new-consistent-hashing-algorithm-9f1bd75709ed?ref=highscalability.com" rel="noopener noreffer ">Improving load balancing with a new consistent-hashing algorithm</a> (2016), Vimeo Engineering Blog</p><p>[9] Stanislav Vishnevskiy, <a href="https://discord.com/blog/how-discord-scaled-elixir-to-5-000-000-concurrent-users?ref=highscalability.com" rel="noopener noreffer ">How discord scaled elixir to 5,000,000 concurrent users</a> (2017), discord.com</p><p>[10] Mohit Vora, Andrew Berglund, Videsh Sadafal, David Pfitzner, and Ellen Livengood, <a href="https://netflixtechblog.com/distributing-content-to-open-connect-3e3e391d4dc9?gi=294346617760&ref=highscalability.com" rel="noopener noreffer ">Distributing Content to Open Connect</a> (2017), netflixtechblog.com</p><p>[11] <a href="https://www.last.fm/user/RJ/journal/2007/04/10/rz_libketama_-_a_consistent_hashing_algo_for_memcache_clients?ref=highscalability.com" rel="noopener noreffer ">libketama - a consistent hashing algo for Memcache clients</a> (2007)</p><p>[12] Ben Appleton, Michael O’Reilly, <a href="https://arxiv.org/abs/1505.00062?ref=highscalability.com" rel="noopener noreffer ">Multi-Probe Consistent Hashing</a> (2015), Google</p><p>[13] Vahab Mirrokni, Mikkel Thorup, and Morteza Zadimoghaddam, <a href="https://arxiv.org/abs/1608.01350?ref=highscalability.com" rel="noopener noreffer ">Consistent Hashing with Bounded Loads</a> (2017), Google</p>
</section>
</article>
<div class="gh-comments gh-canvas">
<script defer src="https://cdn.jsdelivr.net/ghost/comments-ui@~1.6/umd/comments-ui.min.js" data-locale="en" data-ghost-comments="https://highscalability.com/" data-api="https://high-scalability.ghost.io/ghost/api/content/" data-admin="https://high-scalability.ghost.io/ghost/" data-key="e7388272b3fcb33a0abfc2f95c" data-title="null" data-count="true" data-post-id="65bceaca3953980001659b52" data-color-scheme="auto" data-avatar-saturation="60" data-accent-color="#35cea0" data-comments-enabled="all" data-publication="High Scalability" crossorigin="anonymous"></script>
</div>
</main>
<section class="gh-container is-grid gh-outer">
<div class="gh-container-inner gh-inner">
<h2 class="gh-container-title">Read more</h2>
<div class="gh-feed">
<article class="gh-card post">
<a class="gh-card-link" href="/untitled-2/">
<figure class="gh-card-image">
<img
srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w160/format/webp/2024/05/pasted-image-0-2.png 160w,
https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w320/format/webp/2024/05/pasted-image-0-2.png 320w,
https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w600/format/webp/2024/05/pasted-image-0-2.png 600w,
https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w960/format/webp/2024/05/pasted-image-0-2.png 960w,
https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w1200/format/webp/2024/05/pasted-image-0-2.png 1200w,
https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w2000/format/webp/2024/05/pasted-image-0-2.png 2000w"
sizes="320px"
src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w600/2024/05/pasted-image-0-2.png"
alt="Kafka 101"
loading="lazy"
>
</figure>
<div class="gh-card-wrapper">
<h3 class="gh-card-title is-title">Kafka 101</h3>
<p class="gh-card-excerpt is-body">This is a guest article by Stanislav Kozlovski, an Apache Kafka Committer. If you would like to connect with Stanislav, you can do so on Twitter and LinkedIn.
Originally developed in LinkedIn during 2011, Apache Kafka is one of the most popular open-source Apache projects out there. So far</p>
<footer class="gh-card-meta">
<!--
-->
<span class="gh-card-author">By ByteByteGo</span>
<time class="gh-card-date" datetime="2024-05-09">09 May 2024</time>
<!--
--></footer>
</div>
</a>
</article>
<article class="gh-card post">
<a class="gh-card-link" href="/capturing-a-billion-emo-j-i-ons/">
<figure class="gh-card-image">
<img
srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w160/format/webp/2024/03/1-rSRWALA4XzOdDcn-5vv7Zw.gif 160w,
https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w320/format/webp/2024/03/1-rSRWALA4XzOdDcn-5vv7Zw.gif 320w,
https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w600/format/webp/2024/03/1-rSRWALA4XzOdDcn-5vv7Zw.gif 600w,
https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w960/format/webp/2024/03/1-rSRWALA4XzOdDcn-5vv7Zw.gif 960w,
https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w1200/format/webp/2024/03/1-rSRWALA4XzOdDcn-5vv7Zw.gif 1200w,
https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w2000/format/webp/2024/03/1-rSRWALA4XzOdDcn-5vv7Zw.gif 2000w"
sizes="320px"
src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w600/2024/03/1-rSRWALA4XzOdDcn-5vv7Zw.gif"
alt="Capturing A Billion Emo(j)i-ons"
loading="lazy"
>
</figure>
<div class="gh-card-wrapper">
<h3 class="gh-card-title is-title">Capturing A Billion Emo(j)i-ons</h3>
<p class="gh-card-excerpt is-body">This blog post was written by Dedeepya Bonthu. This is a repost from her Medium article, approved by the author.
In stadiums, sports fans love to express themselves by cheering for their favorite teams, holding up placards and team logos. Emoji’s allow fans at home to rapidly express themselves,</p>
<footer class="gh-card-meta">
<!--
-->
<span class="gh-card-author">By ByteByteGo</span>
<time class="gh-card-date" datetime="2024-03-26">26 Mar 2024</time>
<!--
--></footer>
</div>
</a>
</article>
<article class="gh-card post">
<a class="gh-card-link" href="/brief-history-of-scaling-uber/">
<figure class="gh-card-image">
<img
srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w160/format/webp/2026/03/1704993859593.png 160w,
https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w320/format/webp/2026/03/1704993859593.png 320w,
https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w600/format/webp/2026/03/1704993859593.png 600w,
https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w960/format/webp/2026/03/1704993859593.png 960w,
https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w1200/format/webp/2026/03/1704993859593.png 1200w,
https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w2000/format/webp/2026/03/1704993859593.png 2000w"
sizes="320px"
src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w600/2026/03/1704993859593.png"
alt="Brief History of Scaling Uber"
loading="lazy"
>
</figure>
<div class="gh-card-wrapper">
<h3 class="gh-card-title is-title">Brief History of Scaling Uber</h3>
<p class="gh-card-excerpt is-body">This blog post was written by Josh Clemm, Senior Director of Engineering at Uber Eats. This is a repost from his LinkedIn article, approved by the author.
On a cold evening in Paris in 2008, Travis Kalanick and Garrett Camp couldn&#39;t get a cab. That&#39;s when</p>
<footer class="gh-card-meta">
<!--
-->
<span class="gh-card-author">By ByteByteGo</span>
<time class="gh-card-date" datetime="2024-03-14">14 Mar 2024</time>
<!--
--></footer>
</div>
</a>
</article>
<article class="gh-card post">
<a class="gh-card-link" href="/behind-aws-s3s-massive-scale/">
<figure class="gh-card-image">
<img
srcset="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w160/format/webp/2024/03/7.png 160w,
https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w320/format/webp/2024/03/7.png 320w,
https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w600/format/webp/2024/03/7.png 600w,
https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w960/format/webp/2024/03/7.png 960w,
https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w1200/format/webp/2024/03/7.png 1200w,
https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w2000/format/webp/2024/03/7.png 2000w"
sizes="320px"
src="https://storage.ghost.io/c/64/72/647204cd-7ad2-4539-b98a-3489074b932d/content/images/size/w600/2024/03/7.png"
alt="Behind AWS S3’s Massive Scale"
loading="lazy"
>
</figure>
<div class="gh-card-wrapper">
<h3 class="gh-card-title is-title">Behind AWS S3’s Massive Scale</h3>
<p class="gh-card-excerpt is-body">This is a guest article by Stanislav Kozlovski, an Apache Kafka Committer. If you would like to connect with Stanislav, you can do so on Twitter and LinkedIn.
AWS S3 is a service every engineer is familiar with.
It’s the service that popularized the notion of cold-storage to</p>
<footer class="gh-card-meta">
<!--
-->
<span class="gh-card-author">By ByteByteGo</span>
<time class="gh-card-date" datetime="2024-03-06">06 Mar 2024</time>
<!--
--></footer>
</div>
</a>
</article>
</div>
</div>
</section>
<footer class="gh-footer gh-outer">
<div class="gh-footer-inner gh-inner">
<section class="gh-footer-signup">
<h2 class="gh-footer-signup-header is-title">
High Scalability
</h2>
<p class="gh-footer-signup-subhead is-body">
Building bigger, faster, more reliable websites.
</p>
<form class="gh-form" data-members-form>
<input class="gh-form-input" id="footer-email" name="email" type="email" placeholder="jamie@example.com" required data-members-email>
<button class="gh-button" type="submit" aria-label="Subscribe">
<span><span>Subscribe</span> <svg xmlns="http://www.w3.org/2000/svg" width="32" height="32" fill="currentColor" viewBox="0 0 256 256"><path d="M224.49,136.49l-72,72a12,12,0,0,1-17-17L187,140H40a12,12,0,0,1,0-24H187L135.51,64.48a12,12,0,0,1,17-17l72,72A12,12,0,0,1,224.49,136.49Z"></path></svg></span>
<svg xmlns="http://www.w3.org/2000/svg" height="24" width="24" viewBox="0 0 24 24">
<g stroke-linecap="round" stroke-width="2" fill="currentColor" stroke="none" stroke-linejoin="round" class="nc-icon-wrapper">
<g class="nc-loop-dots-4-24-icon-o">
<circle cx="4" cy="12" r="3"></circle>
<circle cx="12" cy="12" r="3"></circle>
<circle cx="20" cy="12" r="3"></circle>
</g>
<style data-cap="butt">
.nc-loop-dots-4-24-icon-o{--animation-duration:0.8s}
.nc-loop-dots-4-24-icon-o *{opacity:.4;transform:scale(.75);animation:nc-loop-dots-4-anim var(--animation-duration) infinite}
.nc-loop-dots-4-24-icon-o :nth-child(1){transform-origin:4px 12px;animation-delay:-.3s;animation-delay:calc(var(--animation-duration)/-2.666)}
.nc-loop-dots-4-24-icon-o :nth-child(2){transform-origin:12px 12px;animation-delay:-.15s;animation-delay:calc(var(--animation-duration)/-5.333)}
.nc-loop-dots-4-24-icon-o :nth-child(3){transform-origin:20px 12px}
@keyframes nc-loop-dots-4-anim{0%,100%{opacity:.4;transform:scale(.75)}50%{opacity:1;transform:scale(1)}}
</style>
</g>
</svg> <span>Email sent</span>
</button>
<p data-members-error></p>
</form>
</section>
<div class="gh-social-links">
<a href="https://x.com/ghost" target="_blank" rel="noopener" aria-label="X">
<svg viewBox="0 0 24 24" fill="currentColor"><g><path d="M18.244 2.25h3.308l-7.227 8.26 8.502 11.24H16.17l-5.214-6.817L4.99 21.75H1.68l7.73-8.835L1.254 2.25H8.08l4.713 6.231zm-1.161 17.52h1.833L7.084 4.126H5.117z"></path></g></svg> </a>
<a href="https://www.facebook.com/ghost" target="_blank" rel="noopener" aria-label="Facebook">
<svg class="icon" viewBox="0 0 24 24" xmlns="http://www.w3.org/2000/svg" fill="currentColor"><path d="M23.9981 11.9991C23.9981 5.37216 18.626 0 11.9991 0C5.37216 0 0 5.37216 0 11.9991C0 17.9882 4.38789 22.9522 10.1242 23.8524V15.4676H7.07758V11.9991H10.1242V9.35553C10.1242 6.34826 11.9156 4.68714 14.6564 4.68714C15.9692 4.68714 17.3424 4.92149 17.3424 4.92149V7.87439H15.8294C14.3388 7.87439 13.8739 8.79933 13.8739 9.74824V11.9991H17.2018L16.6698 15.4676H13.8739V23.8524C19.6103 22.9522 23.9981 17.9882 23.9981 11.9991Z"/></svg> </a>
</div>
<div class="gh-footer-bar">
<span class="gh-footer-logo is-title">
High Scalability
</span>
<nav class="gh-footer-menu">
<ul class="nav">
<li class="nav-sign-up"><a href="#/portal/">Sign up</a></li>
</ul>
</nav>
<div class="gh-footer-copyright">
Powered by <a href="https://ghost.org/" target="_blank" rel="noopener">Ghost</a>
</div>
</div>
</div>
</footer>
</div>
<div class="pswp" tabindex="-1" role="dialog" aria-hidden="true">
<div class="pswp__bg"></div>
<div class="pswp__scroll-wrap">
<div class="pswp__container">
<div class="pswp__item"></div>
<div class="pswp__item"></div>
<div class="pswp__item"></div>
</div>
<div class="pswp__ui pswp__ui--hidden">
<div class="pswp__top-bar">
<div class="pswp__counter"></div>
<button class="pswp__button pswp__button--close" title="Close (Esc)"></button>
<button class="pswp__button pswp__button--share" title="Share"></button>
<button class="pswp__button pswp__button--fs" title="Toggle fullscreen"></button>
<button class="pswp__button pswp__button--zoom" title="Zoom in/out"></button>
<div class="pswp__preloader">
<div class="pswp__preloader__icn">
<div class="pswp__preloader__cut">
<div class="pswp__preloader__donut"></div>
</div>
</div>
</div>
</div>
<div class="pswp__share-modal pswp__share-modal--hidden pswp__single-tap">
<div class="pswp__share-tooltip"></div>
</div>
<button class="pswp__button pswp__button--arrow--left" title="Previous (arrow left)"></button>
<button class="pswp__button pswp__button--arrow--right" title="Next (arrow right)"></button>
<div class="pswp__caption">
<div class="pswp__caption__center"></div>
</div>
</div>
</div>
</div>
<script src="https://highscalability.com/assets/built/source.js?v=uMULaSMoSjM4gY5C"></script>
</body>
</html>