Curated summary
Next Generation DB Ingestion at Pinterest
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 INTOapplies 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.
Related reading
Continue with another curated summary.
Migrating Data Ingestion Systems at Meta Scale
Read originalFrom Hive to Iceberg: The Secret to 12x Faster Data Reflection
Read originalAWS Weekly Roundup: Price reduction of GPT models in Bedrock, CloudWatch managed collectors for Prometheus metrics, and more (August 3, 2026) | Amazon Web Services
Read originalAWS Weekly Roundup: Agentic CX designer for Amazon Connect Customer, EC2 AMI Watermarks, Open Governance for MySQL, and more (June 29, 2026) | Amazon Web Services
Read original