Interview conceptSystem Design

Distributed Query, Storage, And Metadata Systems

Asked of: Software Engineer

Last updated

What's being tested

Candidates must show practical distributed-systems engineering: designing scalable storage and query architectures, managing strong metadata and consistency, and trading off performance vs. cost. Interviewers probe how you partition data, maintain correctness (snapshots/transactions), surface metadata for planning, and propagate updates/evictions in a multi-tenant environment. Snowflake cares because these decisions determine ruler-scale query latency, storage efficiency, and operational complexity.

Core knowledge

  • Content-addressable storage (CAS): store chunks by fingerprint (e.g., SHA-256) to enable deduplication; collisions negligible if fingerprint length n≥256n\geq256, collision prob ≈1/2n\approx 1/2^n.

  • Chunking strategies: fixed-size vs. variable-size (content-defined) chunking; variable chunking (e.g., Rabin fingerprints) improves cross-file deduplication but increases chunking CPU and metadata cardinality.

  • Index & metadata services: separate metadata service that tracks fingerprints → locations, reference counts, and object manifests; must be highly-available and low-latency (cache hot prefixes).

  • Reference management & GC: choose reference counting (immediate reclamation, distributed counters) or periodic garbage collection (mark-sweep with leases); reference counting scales poorly at extreme concurrency without sharding.

  • Consistency & transactions: provide atomic commit of manifests + metadata update (use two-phase commit, or single-writer leases + idempotent operations) and support MVCC/snapshot isolation for queries reading stable views.

  • Deduplication metadata scale: expect metadata entries ≈ number of unique chunks; for 1PB with 8KB avg chunk, ~1.3e8 chunks → store compacted bloom filters or partitioned key-value indices (e.g., LSM-based stores).

  • Data layout & locality: use consistent hashing to distribute chunks across storage nodes; colocate index shards with compute cache to reduce lookup latency. Replicate metadata via consensus (Raft) for durability.

  • Network & RPC design: use multiplexed RPC (gRPC) and batched lookups for manifest resolution; pipeline chunk fetches to hide latency. Limit synchronous cross-node ops on hot paths.

  • Cache & eviction: multi-layer caching: client-side result cache, cluster SSD layer, and cold object store (S3); eviction policies combine recency and reference count/usage cost.

  • Materialized view maintenance: for DAG-based views, store versioned intermediate results, use incrementality when possible, and invalidate downstream caches via a dependency graph (topological propagation).

  • Query execution: split planning (cost-based optimizer) and execution (distributed operator graph) with shuffle-aware partitioning; model network as constrained resource and aim to minimize data movement.

  • Cost & durability tradeoffs: choose erasure coding for storage cost-efficiency at higher CPU on reads; replication for low-latency reads and simpler GC.

Worked example — Design an object store with deduplication

Frame: ask expected SLAs (read/write latency), workload mix (many small files vs. large blobs), durability target (nines), storage backend (S3 vs raw disks), and per-tenant isolation requirements. Skeleton pillars: (1) chunking & fingerprinting pipeline that computes SHA-256 per chunk, (2) metadata index mapping fingerprints → locations + reference counts, (3) write path that atomically persists chunks and updates manifests, (4) read path that reconstructs objects from chunk pointers with caching, and (5) GC/compaction to reclaim unreferenced chunks. Explicit tradeoff: choosing reference counting gives prompt reclamation but requires strongly-consistent distributed counters (sharded updates + batching), whereas periodic GC (lease-based) simplifies writes at the cost of delayed reclaim and higher storage footprint. Close by noting operational concerns: monitor fingerprint index cardinality, handle hot-chunk hotspots (replicate or cache), and if time permitted add: cross-tenant dedupe policy, chunk-size auto-tuning, and multi-level metadata caches.

A second angle — Design cache for DAG-based query views

Same primitives apply but framed around compute results rather than raw bytes: represent each view node as a versioned artifact with a manifest linking to underlying shard files (which themselves may reuse deduplicated chunks). Key differences: invalidation is driven by upstream data changes and query semantics (append-only vs. arbitrary updates), so prefer incremental maintenance and dependency tracking. Use a dependency DAG to propagate invalidation and schedule recomputation; maintain both fine-grained (per-partition) and coarse-grained (full-view) artifacts to trade recompute cost vs. storage. For correctness, ensure readers obtain a consistent snapshot of the DAG—either via MVCC or via commit tokens—so cached intermediate artifacts map to the correct logical view version. Finally, optimize for small, frequent updates with delta encoding and a write-optimized store for intermediates.

Common pitfalls

Pitfall: confusing deduplication savings with real user-visible cost — dedupe reduces stored bytes but increases CPU, metadata size, and lookup latency; always quantify net cost versus naive replication.

Many candidates propose global, synchronous reference counting without addressing write contention. A distributed counter requires sharding and batching; otherwise metadata becomes the performance bottleneck under high ingest.

Pitfall: over-indexing metadata — storing per-chunk global indexes in memory without partitioning leads to OOM and GC pressure; instead partition and use compact probabilistic filters for fast negatives.

A communication mistake is skipping SLA clarifications: don't design a system assuming synchronous semantics if the product tolerates eventual visibility. State your consistency model and how it affects both user-visible correctness and operational complexity.

Pitfall: ignoring failure modes during compaction/GC — naive delete-on-zero can remove chunks concurrently referenced by a delayed writer; prefer atomic manifest reference updates and lease-based GC windows to avoid data loss.

Depth mistake: optimizing for average-case reading cost while neglecting pathological workloads (hot files or small-random-writes). Call out hot-spot mitigation (replication, LRU hot-cache, backpressure).

Connections

Interviewers may pivot to adjacent system topics such as query optimizers (cost models, cardinality estimation), storage engines (LSM vs B-tree, compaction strategies), or distributed transaction systems (two-phase commit, distributed snapshot algorithms). Be ready to map design choices across those domains.

Further reading

Practice questions

Related concepts