netflix

High-Throughput Graph Abstraction at Netflix: Part I (opens in new tab)

Netflix’s Graph Abstraction is designed for OLTP graph workloads requiring millions of operations per second and millisecond-level latency, rather than open-ended analytical exploration. Built on existing Netflix abstractions, it supports real-time and optional historical graph views while handling nearly 10 million operations per second across 650 TB of data. Its core design emphasizes strong schemas, efficient traversal planning, low-latency caching, and controlled trade-offs such as eventual consistency and bounded query depth.

OLTP Graph Use Cases

  • Netflix distinguishes between:
    • OLAP workloads, which prioritize large-scale exploration using RDF/SPARQL, property graphs, Gremlin, openCypher, or SQL.
    • OLTP workloads, which require extremely high throughput, low latency, and global availability.
  • OLTP queries may restrict traversal starting points, depth, or complexity to meet performance goals.
  • Key applications include:
    • Real-Time Distributed Graph, modeling dynamic relationships and interactions across Netflix.
    • Social Graph, supporting social connections in Netflix Gaming.
    • Service Topology, enabling real-time and historical analysis of internal services during incidents.

Architecture and Netflix Data Abstractions

  • The Graph Abstraction builds on existing platform components rather than implementing storage and caching independently.
  • Key-Value (KV) Abstraction provides the latest state of nodes and edges and serves as the real-time index.
  • TimeSeries (TS) Abstraction can be added for historical graph views.
  • EVCache delivers low-millisecond latency, with additional specialized caching layers under experimentation.
  • The Data Gateway Control Plane manages:
    • Graph schemas
    • Dataset provisioning and deletion
    • KV and TS configuration

Property Graph Model

  • Graphs contain typed nodes and edges, each with associated properties.
  • Properties are strongly typed to support:
    • Efficient filtering
    • Consistent data exports
    • Validation during writes
  • Edges may be:
    • Unidirectional, representing one-way relationships
    • Bidirectional, representing relationships traversable in both directions

Namespaces and Provisioning

  • Data is isolated into logical units called namespaces.
  • Each namespace maps to a physical storage layer and may use dedicated or shared hardware.
  • Provisioning automation selects an appropriate hardware configuration based on:
    • Required throughput
    • Latency targets
    • Dataset size
    • Workload criticality

Graph Schema and Query Optimization

  • Every namespace has an explicit schema defining:
    • Node and edge types
    • Valid properties and their types
    • Allowed relationships
    • Edge directions
  • Schemas are represented through edge mappings, such as an account owns profile relationship or a bidirectional profile linked_to device relationship.
  • Property definitions can specify types such as TIMESTAMP and STRING.
  • Servers load schemas into an in-memory metadata graph, enabling:
    • Rejection of invalid nodes, edges, and properties
    • Faster traversal-path planning
    • Deduplication of bidirectional edge traversals
    • Removal of impossible paths and incompatible filters
  • Servers periodically poll the Control Plane so schema changes are reflected without requiring manual updates.
  • Planned improvements include:
    • Using edge cardinality to reduce query fanout
    • Generating type-safe data-access layers
    • Making the Gremlin-like API schema-aware

Real-Time Indexing with Key-Value Storage

  • KV stores the real-time representation of all graph nodes and edges.
  • Each namespace corresponds to a table, partitioned into records by unique IDs.
  • Records contain multiple sorted key-value items, effectively forming a map of sorted maps.
  • Writes to the same ID and key are idempotent, allowing safe retries and request hedging.
  • KV uses timestamp-based tokens to enforce Last-Write-Wins (LWW) semantics.
  • The post begins discussing the two-tier partitioning strategy for node storage, but the provided content ends before that design is explained.

Netflix’s approach demonstrates that high-throughput graph serving depends on specialized constraints and platform integration rather than unrestricted graph querying. Strong schemas, bounded traversals, KV-based indexing, automated provisioning, and low-latency caching together provide a practical foundation for production-scale OLTP graph workloads.