DAILYSPHERE MEDIA•Today's Briefing
Advertisement
Leaderboard (728×90 / 970×90) — Reserved Ad Space

Partition Finalization in Pinterest’s Next-Generation DB Ingestion Framework

Advertisement
In-Content (728×90 / 300×250) — Reserved Ad Space

Qianrui Zhang | Sr Software Engineer, Logging Platform
Kanchi Masalia | Software Engineer II, Stream Processing Platform
Liang Mou | Sr Staff Software Engineer, Logging Platform
Yi Pan | Principal Engineer, Agent Platform

Introduction

This is the third post in our series on Pinterest’s next-generation database ingestion framework. Part 1 introduced the DB ingestion framework built on Kafka, Flink, Spark, and Iceberg, and Part 2 covered automated schema evolution. This post tackles another challenge in migrating downstream customers to the new ingestion framework: knowing when data is complete enough to read.

We’ll walk through that data completeness challenge and how we solve it: what “partition finalization” means, why it matters to downstream customers, how we built a unified mechanism to generate partition finalization markers, and how those markers are consumed. We’ll also look at how this mechanism generalizes beyond DB ingestion.

Background & Motivation

Before explaining partition finalization, it helps to define what a partition is and why a downstream consumer cares when one partition is “done.”

What is a Partition

A time partition is a subset of a table whose rows are grouped by a time column, e.g. all rows whose timestamp falls within a given hour. Grouping data this way lets a consumer read just the slice it needs, such as the 1 AM hour, instead of scanning the whole table.

How that grouping is physically represented depends on the table format:

  • Hive uses explicit, directory-based partitioning: each partition maps to a physical directory in S3.
  • Iceberg uses hidden partitioning: the partitioning is defined in the metadata layer rather than by directory layout, so consumers don’t need to know the physical file layout to query a partition.

What Finalization Means

Now consider a typical consumer of these tables: a downstream batch job that processes a single hour’s partition and runs only once for that hour. The consumer faces the question: when should my job run?

The answer is straightforward: Because the job runs only once, it should run only after the target partition is ‘finalized’, i.e. data in that partition is less likely to change. If it runs too early and the partition is still being updated after the run, it will miss later updates in that partition. And to know when a partition is finalized, it relies on the producer to put some signals on the data. For example, the producer can emit a lightweight marker in the table location when it thinks the partition is finalized, and the consumer will use a sensor to check that marker.

Challenge: Finalization in the Stream-based DB Ingestion Framework

In the previous batch-based system, finalization was trivial. The data was produced by a single batch job that read the whole source DB and wrote it in one shot, so once that job finished and there is data in a partition, we can consider that as finalized and downstream jobs could start consuming.

The new framework is stream-based. A Flink job writes change events continuously, committing many small batches into a partition over time. Because a record’s event time (when the change happened in the source) can lag when we process it, a partition can keep receiving data after its wall-clock hour has passed and no single moment marks it “done”, and the presence of data in a partition no longer means the partition is complete.

So finalization is no longer a byproduct of a job finishing. We have to infer it from the stream itself uniformly across the thousands of pipelines the framework supports, and a new mechanism is needed to support this.

Our Solution: EventTime Based Partition Finalization

We introduced a partition finalization mechanism into the streaming layer of the pipeline. The same Flink job that writes the CDC table also measures how far its event time has progressed and publishes a finalization marker into the destination Iceberg table, where downstream jobs can wait on it via a sensor before they start consuming.

Architecture Overview

There is no change in the underlying data flow as in Part 1 (see below figure, blue components are added for partition finalization): Kafka carries CDC events, Flink writes them into the CDC Iceberg table, and Spark upserts into the base table. The partition finalization logic mainly happens inside the Flink-to-Iceberg sink, plus one piece of published metadata that can be carried over in the Iceberg table.

Partition Finalization Logic

