How Dragonboat's Multi-Group Raft Implementation Scales to Thousands of Concurrent Raft Shards

Dragonboat scales to thousands of concurrent Raft shards by isolating each shard as an independent node while sharing transport, LogDB, and worker pools across the NodeHost process to minimize per-shard overhead.

Dragonboat is an open-source Go library (lni/dragonboat) that implements the multi-group Raft consensus algorithm, allowing a single process to host thousands of independent Raft groups (shards) simultaneously. Unlike traditional Raft deployments that run one group per process, Dragonboat's architecture deliberately minimizes per-shard resource consumption through strategic resource sharing and lock-free data structures. This design enables a single NodeHost instance to manage over 10,000 concurrent shards on modest hardware without exhausting file descriptors, goroutines, or memory.

Architectural Foundation: Shard Isolation with Shared Resources

At the core of Dragonboat's scaling capability is the multi-group Raft model, where each shard operates as an independent Raft replica set within a shared NodeHost process. In nodehost.go, the NodeHost struct maintains a sync.Map called mu.shards (lines 71-73) that stores each shard's node object.

This map provides lock-free lookups for the hot request path. When routing requests, NodeHost.getShard performs an RLock-protected read against this map, enabling concurrent access to thousands of shards without contention on a single mutex. Each shard maintains its own Raft state machine, proposal queue, and timeout tick independently, ensuring that computational work for one shard never blocks another.

Shared Transport Layer

All shards within a NodeHost share a single transport layer instance defined in nodehost.go (lines 84-85) as transport.ITransport. Rather than opening separate TCP sockets for each shard, Dragonboat pools connections at the transport level.

This design means a node serving thousands of shards maintains a constant number of network connections relative to its peers, not proportional to the shard count. The shared transport eliminates the socket exhaustion that would otherwise occur with thousands of independent Raft processes.

LogDB Multiplexing and Sharding Strategies

Dragonboat provides two complementary strategies to manage on-disk storage for thousands of shards without exhausting file descriptors.

Multiplexed Log Storage

By default, Dragonboat uses the tan storage engine which implements multiplexed log databases. In internal/tan/db_keeper.go (lines 84-107), the multiplexedKeeper groups shards by shardID % 16 (the default group size), assigning each group to a single TAN database instance.

As implemented in internal/tan/collection.go (lines 31-44), this multiplexing reduces file descriptor usage from one per shard to one per group, while maintaining logical isolation between shard logs. This allows thousands of shards to share a minimal set of on-disk log files.

Configurable LogDB Sharding

For deployments requiring physical isolation, Dragonboat supports the LogDBShardCount configuration option in config.Config. When set greater than zero, engine.go initializes a sharded LogDB via internal/logdb/sharded.go, mapping ranges of shard IDs to independent LogDB instances.

This flexibility allows operators to balance between the extreme efficiency of multiplexed storage (suitable for most workloads) and the isolation of dedicated storage instances when needed.

Worker Pool and Goroutine Management

Dragonboat avoids the "goroutine explosion" common in high-concurrency systems by using fixed-size worker pools rather than per-shard goroutines. In engine.go (lines 150-176), the engine initializes worker pools sized proportionally to CPU cores, not shard counts.

Raft message handling, snapshot generation, and background compaction tasks are dispatched to these shared pools. This ensures that adding thousands of shards increases memory usage only for state storage, not for thread management, keeping the process lightweight and scheduler-friendly.

Lock-Free Request Processing

To minimize allocation churn and contention on the request path, Dragonboat implements configurable request pools. In nodehost.go (lines 90-92), the requestPools field holds a slice of sync.Pool instances, with the default count matching the number of CPU cores (settings.Soft.NodeHostRequestStatePoolShards).

Each pool contains reusable RequestState objects, eliminating per-request allocations and ensuring the hot path remains constant-time regardless of shard count. This pooling strategy significantly reduces GC pressure when processing proposals across thousands of active shards.

Batched Snapshots and Compaction

