Referential Integrity Across Two Databases at Exabyte Scale

▶ Watch (0:30)

Google Photos holds trillions of objects and must keep a relational metadata database perfectly in sync with an object storage system. Every object referenced by metadata must exist in object storage. Every unreferenced object must be cleared. That constraint is a foreign key, but it spans two separate databases at trillions-of-objects scale. No off-the-shelf tooling handles it. The team built a high-concurrency pipeline to scan both sides and verify consistency, covering both data-loss detection and garbage collection.

Why 16 Parallel Partitions Still Broke the SLA

▶ Watch (4:02)

The team split the user namespace into 16 equally sized partitions and ran 16 independent integrity pipelines. The design promised balanced, predictable runtimes. After deployment, runtimes were not predictable. Audits showed less than 1% data variation across all 16 partitions. The data was genuinely uniform, but uniform data did not produce uniform performance. Every partition showed high variance, with some tails breaking the 3-day SLA. The join stage and post-processing both checked out. The bottleneck was the read stages, with P95 latencies over 2 seconds and P99 latencies over 10 seconds.

Stragglers: The Math Behind Why One Slow Shard Governs Everything

▶ Watch (13:00)

A batch pipeline’s runtime is set by its slowest shard. With 1,000 workers and 1,000 shards, even 1% of shards running slow forces every other worker to sit idle. The fix is over-sharding. Research puts the ideal shard count at 3 x n x log(n), where n is the number of workers. For 100 workers that is roughly 1,300 shards. The team found they requested 450k shards per partition from the object storage system but received only 30k, leaving large unsplittable blocks to dominate runtime.

Liquid Sharding and the Cost of Cheap Compute

▶ Watch (16:30)

The root cause was a legacy MapReduce bridge that could not split shards further. The fix required modernizing the reader to support liquid sharding, a technique where the framework splits oversized work chunks at runtime when workers go idle. This cut hours of tail latency down to minutes. A second problem compounded the first: the pipeline consumed hundreds of thousands of CPU units and nearly a petabyte of RAM on best-effort instances. Heavy shards got preempted frequently, forcing restarts. The waste ratio in the read stages exceeded 70%.

Three Engineering Truths from Operating at This Scale

▶ Watch (19:47)

Kakkad distilled three conclusions. First, P99 tail latency governs the entire pipeline’s end time, so systems must be architected around the worst case, not the average. Second, software partition boundaries are not enough. Infrastructure alignment matters, and physical storage characteristics will override logical designs. Third, best-effort compute is not cheap when preemption rates are high. If work chunks take longer than the preemption window, the cost gains from spot pricing are outweighed by lost and repeated work.

Q&A

Could increasing the number of partitions beyond 16 make individual jobs small enough to reduce preemption losses? More partitions increase worst-case end-to-end completion time and maintenance complexity, so the team chose to fix the shard splits inside each partition instead. ▶ 22:44

Why not run the whole thing as one batch job instead of 16 partitions? The underlying batch framework hit hard limits on shuffle temporary files at full scale, making the multi-partition split a hard requirement, not a preference. ▶ 27:00

Notable Quotes

uniform data does not equal uniform performance. Yashraj Kakkad · ▶ 06:47

our P95 latencies were like over 2 seconds. Our P99 latencies were like over 10 seconds. Yashraj Kakkad · ▶ 08:44

this this can actually take like hours of tail latency down to minutes. Yashraj Kakkad · ▶ 17:41

cheap compute is not really cheap. Yashraj Kakkad · ▶ 19:31

Key Takeaways

  • P99 tail latency, not average latency, sets the runtime of any batch pipeline.
  • Over-shard by a factor of 3 x n x log(n) to prevent stragglers from dominating.
  • Best-effort compute adds cost when preemptions force heavy shards to restart from scratch.

About the Speaker(s)

Yashraj Kakkad is a Software Engineer at Google, building and operating petabyte-scale storage infrastructure for Google Photos. With a background in distributed systems and data-intensive applications, he focuses on performance, reliability, and debugging the strange failure modes that appear at extreme scale.