pinterest3 min read

Curated summary

Next Generation DB Ingestion at Pinterest

Read original(opens in new tab)

Pinterest replaced fragmented, batch-oriented database ingestion with a unified Change Data Capture (CDC) framework. The new architecture uses Debezium/TiCDC, Kafka, Flink, Spark, and Iceberg to process only changed records, reducing latency from over 24 hours to minutes while lowering infrastructure costs. It also provides native row-level deletion, scalable operations, and improved compliance.

Problems with the Legacy System

  • Batch workflows often delayed updates by more than 24 hours.
  • Full-table processing was inefficient because many tables changed by less than 5% each day.
  • Lack of row-level deletion support complicated data compliance.
  • Multiple independently maintained pipelines created operational complexity and inconsistent data quality.

Unified CDC-Based Architecture

  • Supports MySQL, TiDB, and KVStore.
  • Captures database changes through a generic CDC service and publishes them to Kafka, typically in under one second.
  • Flink processes events in near real time and stores them in append-only CDC Iceberg tables on S3.
  • Spark jobs run periodically—often every 15 minutes—to merge recent changes into base Iceberg tables.
  • A bootstrap pipeline initializes base tables from historical database dumps.
  • Maintenance jobs handle compaction and snapshot expiration.
  • The framework is designed for at-least-once processing, petabyte-scale data, thousands of pipelines, and YAML-based configuration.

CDC Tables and Base Tables

  • CDC tables act as time-series ledgers containing every change event.
  • CDC data typically becomes available within five minutes.
  • Base tables mirror the current state of the source database while retaining historical records.
  • Base-table latency is generally between 15 minutes and one hour.

Upserting Changes into Base Tables

  • Spark first identifies the newest event for each primary key.
  • Events are ranked by timestamp and GTID, then deduplicated.
  • Iceberg’s MERGE INTO applies the resulting changes:
    • Deletes matching records when the event represents a deletion.
    • Updates existing records.
    • Inserts new records unless the event is a deletion.
  • The process uses a recent CDC window and a processing watermark to avoid reprocessing unnecessary data.

Choosing Merge-on-Read

  • Pinterest standardized on Iceberg’s Merge-on-Read (MOR) strategy.
  • Copy-on-Write (COW) was rejected for most workloads because:
    • It requires more computation during writes.
    • It produces substantially larger replacement files, increasing storage costs.
  • MOR better balances update performance and storage efficiency for frequent incremental changes.

Partitioning for Faster Upserts

  • Large base tables can be partitioned using a hash bucket of the primary key.
  • For example, bucket(100, id) distributes records across 100 partitions.
  • This allows Spark to process partitions in parallel and reduces the data scanned or rewritten during merges.
  • Iceberg tables are configured with format version 2, identifier fields, merge-on-read update and delete modes, and target file sizes.

Small-File Challenge

  • Bucketing improved parallelism but caused each upsert to generate many small files within partitions.
  • The article indicates that Pinterest investigated this bottleneck and introduced further optimizations, though the supplied excerpt ends before describing them.

Pinterest’s CDC-based design provides a substantially faster and more efficient alternative to full-table batch ingestion. Teams adopting a similar system should combine incremental CDC processing with partitioning, merge-on-read storage, bootstrapping, and ongoing file-maintenance strategies.

Continue with another curated summary.