BigBFT: Scaling BFT without Compromising Fault Tolerance via State Sharding
Abstract
Scalability remains a major challenge for Byzantine fault tolerance (BFT) systems, whose throughput is often limited by sequential transaction processing at each node. Prior work has attempted to address this challenge by full sharding or by replacing total ordering with serializable concurrent execution, but these approaches introduce coordination overhead, complex protocols, or weakened security and isolation guarantees. In this work, we propose BigBFT, a BFT transactional store that scales by sharding the state but not the execution. BigBFT keeps a single consensus protocol to order transactions globally, while multiple node shards replicate state partitions and process state access in parallel. In this design, the consensus protocol circumvents complicated concurrency control in prior works. Unlike Basil and sharded blockchains, BigBFT scales without weakening fault tolerance or absence of Byzantine behavior. A preliminary prototype outperforms a fully replicated BFT store, and we expect the advantage to grow as the design is fully realized.