We’ll start with how a Flink-to-Iceberg pipeline works in general. As Flink processes the event stream, it first writes incoming records into data files. Those files aren’t visible to readers right away; they become visible during Flink’s checkpoint process, which runs every few minutes (configurable) to persist the job’s progress for fault tolerance. Every checkpoint produces a new snapshot in the Iceberg table: an atomic commit that also carries a small summary of metadata describing it. It works much like a Git commit, with each snapshot recording a new visible state of the table on top of the last.

With this process in mind, our partition finalization logic executes in the following procedures:

Event time extraction

Everything starts from a single value per record: the CDC event’s event time. Each pipeline configures which column carries it, and an extractor reads that column and normalizes it to a timestamp. Because the extractor is pluggable, the same machinery works for any table by simply pointing it at that table’s event-time column.

Fact collection: event time statistics

Between checkpoints, we collect the event time of every record processed within that window and track statistics over them. The core facts are the minimum and maximum event time seen in that window. On top of those, we also capture percentiles, e.g. the 99th and 95th percentile of event time, so we retain the shape of the distribution rather than just its extremes.

Fact storage: Iceberg snapshot summary

When a checkpoint commits, these statistics are written into the Iceberg snapshot’s summary, the per-snapshot metadata Iceberg already maintains for each commit. Every commit therefore carries its own self-describing record of what event times it contained, with no external store to keep in sync.

To keep the footprint small, we store the percentiles in a t-digest, a compact sketch that approximates a distribution in a small, fixed amount of space. The summary stays around a few hundred bytes whether a checkpoint saw a thousand records or millions, which is what makes it affordable to attach to every commit. The sketch is also mergeable, so partial statistics computed in parallel combine into a single summary at checkpoint time.

From fact to opinion: partition finalization watermark

After each commit, the framework turns those facts (the per-commit event-time statistics, such as the minimum event time seen in each checkpoint) into an opinion: how far the partition is finalized. It runs an algorithm over the statistics from recent commits and produces a watermark that marks the point up to which the partition is considered complete. By default, it takes the minimum of the per-commit minimum event times across the last X commits, then rounds down to the hour boundary.

The watermark is kept monotonically increasing. If late-arriving data would pull it backward, the value holds instead of regressing, so a finalized hour stays finalized. When that happens, the application can also send an alert to the data consumer, flagging that late data arrived after the partition was finalized.

Opinion storage: Iceberg table property

Unlike the facts which live on individual snapshots, the opinion is a single table-level value, so we store the partition finalization watermark as an Iceberg table property. Downstream jobs read it with an ordinary metadata lookup and gate their processing on it, without the need of any side channel or separate service.

Finalization Metadata Propagation

Since the statistics and finalization watermarks all live in Iceberg metadata, they are easy to propagate. Finalization is computed on the CDC table, but many consumers read the base table produced by the Spark upsert job, which carries the data but not the signal.

We close that gap in the upsert job: on each run it reads the finalization watermark from the CDC table and writes it onto the base table in the same commit. The signal follows the data to where consumers read it, and because these are ordinary Iceberg properties, propagation is just copying metadata during a commit the pipeline already makes without the need of side channel or services. Downstream jobs can then query those table metadata via sensors.

Underlying Implementation: Extensible Iceberg-Flink Sink

The above partition finalization logic is built entirely on top of the Flink-to-Iceberg sink. We needed two capabilities the standard sink did not have: (1) a way to collect per-record metrics inside the writer and merge them at commit time, and (2) a way to run custom logic immediately after every successful Iceberg commit. Rather than hard-coding partition finalization into the sink, we introduced two generic extension points — Accumulators and CommitProcessor — that any Flink-to-Iceberg pipeline can use.

Operator Topology

The standard sink has two Flink operators in sequence: a writer (one instance per parallelism unit) that flushes records to data files, and a committer (single instance) that makes those files visible by appending a new snapshot to the Iceberg table. The extensible sink keeps exactly the same topology and adds no extra operators.

When either extension point is configured, the writer and committer swap in extensible variants that carry accumulator and commit processor logic alongside. When neither is configured, the pipeline falls back to the standard operators with zero overhead.

Accumulators: collecting facts per record

