Files
nexus/sreweekly/articles/299/08-shardz.html
2026-09-12 17:23:01 +08:00

397 lines
17 KiB
HTML

<!DOCTYPE html>
<html lang="en">
<head>
<meta charset="UTF-8">
<meta name="viewport" content="width=device-width, initial-scale=1.0">
<title>Shardz &middot; rakyll.org</title>
<link rel="apple-touch-icon-precomposed" sizes="144x144" href="/apple-touch-icon-144-precomposed.png">
<link rel="shortcut icon" type="image/png" href="/favicon.png">
<link href="https://rakyll.org/index.xml" rel="alternate" type="application/rss+xml" title="rakyll.org" />
<link rel="preconnect" href="https://fonts.googleapis.com">
<link rel="preconnect" href="https://fonts.gstatic.com" crossorigin>
<link href="https://fonts.googleapis.com/css2?family=JetBrains+Mono:wght@400;700&display=swap" rel="stylesheet">
<style>
* {
margin: 0;
padding: 0;
box-sizing: border-box;
}
:root {
--green: #2b2f33;
--green-dim: #c4c8cc;
--text: #2b2f33;
--text-dim: #6b7075;
--bg: #f4f6f4;
--panel: #eceef0;
}
html {
scroll-behavior: smooth;
}
body {
background: var(--bg);
color: var(--text);
font-family: 'JetBrains Mono', ui-monospace, 'SF Mono', Menlo, Consolas, monospace;
line-height: 1.7;
min-height: 100vh;
background-image:
linear-gradient(rgba(43, 47, 51, 0.05) 1px, transparent 1px),
linear-gradient(90deg, rgba(43, 47, 51, 0.05) 1px, transparent 1px);
background-size: 32px 32px;
}
.topnav a {
position: fixed;
top: 20px;
z-index: 100;
color: var(--green);
text-decoration: none;
padding: 0.5rem 0.9rem;
border: 1px solid var(--green);
transition: all 0.3s;
display: inline-flex;
align-items: center;
gap: 0.5rem;
background: rgba(255, 255, 255, 0.7);
}
.topnav a.home { left: 20px; }
.topnav a.gh { right: 20px; padding: 0.5rem; }
.topnav a:hover {
background: var(--green);
color: var(--bg);
}
.topnav .gh svg { width: 22px; height: 22px; fill: var(--green); transition: fill 0.3s; }
.topnav .gh:hover svg { fill: var(--bg); }
.content {
position: relative;
z-index: 10;
max-width: 760px;
margin: 0 auto;
padding: 6.5rem 1.5rem 5rem;
}
.post {
padding: 0;
}
.post-date {
display: block;
color: var(--text-dim);
font-size: 0.9rem;
margin-bottom: 2rem;
}
.post-date::before { content: "// "; }
.post h1 {
font-size: 2.4rem;
line-height: 1.2;
color: var(--green);
margin-bottom: 0.75rem;
}
.post h2, .post h3, .post h4 {
color: var(--green);
margin: 2.2rem 0 0.8rem;
line-height: 1.3;
}
.post h2 { font-size: 1.5rem; }
.post h3 { font-size: 1.2rem; }
.post h2::before, .post h3::before { content: "# "; color: var(--text-dim); }
.post p, .post ul, .post ol, .post blockquote { margin-bottom: 1.3rem; }
.post ul, .post ol { padding-left: 1.5rem; }
.post li { margin-bottom: 0.4rem; }
.post a {
color: var(--green);
text-decoration: none;
border-bottom: 1px solid var(--green-dim);
transition: all 0.2s;
}
.post a:hover {
background: var(--green);
color: var(--bg);
border-color: var(--green);
}
.post strong { color: #1a1d20; }
.post em { color: var(--text); }
.post blockquote {
border-left: 2px solid var(--green-dim);
padding-left: 1rem;
color: var(--text-dim);
font-style: italic;
}
.post code {
font-family: 'JetBrains Mono', ui-monospace, 'SF Mono', Menlo, Consolas, monospace;
color: var(--green);
font-size: 0.92em;
}
.post pre {
background: var(--panel);
border-left: 3px solid var(--green);
padding: 1.1rem 1.2rem;
overflow-x: auto;
margin-bottom: 1.5rem;
font-size: 0.9rem;
line-height: 1.55;
}
.post pre code {
background: none;
border: none;
padding: 0;
color: var(--text);
font-size: inherit;
}
.post img { max-width: 100%; height: auto; padding: 10px; background: #fff; border: 1px solid var(--green-dim); }
.post hr {
border: none;
border-top: 1px dashed var(--green-dim);
margin: 2rem 0;
}
.post table {
width: 100%;
border-collapse: collapse;
margin-bottom: 1.5rem;
}
.post th, .post td {
border: 1px solid var(--green-dim);
padding: 0.5rem 0.75rem;
text-align: left;
}
.post th { color: var(--green); }
.post-footer {
margin-top: 2.5rem;
padding-top: 1.5rem;
border-top: 1px dashed var(--green-dim);
color: var(--text-dim);
font-size: 0.9rem;
}
.post-footer a { color: var(--green); text-decoration: none; }
.post-footer a:hover { text-decoration: underline; }
.cursor {
display: inline-block;
width: 0.6rem;
height: 1rem;
background: var(--green);
margin-left: 0.2rem;
vertical-align: middle;
animation: blink 1s step-end infinite;
}
@keyframes blink { 50% { opacity: 0; } }
@media (max-width: 600px) {
.content { padding: 5.5rem 1rem 3rem; }
.post { padding: 1.5rem 1.25rem 2rem; }
.post h1 { font-size: 1.8rem; }
}
</style>
</head>
<body>
<nav class="topnav">
<a class="home" href="https://rakyll.org/">&larr; home</a>
<a class="gh" href="https://github.com/rakyll" aria-label="GitHub">
<svg viewBox="0 0 16 16" version="1.1" aria-hidden="true">
<path d="M8 0c4.42 0 8 3.58 8 8a8.013 8.013 0 0 1-5.45 7.59c-.4.08-.55-.17-.55-.38 0-.27.01-1.13.01-2.2 0-.75-.25-1.23-.54-1.48 1.78-.2 3.65-.88 3.65-3.95 0-.88-.31-1.59-.82-2.15.08-.2.36-1.02-.08-2.12 0 0-.67-.22-2.2.82-.64-.18-1.32-.27-2-.27-.68 0-1.36.09-2 .27-1.53-1.03-2.2-.82-2.2-.82-.44 1.1-.16 1.92-.08 2.12-.51.56-.82 1.28-.82 2.15 0 3.06 1.86 3.75 3.64 3.95-.23.2-.44.55-.51 1.07-.46.21-1.61.55-2.33-.66-.15-.24-.6-.83-1.23-.82-.67.01-.27.38.01.53.34.19.73.9.82 1.13.16.45.68 1.31 2.69.94 0 .67.01 1.3.01 1.49 0 .21-.15.45-.55.38A7.995 7.995 0 0 1 0 8c0-4.42 3.58-8 8-8Z"></path>
</svg>
</a>
</nav>
<main class="content">
<article class="post">
<h1>Shardz</h1>
<span class="post-date">Fri, Dec 3, 2021</span>
<p>Shard coordination has been one of the bigger challenges to design
sharded systems especially for engineers with little experience in the subject.
Companies like Facebook have been using general purpose shard coordinators,
e.g. <a href="https://engineering.fb.com/2020/08/24/production-engineering/scaling-services-with-shard-manager/">Shard Manager</a>, and suggesting that general
purpose sharding policies have been widely successful. A general
purpose sharding coordinator is not a solution to all advanced sharding needs,
but a starting point for average use cases.
It&rsquo;s a framework to think about core abstractions
in sharded systems and providing a protocol that orchestrate sharding
decisions. In the last few weeks, I&rsquo;ve been working on an extendible general purpose
shard coordinator, Shardz. In this article, I will explain
the main concepts and the future work.</p>
<h2 id="sharding">Sharding</h2>
<p>Sharding is a concept of sharding load to different nodes in a cluster
to horizontally scale systems. We can talk about two commonly used approaches
to sharding strategies.</p>
<h3 id="hashing">Hashing</h3>
<p>Hashing is used to map an arbitrary key to a node in a cluster. Keys and
hash functions are chosen carefully for load balancing to be effective.
With this approach, each incoming key is hashed and its modulo is
calculated to associate the key with one of the available nodes. Then,
the request is routed to that node to be served.</p>
<pre tabindex="0"><code>hash_function(key) % len(nodes)
</code></pre><p>This approach, even though its shortcomings, offer an easy to implement way
of designating a location to the incoming key. Consistent hashing can improve
the excessive migration of shards upon node failure or addition.
This approach is still commonly used in databases,
key/value caches, and web/RPC servers.</p>
<p>One of the significant shortcomings of hashing is that the destination node is determined
when handling an incoming request. The nodes in the system don&rsquo;t know which
shards they are serving in advance unless they can calculate the key ranges they need to serve. This situation makes
it hard for nodes to preload data or it makes it harder to implement
strict replication policies. The other difficulty is that existing shards
cannot be easily visualized for debugging purposes unless you can query
where key ranges live.</p>
<h3 id="lookup-tables">Lookup tables</h3>
<p>An alternative to hashing is to generate lookup tables. In this approach,
you keep track of available nodes and partitions you want to serve. Lookup
tables are regenerated each time a node or a partition is added or removed,
and is the source of the truth which node a partition should be served at.
Lookup tables can also help enforcing replication policies, e.g. ensure there
are at least two replicas being served at different nodes.</p>
<p>Lookup tables are generally easy to implement when you have homogenous nodes.
In heterogenous cases, node capacity can be used to
influence the distribution of partitions.
Shardz chose to only tackle the homogenous nodes for now but it&rsquo;s trivial to
advance the coordinator to consider nodes with different capacity in the future.</p>
<h2 id="p-x-n">P x N</h2>
<p>Sharding with lookup tables is essentially a P x N problem where
you have P partitions and N nodes. The goal is to distribute P
partitions on the available nodes at any time. Most systems require
replication for fault tolerance, so we decided to identify a P with
a unique identifier and its replica number. At any times, a P should be
replicated on Ns based on how many replicas are required.</p>
<pre tabindex="0"><code>type P struct { ID string, Replica int }
type N struct { Host string }
</code></pre><p><img src="/img/shardz-pxn.png" alt="An illustration of PxN"></p>
<p>We expect users to manage partitions separately and talk to the Shardz
coordinator to schedule them on the available nodes. The partitions
should be uniquely identified within the same coordinator.</p>
<h2 id="sharder">Sharder</h2>
<p>A core concept in Shardz is the Sharder interface. A sharder allows
users to register/unregister Ps and Ns. Then, users can query where
a partition is being served. Sharder interface can be implemented with
hashing or lookup tables. We don&rsquo;t enforce any implementation details
but expect the following interface to be satisfied.</p>
<pre tabindex="0"><code>type Sharder interface {
RegisterNodes(n ...N) error
UnregisterNodes(n ...N) error
RegisterPartitions(ids ...string) error
UnregisterPartitions(ids ...string) error
Partitions(n N) ([]P, error)
Nodes(partitionID string) ([]N, error)
}
</code></pre><p>Sharders can be extended to satisfy custom needs such as implementing
custom replication policies, scheduling replicas in different availability
zones, finding the correct VM type/size to schedule the partition.</p>
<h2 id="work-stealing-replica-sharder">Work stealing replica sharder</h2>
<p>Shardz comes with general use Sharder implementations. The default sharder
is a work stealing sharder that enforce a minimum number of replicas for
each partition. ReplicaSharder implements a trivial approach to find
the next available node by looking at the least overloaded node that
doesn&rsquo;t already serve a replica of the partition. If this approach fails,
we will look for the first node that doesn&rsquo;t already serve a replica of the partition.
If everything fails, we pick a random node and schedule the partition on it.</p>
<pre tabindex="0"><code>func nextNode(p P) {
// Find the node with the lowest load.
// If the node is not already serving p.ID, return the node.
// Find all the available nodes.
// Randomly loop for a while until you find a node
// that doesn&#39;t serve p.ID (a replica).
// Return the node if something is found in acceptable number of iterations.
// Randomly pick a node and return.
}
</code></pre><p>Additional to the partition positioning, work stealing is triggered when
a new node joins or a partition is unregistered. Work stealing can be triggered
as many as times possible, e.g. routinely to avoid unbalanced load on nodes.</p>
<pre tabindex="0"><code>func stealWork() {
// Calculate the average load on nodes.
// Find nodes with 1.3x more load than average.
// Calculate how many partitions need to be removed
// from the overloaded node. Remove the partitions.
// For each removed partition, call nextNode to find
// a new node to serve the partition.
}
</code></pre><h2 id="protocol">Protocol</h2>
<p>Shardz have been heavily influenced by <a href="https://engineering.fb.com/2020/08/24/production-engineering/scaling-services-with-shard-manager/">Shard Manager</a> when it comes to making it easier for servers to report their status and coordinate with the manager.</p>
<p><img src="/img/shardz-coordinator.png" alt="An illustration of PxN"></p>
<p>A worker node creates a server that implements the protocol. ServeFn and
DeleteFn are the hooks when coordinator notifies the worker which partitions it should
serve or stop serving.</p>
<pre tabindex="0"><code>import &#34;shardz&#34;
server, err := shardz.NewServer(shardz.ServerConfig{
Manager: &#34;shardz.coordinator.local:7557&#34;, // shardz coordinator endpoint
Node: &#34;worker-node1.local:9090&#34;,
ServeFn: func(partition string) {
// Do work to serve partition.
},
DeleteFn: func(partition string) {
// Do work to stop serving partition.
},
})
if err != nil {
log.Fatal(err)
}
http.HandleFunc(&#34;/shardz&#34;, server.Handler)
log.Fatal(http.ListenAndServe(listen, nil))
</code></pre><p>At startup, worker node automatically pings the coordinator to register itself.
Then Shardz coordinator will ping the worker back periodically to check its status.
The partitions that need to be served by the worker is periodically distributed
and the ServeFn and DeleteFn functions are triggered automatically if
partitions changed.</p>
<p>At a graceful shutdown, worker node automatically reports that it&rsquo;s going
away and give the coordinator a change to redistribute its partitions.</p>
<h2 id="fault-tolerance">Fault tolerance</h2>
<p>Shardz is designed to run in a clustered mode where there will
be multiple replicas of the coordinator at any time. The coordinator will have
a single leader that is responsible for sharding decision and propagating
them to others. If leader goes away, another replica becomes the leader.
ZooKeeper is used to coordinate the Shardz coordinators.</p>
<h2 id="future-work">Future work</h2>
<p>Shardz is still in the early stages but it&rsquo;s a promising concept
to build a reusable multi purpose sharding coordinator. It has potential
to lower to entry barrier to design sharded systems. The next steps
for the project:</p>
<ul>
<li>Fault tolerance. The coordinator still has work to be finished when running
in the cluster mode.</li>
<li>New Sharder implementations. Project aims to provide multiple
implementation with different policies to meet the needs. I desire to open
source this project to allow community to contribute.</li>
<li>Language support for protocol implementation. We only have support
for Go and should expand the server implementation to capture more languages
to make it easy just to import a library to add Shardz support to any process.</li>
<li>Visualization and management tools. Dashboards and custom control planes
to monitor, debug and manage shards.</li>
<li>An ecosystem that speaks Shardz. It&rsquo;s an ambitious goal but a unified
control plane for shard management would benefit the entire industry.</li>
</ul>
<div class="post-footer">
<a href="https://rakyll.org/archive/">cd ~</a><span class="cursor"></span>
</div>
</article>
</main>
</body>
</html>