Pinterest’s Data Platform at Scale

▶ Watch (0:03)

Pinterest processes hundreds of petabytes of data daily across thousands of scheduled batch workloads. The Mocha platform sits at the center of this, built entirely on EKS. Spark and Flink handle batch and streaming workloads. Presto, Trino, and Spark SQL serve the query layer. Apache Iceberg acts as the metadata layer over S3, providing a standardized table format and fast snapshot switching. Spark alone consumes about one-third of Pinterest’s entire compute budget.

Multi-Tier Queue Management and Custom Schedulers

▶ Watch (0:36)

Mocha runs three job tiers on shared clusters. Tier-one jobs have guaranteed resources set to their historical usage, preventing preemption. Tier-two guarantees are set at 80% of historical usage. Tier-three runs best-effort. The queue structure mirrors Pinterest’s org chart: five top-level orgs, each with projects underneath, each project with three tier queues. For Spark, Pinterest uses Unicorn, a YARN-compatible scheduler with fair sharing and a max-application cap per queue. PyTorch and Ray workloads use Volcano, which adds topology-aware scheduling across availability zones and co-location controls for training and serving jobs.

The Tides Problem: Idle Reserved Capacity

▶ Watch (10:36)

Pinterest reserves EC2 capacity from AWS for online services to get guaranteed supply and long-term pricing discounts. When online traffic drops during off-peak hours, that reserved capacity sits idle. Running Spark workloads directly on the same shared cluster as online services hit EKS control plane and API server scaling limits. The chosen solution was to rebalance EC2 instances from online clusters to dedicated offline Spark clusters, keeping full network and disk isolation between the two environments.

The Global Scheduler and Tides Controller Architecture

▶ Watch (13:14)

The Tides controller watches ODCR capacity metrics. When online traffic falls, it scales up one of three Tides cluster types per availability zone. When traffic rises, it scales them back down. The global scheduler queries the job submission database, monitors Tides cluster capacity, and routes tier-two and tier-three jobs to Tides while keeping tier-one jobs on the static Mocha clusters. A global holding mechanism queues jobs when capacity is tight or when a job is unlikely to finish before the next scale-down event.

Protecting Jobs During Scale-Down

▶ Watch (16:19)

A Spark driver pod is a single point of failure. If it is evicted, the entire job fails. Tides addresses this with two separate node groups per cluster. The static node group hosts only driver pods and never scales down. The dynamic node group holds executor pods and is controlled by the Tides controller. Remote shuffle storage on a separate cluster means executor evictions do not corrupt shuffle data, so jobs survive scale-down events without losing intermediate results.

Automatic Routing and the Path to 90% Utilization

▶ Watch (20:20)

Asking users to opt in with a tides=true flag failed because users had no strong motivation to do so. Pinterest replaced opt-in with automatic routing inside the global scheduler. It selects non-critical jobs based on priority, historical resource usage, and wait time, then routes them to Tides without user involvement. A weighted score also reorders jobs when capacity returns, replacing strict FIFO. Current Tides cluster utilization sits at 70%. The team plans to add an ML model to predict available online-service capacity, with a target of 90%.

Notable Quotes

Spark itself is used about onethird of the entire Pinterest computer resource. So it it cost a lot. Rainie Li · ▶ 05:40

user doesn’t have the strong motivations to onboard to tight cluster but from the platform’s um perspective we still want to uh save the infra cost. Rainie Li · ▶ 20:43

we are able to uh reach tight cluster utilization to se uh 70% so far Rainie Li · ▶ 21:37

Hopefully we can reach out to 90% with this approach. Rainie Li · ▶ 22:55

Key Takeaways

  • Spark consumes one-third of Pinterest’s compute budget, making scheduling decisions expensive.
  • Separate static and dynamic node groups keep driver pods alive during Tides scale-down events.
  • Automatic job routing replaced opt-in flags and pushed Tides utilization to 70%.

About the Speaker(s)

Rainie Li is a Senior Engineering Manager at Pinterest, leading the Data Processing Infrastructure and Big Data Storage Platform teams. Before Pinterest she held engineering roles at Microsoft Azure and AWS Rekognition. She holds a Master’s degree in Electrical and Computer Engineering from Duke University.

Ang Zhang is Senior Director of Engineering at Pinterest and leads the big data platform organization.