397 lines
17 KiB
HTML
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 · 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/">← 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’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’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’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’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’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’t already serve a replica of the partition. If this approach fails,
|
|
we will look for the first node that doesn’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'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 "shardz"
|
|
|
|
server, err := shardz.NewServer(shardz.ServerConfig{
|
|
Manager: "shardz.coordinator.local:7557", // shardz coordinator endpoint
|
|
Node: "worker-node1.local:9090",
|
|
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("/shardz", 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’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’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’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>
|