Building a Next-Generation Key-Value Store at Airbnb (opens in new tab)
Airbnb rebuilt Mussel, its key-value store for derived data, from a complex EC2-based system into a cloud-native NewSQL platform. Mussel v2 combines bulk ingestion, streaming writes, low-latency reads, flexible consistency, and automated operations while supporting more than 100 existing use cases. A gradual, reversible blue/green migration moved production workloads without data loss or customer-visible downtime.
Why Airbnb Rebuilt Mussel
- New use cases—including real-time fraud detection, personalization, and dynamic pricing—required both streaming updates and large-scale bulk ingestion.
- Mussel v1 had become difficult to operate and scale:
- Node changes required multi-step Chef scripts on EC2.
- Static hash partitioning created hotspots and latency spikes.
- Consistency options were limited.
- Resource consumption and costs were difficult to track.
- Mussel v2 provides Kubernetes-based automation, dynamic range sharding, configurable consistency, namespace tenancy, quotas, and usage dashboards.
Mussel v2 Architecture
Stateless Dispatcher
- A horizontally scalable Kubernetes service translates client requests into backend queries and mutations.
- It supports:
- Dual writes and shadow reads during migration
- Retries, rate limiting, and dynamic throttling
- Service-mesh security and discovery
- Point lookups, range queries, prefix queries, and low-latency stale reads
- Each dataname maps to a logical table, simplifying access patterns.
Kafka-Based Write Pipeline
- Writes are first persisted to Kafka for durability.
- The Replayer and Write Dispatcher apply them to the backend in order.
- Kafka absorbs traffic bursts and supports consistency, migrations, bootstrapping, and upgrades.
- Airbnb plans to eventually rely more directly on the distributed database for ingestion and replication to reduce latency and operational complexity.
Bulk Loading
- Mussel retains support for both:
- Merge jobs, which add data to existing tables
- Replace jobs, which swap in a new dataset
- Existing Airflow onboarding workflows transform warehouse data into a standard format and upload it to S3.
- A stateless controller coordinates ingestion, while Kubernetes StatefulSet workers load data in parallel.
- Deduplication, delta merges, and insert-on-duplicate-key-ignore improve throughput and reduce unnecessary writes.
Scalable Data Expiration
- Mussel v1 depended on storage-engine compaction for TTL expiration, which became inefficient at scale.
- V2 uses a topology-aware expiration service:
- Namespaces are divided into range-based subtasks.
- Multiple workers scan and delete expired records concurrently.
- Scheduling limits interference with live queries.
- Max-version enforcement and targeted deletes help manage write-heavy tables.
- The result is faster, more visible, and more scalable retention management.
Blue/Green Migration
- The migration had to handle massive datasets, thousands of tables, and mission-critical traffic with zero data loss and no availability impact.
- Because v1 lacked table-level snapshots and CDC, Airbnb built a custom migration pipeline.
- Tables were selected and migrated individually according to usage and risk.
Migration Stages
- Blue: All production traffic continued serving from v1.
- Shadowing: Bootstrapped v2 tables processed parallel reads and writes, but v1 still served responses.
- Reverse: V2 served live traffic while v1 remained available as a fallback.
- Cutover: After validation, traffic was permanently moved to v2 one dataname at a time.
- Automatic circuit breakers and fallback logic enabled rapid rollback if v2 showed errors or replication lag.
- Kafka’s replication stream maintained eventual consistency between the two systems throughout the transition.
Practical Takeaway
Mussel v2 demonstrates that large datastore rearchitectures can be made safe through incremental migration, durable event logs, shadow traffic, and reversible per-table cutovers. The key recommendation is to combine a more scalable backend with strong operational automation and migration tooling, rather than attempting a single disruptive replacement.