An accumulator is a lightweight aggregator that runs inside each writer subtask and builds up a summary of the records it has seen since the last checkpoint. Between checkpoints it records and aggregates each record, extracting a field value and updating running statistics such as a count, a min/max, or a custom metric. At checkpoint time, each subtask takes a snapshot of its accumulator’s current state and resets it for the next window.

Those per-subtask snapshots travel to the committer alongside the data-file manifests. The committer, which sees results from every subtask, combines the individual snapshots into a single merged result. Because the subtasks may complete in any order, the merge operation must be order-independent: two snapshots combined must produce the same result regardless of which is folded into which.

Once merged, the accumulator’s value is written into the Iceberg snapshot summary under a namespaced key, sitting alongside the other metadata Iceberg already records per commit.

Commit processor: acting after each commit

A commit processor adds two hooks around the Iceberg commit: one that runs just before the snapshot is written, and one that runs just after it is durable. The pre-commit hook can set additional snapshot properties or abort the commit by throwing, which causes Flink to retry the checkpoint. The post-commit hook is where partition finalization lives: it reads the statistics just written to the snapshot summary, runs the watermark algorithm, and updates the table property if the watermark advances.

Both hooks receive the same context: a handle to the live table, the checkpoint identifier, and the merged accumulators from all writer subtasks.

Putting it together

The event-time accumulator collects the minimum, maximum, and a compact percentile sketch of event times across all records in a checkpoint window. At commit time those per-subtask observations are merged into a single global summary and written into the snapshot. The commit processor then reads that summary, together with recent snapshot history, runs the non-regressing watermark algorithm, and writes the result back to the table, all within the commit cycle without any external service.

Configurable Finalization Policies

Because each summary keeps the full distribution rather than a single number, the field the commit processor reads becomes a knob each table can turn. This is where the completeness-versus-latency trade-off is handled.

Most tables use the conservative default. Tables that value freshness and can accept a rare late record opt into a percentile instead, accepting that a small tail of records may arrive after their hour was marked safe to read.

Generalizing Beyond DB Ingestion

Nothing above is specific to database ingestion. The accumulator and commit processor only need a way to read an event time from each record; everything else, e.g. the sketch, the windowed lookback, the non-regressing watermark are generic. Any Flink-to-Iceberg pipeline that needs a completeness signal can adopt partition finalization by plugging in its own event-time extractor. And we are actively working on adopting this partition finalization mechanism into Pinterest’s next-generation Stream Ingestion framework.

Conclusion & What’s Next

Moving to streaming ingestion gives us fresh data in minutes, but it costs us the implicit completeness that batch loads provided for free. Partition finalization mechanism restores that guarantee: a cheap, mergeable event-time sketch on every commit, a non-regressing “safe to read” watermark derived from it, and a policy knob to trade-off completeness against freshness, and all of those are reusable beyond our DB ingestion framework.

In the following post, we will cover the incremental processing logic built on top of this ingestion framework and how we use it to support downstream use cases efficiently. Stay tuned for the future post: Incremental Processing on CDC Pipeline.

Acknowledgments

Huge thanks to my teammates Yisheng Zhou, Vi Nguyen, and Artem Tetenkin for building the Next Generation DB Ingestion at Pinterest together.

This project would not have been possible without the significant contributions and support of the following partners:

  • Storage Services: Leonardo Marques Maciel Silva, Yuan Gao, Liqi Yi
  • Storage Foundations: Tailin Lyu, Yu Su, Istvan Podor, John Grass
  • Streaming Processing Platform: Kevin Browne
  • Batch Processing Platform: Carlos Benavides
  • Big Data Storage: Pucheng Yang, Mian Luo, Jenny Wang
  • Data Warehouse: Bryant Xiao, Anita Narra, Pritam Pan
  • Logging: Vahid Hashemian, Jeff Xiang, Jesus Zuniga

Special gratitude goes to Shardul Jewalikar, Ang Zhang, and Roger Wang for their continuous guidance, feedback, and support throughout the project.


Partition Finalization in Pinterest’s Next-Generation DB Ingestion Framework was originally published in Pinterest Engineering Blog on Medium, where people are continuing the conversation by highlighting and responding to this story.