Snapshot operations, which are I/O-intensive, are optimized through asynchronous batching. The snapshotter.go implementation triggers snapshots per-shard, but engine.go (lines 550-620) runs a snapshotWorker that debounces and merges concurrent snapshot requests.

When thousands of shards request snapshots simultaneously, the worker coalesces these operations, preventing I/O spikes and ensuring storage bandwidth remains sustainable. Similarly, log compaction runs asynchronously through the shared worker pool, avoiding thundering herd problems.

Practical Implementation: Starting Thousands of Shards

The following example demonstrates creating a NodeHost configured for high-density shard deployment:

// Create a NodeHost configured for thousands of shards
cfg := config.NodeHostConfig{
    WALDir:        "/var/lib/dragonboat/wal",
    NodeHostDir:   "/var/lib/dragonboat/nodehost",
    // Configure 4000 LogDB shards for physical isolation, or 0 for multiplexed mode
    LogDBShardCount: 4000,
    RaftAddress:   "10.0.0.1:63001",
}
nh, err := dragonboat.NewNodeHost(cfg)
if err != nil {
    panic(err)
}

// Start 1000 independent Raft shards
for i := uint64(0); i < 1000; i++ {
    rc := config.Config{
        ShardID:   i,
        ReplicaID: 1,
    }
    smFactory := func(shardID, replicaID uint64) sm.IStateMachine {
        return mykv.New()
    }
    if err := nh.StartReplica(nil, false, smFactory, rc); err != nil {
        panic(err)
    }
}

Because the NodeHost internally utilizes the multiplexed TAN DB, shared transport, and fixed worker pools, this loop creates thousands of shards without requiring OS-level tuning for file descriptors or thread limits.

Summary

  • Shard isolation: Each shard maintains independent state in mu.shards (a sync.Map in nodehost.go), eliminating cross-shard contention while enabling lock-free lookups via getShard.
  • Shared resources: Transport (ITransport), worker pools (engine.go), and LogDB instances are shared across all shards, keeping resource usage proportional to CPU cores rather than shard count.
  • Storage multiplexing: The default multiplexedKeeper (internal/tan/db_keeper.go) groups shards by ID modulo 16, reducing file descriptors from thousands to a handful.
  • Configurable sharding: The LogDBShardCount option allows physical isolation when needed, routing shards to distinct LogDB instances via internal/logdb/sharded.go.
  • Request pooling: requestPools (nodehost.go) uses sync.Pool per CPU core to eliminate allocations and maintain constant-time request processing.
  • Asynchronous batching: Snapshot and compaction workers (engine.go) merge concurrent operations to prevent I/O saturation across thousands of shards.

Frequently Asked Questions

How many Raft shards can a single Dragonboat NodeHost support?

A single NodeHost can support 10,000 or more concurrent shards on a 64-core machine with modest SSD storage. The limiting factors are typically disk I/O bandwidth and memory for state machine storage, not the Raft implementation itself, due to the shared worker pools and transport layers that keep per-shard overhead to a few megabytes.

What is the difference between multiplexed and sharded LogDB configurations?

The multiplexed mode (default) uses multiplexedKeeper in internal/tan/db_keeper.go to group multiple shards into single database files, minimizing file descriptors. The sharded mode (enabled via LogDBShardCount > 0) creates independent LogDB instances per shard range through internal/logdb/sharded.go, providing physical isolation at the cost of increased file descriptors. Multiplexed mode is preferred for high-density deployments.

Does each shard run its own goroutine?

No. Dragonboat uses fixed-size worker pools initialized in engine.go (lines 150-176) sized to CPU core count. Raft message processing, snapshots, and compactions are handled by these shared pools rather than per-shard goroutines, preventing scheduler overload when running thousands of shards.

How does Dragonboat prevent allocation churn when processing thousands of concurrent proposals?

Dragonboat implements request pooling via the requestPools field in nodehost.go (lines 90-92), which maintains a slice of sync.Pool objects defaulting to the number of CPU cores. These pools reuse RequestState objects, ensuring proposal processing requires no new allocations in the hot path and maintains constant-time performance regardless of shard count.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →