Introduction
Shard Manager is a new service. It runs the sharding of Cadence Matching at Uber. It is still pre version 1, so breaking changes may occur.
Shard Manager manages the shards of an application, such as Cadence Matching. It owns the shard-to-owner mapping, and load balances on top of it.
Why Shard Manager
Cadence historically used Ringpop to manage shards. Ringpop is a traditional consistent hash ring. While mathematically and theoretically beautiful, consistent hashing gives a set of problems:
- Load Balancing - A ring balances by shard count, never by shard load. Moving a shard means every host changing its view of the ring at the same instant, and nothing coordinates that. Methods such as virtual nodes make it possible to shuffle the shards, but do not allow fine-grained control.
- Graceful Handovers - When a shard moves in a traditional consistent hash ring, requests will simultaneously go to both owners, causing availability drops. With a centralized Shard Manager we let the old owner drain before the new one takes over. It also opens the door to warming caches on the new owner before the handover, and truly zero-downtime transfers. Graceful handover is not implemented yet.
- Debuggability and Introspection - In a traditional consistent hash ring the state of the shard assignments is spread across all participants, and there is no authority on what the correct state is. A single misbehaving instance can therefore cause issues across the cluster, and shard assignment issues are extremely hard to debug. With central assignment the mapping is a record we can inspect with
smctl, the Shard Manager CLI. - Operability - A ring cannot be steered. There is no way to take a single hot shard off a host, or to empty a host ahead of a deploy, without changing the membership list itself.
Shard Manager solves all of these problems by centralizing the assignment decision: balancing by real load, introspection with smctl, and operability such as draining shards and hosts.
Running it
Shard Manager needs etcd. The repository ships a compose file:
docker compose -f docker/github_actions/docker-compose.yml up -d etcd
make bins
./shard-manager-server start --services shard-distributor
It reads config/development.yaml by default, which points at localhost:2379, defines a few development namespaces, and serves gRPC on port 7943.
To watch it do something, run the canary. It starts executors and a pinger against a fixed and an ephemeral namespace. It is also a good example of how integration with Shard Manager looks:
make start-shard-manager-canary
Where to go next
- Architecture. How the service, etcd, leader election, and rebalancing fit together.
- cadence-workflow/shard-manager. The source, the client libraries, and
smctl.