Airbnb/k8s

3 posts

airbnb

Safeguarding Dynamic Configuration Changes at Scale (opens in new tab)

Airbnb’s Sitar platform is designed to make runtime configuration changes as safe and reliable as code deployments. It combines Git-based reviews, automated validation, staged rollouts, observability, and fast rollback with a highly available distribution system. Separating decision-making from config delivery, while using local caches, lets teams change behavior quickly without unnecessarily increasing outage risk. ## Requirements for a Modern Configuration Platform - Provides an end-to-end workflow for defining, reviewing, testing, and deploying configuration. - Treats configuration like code: - Versioned and reviewable - Auditable - Governed by ownership and access controls - Supports isolated local and canary testing before production rollout. - Accommodates multiple tenants with different: - Deployment triggers - Guardrails - Rollout strategies - Enables incident responders to make emergency changes while preserving auditability and visibility into who changed what, when, and which users or services were affected. ## Sitar’s Architecture Sitar consists of four major layers: - **Developer-facing layer:** Configs are usually managed through GitHub pull requests. The Sitar portal supports exceptions and administrative operations, including emergency deployments. - **Control plane:** Validates schemas, enforces ownership and authorization, selects rollout targets, manages progressive deployment, and supports rollback and targeted testing. - **Data plane:** Stores config values and versions as the source of truth, then distributes updates reliably and efficiently. - **Agents and client libraries:** An agent sidecar fetches subscribed configs and maintains a local cache. In-process client libraries read from that cache and expose values to application code, with optional fallbacks. A typical change moves from a Git workflow through validation and rollout decisions, into the data plane, and finally to sidecars and application clients. ## Git-Based Configuration Management - GitHub is the default interface because it integrates with Airbnb’s existing CI/CD systems and review practices. - Teams can use pull requests, mandatory reviewers, approval flows, and complete change history. - Related configs are grouped into tenants with defined owners, custom tests, and dedicated continuous-delivery pipelines. - The Sitar portal remains available for teams that need a UI or for urgent changes that must bypass the standard CI/CD process. ## Progressive Rollouts and Rollbacks - CI first checks schema correctness, expected structure, types, and other automated requirements. - Config changes require review and approval before deployment. - After merging, changes roll out gradually: - Start with a limited environment, AWS zone, or percentage of Kubernetes pods. - Evaluate the change at each stage. - Expand only when results are healthy. - Authors and stakeholders are notified when regressions are detected, and bad changes can be rolled back quickly. - Limiting the initial scope reduces the blast radius of configuration errors. ## Separating Control and Data Planes - The control plane decides whether and how a change should be deployed. - The data plane stores and distributes the resulting configuration. - This separation allows rollout policies and authorization logic to evolve independently from storage and delivery infrastructure. - Changes to one layer are less likely to disrupt the other. ## Local Caching and Resilient Clients - Each service runs an agent sidecar alongside its application container. - The sidecar periodically retrieves subscribed configs and persists them locally. - Client libraries read configuration from the local cache for fast, in-process access. - If the configuration backend becomes unavailable or degraded, services can continue using the last known good values. ## Practical Takeaway A reliable dynamic configuration system should combine code-like governance with runtime flexibility. Git reviews, validation, staged deployment, strong observability, plane separation, and local caching allow teams to respond quickly while keeping configuration failures contained and reversible.

airbnb

From Static Rate Limiting to Adaptive Traffic Management in Airbnb’s Key-Value Store (opens in new tab)

Airbnb evolved Mussel’s QoS system from static, per-client QPS limits into adaptive traffic management designed to maximize goodput. The newer approach accounts for the actual cost of requests, prioritizes critical workloads under stress, and detects hot keys or attack traffic before they overwhelm storage. Together, resource-aware quotas and real-time load shedding provide stronger protection against traffic spikes, uneven workloads, and DDoS-like bursts. ## Why Static QPS Limits Fell Short - Mussel is a multi-tenant key-value store serving millions of point and range reads across Airbnb. - Its original Redis-backed limiter assigned each client a fixed requests-per-second quota. - Requests exceeding the quota received HTTP 429 responses. - This model worked when backend effort roughly matched request count. - As usage grew, it could not account for: - The difference between a cheap one-row lookup and a 100,000-row scan. - Hot keys accessed by many clients simultaneously. - Localized storage-shard overload that affected unrelated traffic. - Sudden events such as bot floods, DDoS attacks, or large uploads. ## Resource-Aware Rate Control - Mussel replaced raw request counting with request units (RU), which represent estimated backend work. - RU calculations incorporate: - Fixed per-request overhead. - Rows and payload bytes processed. - Request latency, which distinguishes cached operations from disk-heavy ones. - The system uses calibrated linear formulas for reads and writes, with weights based on compute, network, and disk-I/O measurements. - Dispatchers debit a local token bucket according to each request’s RU cost rather than charging every request equally. - Periodic RU refills preserve simple, static quotas while making them more proportional to actual resource consumption. - Requests are rejected with HTTP 419 when the RU bucket is exhausted. - Load shedding remains separate, allowing latency-based protection to react dynamically without changing the underlying quota-refill mechanism. ## Load Shedding Under Sudden Stress - RU rate limiting smooths normal traffic but may react too slowly to rapidly changing workloads. - Mussel adds a load-shedding layer based on: - Traffic criticality. - A real-time latency ratio. - A CoDel-inspired queue-management policy. - Each dispatcher compares long-term p95 latency with short-term p95 latency. - A ratio near 1.0 indicates stable performance; a drop toward 0.3 signals rapidly increasing latency. - When stress crosses the threshold: - The system raises the effective RU cost for a designated lower-priority client class. - That class’s token bucket drains faster, causing its traffic to back off. - If conditions worsen, the penalty expands to additional classes. - Critical workloads, such as customer support and trust-and-safety traffic, can remain responsive while less important traffic is reduced. - The latency estimate uses the constant-memory P² algorithm, avoiding raw sample storage and cross-node coordination. ## Hot-Key Detection and DDoS Protection - Client-level quotas cannot prevent overload when many clients request the same popular key. - Mussel therefore detects skewed access patterns in real time. - When duplicate requests target a hot key, the system can protect storage by: - Serving responses from cache. - Coalescing identical requests before they reach the backend. - This approach protects the underlying shard whether the traffic comes from legitimate popularity, automation, or a DDoS burst. Mussel’s experience suggests that mature multi-tenant services should move beyond fixed QPS limits. Combining resource-based accounting, priority-aware load shedding, and hot-key mitigation provides a more effective way to preserve reliability while maximizing useful work during unpredictable traffic conditions.

airbnb

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.