Skip to content

SPOTIFY 2026-07-27

Read original ↗

Indexing the Data Lake for Online Point Queries

Summary

Spotify describes Random Access Parquet (RAP), a technique for serving fast point queries by key directly against the Parquet files already sitting in its GCS data lake — without copying that data into a key-value serving store like Bigtable or DynamoDB. Petabytes live in Bigtable for online use; exabytes sit in the lake. Object storage itself is now fast (GCS 30–100ms/request; S3 Express One Zone and GCS Rapid Storage at single-digit-millisecond latency), so the bottleneck for a single-row lookup is no longer the storage medium but the distributed SQL engines (Trino, BigQuery) whose job-scheduling and query-planning overhead adds seconds. RAP replaces scan-and-plan with an external index that maps a key directly to the exact file(s) and row(s), plus precise ranged reads that fetch only the needed bytes. It operates on the same Parquet files used by batch analytics, ML pipelines, and notebooks — "store once, pay once" — collapsing the choice between analytical and interactive access.

Key takeaways

  1. Point queries over the lake are a shared primitive. Online portals paginating a user's listening history and AI agents retrieving a user's data to reason over both need the same thing: fast lookup by key over datasets too large to economically keep resident in a KV store. (point-query-over-data-lake)
  2. The bottleneck moved from storage to the query engine. Cloud object storage now delivers single-digit-to-low-tens-of-ms latency, but Trino / BigQuery add seconds of scheduling + planning even for one row — they are built for analytical throughput, not interactive point lookups. (concepts/oltp-vs-olap)
  3. The chain of dependent reads is the fundamental bottleneck. Finding one user's row in a Parquet file requires a chain: fetch footer → parse row-group metadata → scan the key column → use column/page indexes to locate value pages. Each link costs one round-trip (latency) and reads bytes only to discover where to read next (bandwidth). This structure is the same at every storage tier; only the absolute per-round-trip cost changes. (dependent-read-chain)
  4. An external index collapses the chain into an O(1) lookup. RAP's index maps each key to the exact file and row numbers; the reader resolves rows to page locations via cached file metadata and issues a small number of precise ranged reads in parallel — no dependent load chain. (external-index, external-index-over-existing-parquet)
  5. The index is a compact append-only multimap. Each entry is {key, file (dictionary-encoded ordinal), row numbers, value count (optional, for pagination)}. A single key can appear across many files and partitions. Rule of thumb: indexing terabytes → gigabytes of index; petabytes → terabytes. Large indexes distribute by hash bucketing. New data appends new index fragments; existing fragments are never modified.
  6. Definitive, not probabilistic. Unlike Parquet's built-in PageIndex or Bloom filters — which are probabilistic and only narrow a scan — the external index is definitive: given a key it returns the exact files and rows, eliminating the scan entirely. (concepts/predicate-pushdown contrast.)
  7. RAP works on unmodified Parquet, but write-time preparation makes reads smaller. On stock files, RAP still reads the entire page containing the target row (a 4MB page to extract 100 bytes). Preparation optimizations fall into three buckets — concentrate a key's data, reduce bytes per read, reduce number of reads — and shift Parquet layout tradeoffs: properties that help in-file discovery (fine-grained page indexes, small pages, dictionary encoding for pushdown) matter less; properties that minimize the final read (fewer round-trips, smaller fetches, contiguous data) matter more.
  8. Reducing read count is the single biggest win. Fetching N columns for a key is N parallel reads (N× request capacity, pushes per-key latency toward the tail). Store point-query fields as one blob/Variant column (one read), or interleave columns physically so a RAP reader issues a single contiguous ranged read spanning all columns for a key — while a conventional reader still sees valid Parquet. (column-interleaving-for-single-read)
  9. Covering index / hoisted values → zero storage reads. Since the index builder visits every row at build time, it can hoist small values directly into the index entry (a covering index) and precompute per-key aggregates (event counts, total play time), available at index-lookup speed and enabling predicate pushdown at the index level. (concepts/secondary-index)
  10. Secondary indexes are a serving-layer decision. Multiple lookup dimensions (e.g. buyer_id and seller_id) are supported by building multiple access structures over the same index entries — hash tables for O(1) exact lookup, sorted indexes for range queries — with no pipeline changes and no data rewriting. Space-filling curves (Z-order, Hilbert) are complementary, improving file-layout locality for secondary dimensions.
  11. The economic point: cost of a point query drops to the cost of a cloud storage read. Historical data, long-tail entities, and low-traffic features — previously excluded from KV stores on per-GB cost grounds — all become viable for interactive access. (store-once-serve-both-analytical-and-interactive)

Operational numbers & concrete figures

  • Storage scale: petabytes in Bigtable (online) vs exabytes in the GCS data lake.
  • Object-storage latency: GCS 30–100ms/request; S3 Express One Zone and GCS Rapid Storage at single-digit ms.
  • Example query: "what was I listening to last summer?" → ~90 days of listening history at ~1,000 files/day = 90,000 Parquet files.
  • Candidate-set reduction: key-bucketing (1,000 buckets/day) cuts 90,000 → 90; Bloom filters narrow 90 → ~12 days the user was active.
  • Index sizing rule of thumb: TB indexed → GB of index; PB indexed → TB of index.
  • Unprepared read amplification: a 4MB page read to extract 100 bytes.
  • One-page-per-key overhead: ~20 bytes of page-header overhead per key boundary (negligible with substantial data per key).
  • Storage alignment: ZSTD skippable frames pad to 4KB or 16KB block boundaries to avoid straddle reads.
  • End state: point query reduced to a single ranged read of a few kilobytes, or the storage read eliminated entirely (covering index).

Systems / concepts / patterns extracted

Caveats

  • RAP as described is a Spotify-internal technique/approach; the post does not disclose whether it is open-sourced or the specific index storage backend.
  • Write-time preparation trades analytical flexibility for point-query speed (per-optimization tradeoffs are enumerated in the article's summary table): blob/Variant columns lose per-field pruning; interleaving inflates I/O for single-column scans; ZSTD frame resets rule out delta/RLE encodings (PLAIN/dictionary only); one-page-per-key grows the PageIndex proportional to key count.
  • Figures like "1,000 files/day" are illustrative examples in the post, not necessarily Spotify's exact production listening-history layout.

Source

Last updated · 766 distilled / 2,225 read