Menu
The Pragmatic Engineer·September 30, 2026

Distributed Databases and Storage: Insights from Cockroach Labs CTO Peter Mattis

This article is an interview with Peter Mattis, co-founder and CTO of Cockroach Labs, discussing his journey from Google (Gmail, Colossus) to building CockroachDB. It offers deep insights into the architectural decisions, trade-offs, and underlying data structures like B-trees and LSM-trees crucial for building scalable, reliable, and consistent distributed databases.

Read original on The Pragmatic Engineer

The Ubiquitous B-Tree in System Design

Peter Mattis highlights the pervasive nature of B-trees in storage and database systems. From early Gmail's storage layer for email threads to internal Google data structures that outperformed `std::map` and even CockroachDB's range index, B-trees demonstrate their efficiency for indexed data access. Their spatial locality advantages often make them faster and more memory-efficient than other tree structures like red-black trees, especially in disk-based systems where I/O patterns are critical.

Evolution of Google's Distributed Storage: GFS to Colossus

Google's journey in distributed file systems evolved from GFS to Colossus. GFS initially used three full data replicas for redundancy, which was resource-intensive. Colossus, on the other hand, pioneered the use of Reed-Solomon erasure coding. This innovative approach allowed data to be stored twice (effectively, with parity blocks), increasing redundancy while significantly reducing storage overhead by 33%. This is a classic example of a system design trade-off where computational overhead for encoding/decoding is accepted for improved storage efficiency and fault tolerance.

Designing for Consistency and Durability: Spanner and CockroachDB

The design of CockroachDB was heavily inspired by Google's Spanner, which itself built upon Colossus. A key architectural decision in distributed databases is achieving consensus. Peter explains that at least three replicas are necessary for a distributed database to recover reliably after a node crash, as two replicas introduce ambiguity in determining the last successful write. CockroachDB defaults to three replicas and can scale up to five for critical system tables, balancing latency and fault tolerance.

ℹ️

B-trees vs. LSM-trees

Colossus's append-only nature posed a challenge for B-tree based databases, as B-trees require in-place updates. This constraint led to the adoption of Log-structured merge trees (LSM-trees) for systems like Spanner built on top of Colossus. LSM-trees are highly optimized for write-heavy workloads and append-only storage, where writes are logged sequentially for durability, and data is merged asynchronously in the background. This choice reflects a fundamental trade-off between read and write optimization in storage engine design.

From RocksDB to Pebble: CockroachDB's Storage Engine Evolution

CockroachDB initially leveraged RocksDB, a fork of Google's LevelDB (which popularized LSM-trees). Peter Mattis later developed and open-sourced Pebble, which became CockroachDB's current storage engine. This highlights the continuous evolution and optimization required in core database components to meet specific performance and operational demands of a distributed system.

distributed databaseB-treeLSM-treeCockroachDBGoogle SpannerGoogle Colossusconsistencyfault tolerance

Comments

Loading comments...
Distributed Databases and Storage: Insights from Cockroach Labs CTO Peter Mattis | SysDesAi