apache paimon is a lake table format built for the case every other format treats as an afterthought: a firehose of updates landing in real time and being read, moments later, as a correct changelog. It began life as Flink Table Store, graduated to a top-level Apache project, and its defining choice is architectural — under each table sits a log-structured merge (LSM) tree on object storage, so a high-frequency stream of upserts writes cheaply and a background compaction quietly merges the pieces. You get ACID lake tables on S3, HDFS, or OSS, but tuned for the streaming ingest that Iceberg and Hudi bolt on later.
That is a genuinely different shape from the two formats data engineers usually reach for. Apache Iceberg is snapshot-first: a metadata tree that excels at large batch scans and treats row-level updates as merge-on-read delete files. Apache Hudi is upsert-first but Spark-centric, built around a record index and a copy-on-write / merge-on-read split. Paimon is streaming-first: primary-key tables that accept CDC from Flink, merge engines that decide how two records for the same key combine, and changelog producers that emit a clean +I / -U / +U / -D stream for the next job downstream. This guide walks the four ideas an interviewer will actually probe — the table types and LSM layout, the merge engines, the changelog producers, and snapshots with time travel and tags — and pairs each with a Solution-Tail interview answer: code, a step-by-step trace, an output table, then a concept-by-concept breakdown of why it works.
When you want hands-on reps immediately after reading, drill the streaming practice library →, rehearse the layout decisions on the partitioning practice set →, and sharpen your upsert logic on the merge practice set →.
On this page
- Why Apache Paimon changes the lakehouse in 2026
- Table types & the LSM architecture
- Merge engines — deduplicate, partial-update, aggregation, first-row
- Changelog producers & streaming CDC
- Snapshots, compaction, time travel & tags
- Cheat sheet — Paimon recipes
- Frequently asked questions
- Practice on PipeCode
1. Why Apache Paimon changes the lakehouse in 2026
Paimon is a streaming-first LSM table format — that one fact decides where it beats Iceberg and Hudi
The one-sentence invariant: Paimon stores each table as an LSM tree on object storage, so real-time upserts and CDC land cheaply and are re-readable as a changelog, instead of forcing you to rewrite files on every change. Everything that makes Paimon attractive follows from that. It is still an open lake format — ACID snapshots, ORC/Parquet data files, S3/HDFS/OSS storage, engine-agnostic reads — but the write path is designed for a Flink job pushing millions of updates a minute, not for a nightly batch overwrite.
What Paimon is, in one breath.
- An open table format, not a database. Like Iceberg and Hudi, Paimon is a set of file and metadata conventions over object storage plus catalog integration; there is no always-on server owning your data.
- Born from Flink Table Store. Paimon started inside the Flink project as Flink Table Store, then became a standalone Apache top-level project, which is why its streaming semantics feel native rather than retrofitted.
- Unified streaming and batch. The same table can be written by a streaming Flink job and read either as a bounded batch table (Spark, Trino, StarRocks) or as an unbounded changelog stream (Flink), from one copy of the data.
Why the LSM tree is the whole story.
- Upserts are appends, not rewrites. New records land in small sorted files at level 0; there is no read-modify-write of a big Parquet file per update, so ingest stays cheap even under a heavy CDC load.
- Compaction merges in the background. A separate process merges sorted runs into higher levels, applying the merge engine and reclaiming space, so read amplification stays bounded without blocking writes.
- Primary keys get real point-update semantics. Because each bucket is a sorted LSM tree keyed by the primary key, an update to one row does not touch the rest of the partition.
Where Paimon sits against the alternatives.
- vs Apache Iceberg. Iceberg is snapshot-first and batch-optimized; row-level updates become merge-on-read delete files, which is excellent for occasional updates and huge scans but expensive for high-frequency CDC. Paimon's LSM makes continuous upserts the fast path, and its changelog producers hand downstream streaming jobs a correct change feed that Iceberg does not natively emit.
- vs Apache Hudi. Hudi also targets upserts (copy-on-write / merge-on-read with a record index) but grew up Spark-centric. Paimon is Flink-native first, adds Spark and query-engine reads, and its merge engines (partial-update, aggregation, first-row) cover multi-stream merge patterns Hudi leaves to the application.
-
vs a raw Parquet + manual MERGE. Hand-rolling upserts with
MERGE INTOover Parquet has no incremental changelog, no bucketed concurrency, and rewrites large files per change. Paimon gives you streaming ingest, changelog, and compaction for free.
What interviewers listen for.
- Do you say "Paimon is streaming-first, built on an LSM tree" in the first sentence? — senior signal.
- Do you frame the contrast as "Iceberg is snapshot/batch-first, Paimon is CDC/stream-first" rather than "they're all the same"? — required framing.
- Do you reach for Paimon when the requirement is "a continuously updated table that a downstream Flink job reads as a changelog"? — the whole point.
- Do you mention merge engines and changelog producers as built-ins, not application code? — the differentiator.
Worked example — the same fact loaded two ways
Detailed explanation. The clearest way to feel Paimon's shape is to declare the two table types side by side: a primary-key table that upserts, and an append-only table that only inserts. The DDL is ordinary Flink SQL against a Paimon catalog; the single line that changes everything is whether a PRIMARY KEY is declared. With a key, writes upsert and the LSM tree merges by key. Without one, the table is an immutable log.
Question. Using Flink SQL against a Paimon catalog, declare an orders primary-key table and an order_events append-only table, and state how a second write for the same order_id behaves in each.
Input.
| write | order_id | status |
|---|---|---|
| 1 | 100 | PLACED |
| 2 | 100 | SHIPPED |
Code.
-- primary-key table: second write UPSERTS row 100 to SHIPPED
CREATE TABLE orders (
order_id BIGINT,
status STRING,
PRIMARY KEY (order_id) NOT ENFORCED
) WITH ('bucket' = '4');
-- append-only table: second write ADDS another row for 100
CREATE TABLE order_events (
order_id BIGINT,
status STRING
) WITH ('bucket' = '4', 'bucket-key' = 'order_id');
Step-by-step explanation. orders declares PRIMARY KEY (order_id), so Paimon treats it as a primary-key table: the two writes for order_id = 100 collapse under the default deduplicate merge engine, and the table holds one row, SHIPPED. order_events declares no key, so it is an append-only table: both writes are retained as distinct rows — it is a log. Both set bucket = 4, so rows hash across four LSM trees for parallel writes; the append-only table pins the hash to bucket-key = order_id.
Output.
| table | rows for order_id=100 | final state |
|---|---|---|
orders (PK) |
1 | 100, SHIPPED |
order_events (append-only) |
2 |
100, PLACED + 100, SHIPPED
|
Rule of thumb. Declare a PRIMARY KEY when you want the table to reflect current state (a dimension, a CDC mirror); leave it off when you want an immutable event log. The rest of Paimon — merge engines, changelog producers — only applies to primary-key tables.
2. Table types & the LSM architecture
Primary-key vs append-only, buckets, and the LSM tree are the entire storage model — learn these and the file layout snaps into focus
Paimon has exactly two table types and one storage engine underneath both, and an interviewer who asks "how does Paimon store data?" wants buckets and the LSM tree in the answer. Get the layout crisp and every tuning knob makes sense.
The two table types.
-
Primary-key table — current state by key. Declaring a
PRIMARY KEYmakes each bucket a sorted LSM tree keyed by that primary key. Writes are upserts; a merge engine decides how two records for the same key combine. This is the table you point a CDC stream at. - Append-only table — an immutable log. No primary key; every record is an insert, nothing is updated or deleted in place. Paimon still buckets, sorts, and compacts it, so it behaves like a scalable, queryable queue — ideal for events and raw ODS ingestion.
Buckets — the unit of parallelism.
- A bucket is a horizontal shard within a partition, and each bucket is written and read by one task, so bucket count is your write/read parallelism ceiling per partition.
-
Fixed bucketing (
bucket = N) hashes rows on the bucket key intoNbuckets; simple and stable but re-bucketing means a rewrite. -
Dynamic bucketing (
bucket = -1, primary-key tables only) lets Paimon grow buckets automatically using an index that maps each key to a bucket, at the cost of maintaining that index. -
Bucket key. For primary-key tables the bucket key defaults to the primary key; for append-only tables you set
bucket-keyexplicitly to control data distribution and avoid skew.
The LSM tree inside each bucket.
- Level 0 is the landing zone. Each write flush produces a small sorted run at L0; L0 runs can overlap in key range, which is why reads may merge several of them.
- Higher levels are merged and non-overlapping. Compaction merges L0 runs downward into L1, L2, … where runs no longer overlap, bounding the number of files a read must touch.
-
Data files are columnar. The sorted runs are ORC (default) or Parquet/Avro files;
manifestfiles list them, and asnapshotfile points to the manifest set that defines the table at a commit.
Worked example — how a bucket resolves a key across LSM levels
Detailed explanation. A read for one primary key does not scan a partition; it scans one bucket's LSM tree and merges the sorted runs that could contain the key. Because higher levels are non-overlapping and L0 runs may overlap, a point read touches at most one run per higher level plus any overlapping L0 runs, then applies the merge engine to produce the current row.
Question. A primary-key orders table has bucket = 4. Row order_id = 4611 hashes to bucket 1, whose LSM tree has two overlapping L0 runs and one L1 run. Which files are read to return the current row, and how is the winner chosen?
Input.
| level | runs | contains 4611? |
|---|---|---|
| L0 | run_a (newer), run_b (older) | run_a: yes, run_b: yes |
| L1 | run_c | yes |
Code.
-- point read of one key; planner narrows to bucket 1's LSM tree
SELECT * FROM orders WHERE order_id = 4611;
-- bucket = hash(order_id) % 4 -> bucket 1
Step-by-step explanation. The planner computes the bucket from the hash of the primary key, so only bucket 1 is opened. Within bucket 1, the non-overlapping higher level (L1) contributes at most one run, while the two overlapping L0 runs are both candidates. Paimon merge-reads the three runs, and the merge engine (default deduplicate) keeps the record with the highest sequence number — the newest write of 4611. No other bucket, and no other key range, is touched.
Output.
| files merged for the read | winner |
|---|---|
| L0 run_a, L0 run_b, L1 run_c | newest by sequence number (from L0 run_a) |
Rule of thumb. More buckets = more write parallelism but more small files; size buckets so each holds a few hundred MB to a couple of GB of data, and let compaction keep L0 shallow so point reads merge few runs.
Paimon interview question on choosing a table type and bucket count
Question. You are ingesting a MySQL customers table via Flink CDC (millions of rows, frequently updated) into a Paimon lake table that a downstream job and a BI tool both read. Which table type and bucketing do you choose, and why is a plain append-only table wrong here?
Solution Using a primary-key table with fixed buckets
Code.
CREATE TABLE customers (
customer_id BIGINT,
email STRING,
tier STRING,
updated_at TIMESTAMP(3),
PRIMARY KEY (customer_id) NOT ENFORCED
) PARTITIONED BY (region) WITH (
'bucket' = '8', -- 8 LSM trees per region partition
'merge-engine' = 'deduplicate', -- current state per customer_id
'sequence.field' = 'updated_at' -- newest update wins on out-of-order
);
Step-by-step trace.
| requirement | choice | why |
|---|---|---|
| frequent updates by id | primary-key table | upserts merge in the LSM tree |
| current state for BI |
deduplicate merge engine |
one row per customer_id
|
| out-of-order CDC | sequence.field = updated_at |
newest wins, not arrival order |
| parallel ingest |
bucket = 8 per region
|
8 concurrent write tasks |
- A primary-key table is required because the source mutates rows; an append-only table would keep every historical version and the BI tool would see duplicates per
customer_id. -
merge-engine = deduplicatecollapses all writes for a key to the latest record, giving current-state semantics. -
sequence.field = updated_atprotects correctness when CDC events arrive out of order — the record with the largestupdated_atwins regardless of which arrived last. -
bucket = 8under aregionpartition gives eight independent LSM trees per partition, so eight Flink tasks write concurrently without contending on the same files.
Output:
| table | rows per customer_id | read shape |
|---|---|---|
customers |
exactly 1 (latest) | batch for BI, streaming changelog for the downstream job |
Why this works — concept by concept:
- Primary-key table — declaring a key turns each bucket into a sorted LSM tree, so an update to one customer touches one small run, not the whole partition.
- Bucket = parallelism — fixed buckets fan writes across tasks; sizing them right keeps files large enough to scan efficiently yet small enough to compact cheaply.
- deduplicate + sequence.field — together they guarantee one current row per key even when the CDC stream is reordered, which is the property BI and downstream joins depend on.
- One copy, two read modes — the same primary-key table serves a bounded batch read and an unbounded changelog read, so you avoid a separate serving store.
-
Cost — a keyed upsert is O(log levels) merge work per read plus amortized background compaction, far cheaper than the O(partition) rewrite a Parquet
MERGEwould cost.
Streaming
Topic — streaming
Streaming lakehouse ingest problems
3. Merge engines — deduplicate, partial-update, aggregation, first-row
The merge engine decides what happens when two records share a primary key — pick it at table create and it governs every compaction
Once a table has a primary key, Paimon must answer one question millions of times: two records arrive with the same key — what is the merged result? That answer is the merge engine, set with 'merge-engine' = '...' at table create, and applied every time an LSM read or compaction collapses records for a key. There are four, and an interviewer will expect you to name all four and their use cases.
The four merge engines.
-
deduplicate(default) — last writer wins. For each key, keep the most recent record (by sequence). This is the CDC-mirror / current-state engine; a delete record removes the key. -
partial-update— build a wide row from many streams. Each incoming record updates only the non-null columns it carries; nulls do not overwrite existing values. Several narrow streams (say, one job writingprice, another writingstock) merge into one complete row keyed by the primary key. -
aggregation— combine fields with agg functions. Each non-key field is folded with a configured function —sum,max,min,last_value,product,listagg,bool_or, and more — viafields.<name>.aggregate-function. Ideal for pre-aggregated wide tables and real-time metrics. -
first-row— keep the first record, ignore updates. For each key, retain the earliest record and drop later ones. It de-duplicates an append/ODS stream to "one row per key, first seen," and it only emits insert changelog, which makes it cheap for log ingestion.
Ordering — sequence.field.
- Without it, input order decides the winner, which is fragile for a reordered CDC stream.
-
sequence.fieldnames a column (usuallyupdated_ator an op offset) whose larger value is treated as newer, so the correct record wins regardless of arrival order. - Composite and generated sequences are allowed; you can combine fields or use a per-partition sequence to break ties deterministically.
Partial-update, in detail (the one they drill).
-
Nulls are skips, not overwrites — a record that carries
price=NULLleaves the storedpriceuntouched, so narrow updates never blank out columns another stream owns. - Sequence groups let different column groups have their own ordering column, so a late update to one group does not clobber a newer value in another.
-
Deletes need care — plain partial-update ignores deletes by default; enabling delete handling (or pairing with an aggregation/
sequence-group) defines how a-Dcollapses the wide row.
Worked example — partial-update merges two narrow streams into one wide row
Detailed explanation. The single most Paimon-specific merge behaviour is partial-update. Two independent Flink jobs write to the same key: one owns price, the other owns stock. Each writes only its column; Paimon merges them into one row because nulls are treated as "no change." No join, no upsert SQL — the storage engine does the column-wise merge on compaction.
Question. With a partial-update table keyed by sku, job A writes {sku: X, price: 9.99, stock: null} and job B writes {sku: X, price: null, stock: 42}. What is the stored row?
Input.
| source | sku | price | stock |
|---|---|---|---|
| job A | X | 9.99 | (null) |
| job B | X | (null) | 42 |
Code.
CREATE TABLE product (
sku STRING,
price DECIMAL(10, 2),
stock INT,
PRIMARY KEY (sku) NOT ENFORCED
) WITH (
'merge-engine' = 'partial-update',
'bucket' = '4'
);
-- job A: INSERT INTO product (sku, price) VALUES ('X', 9.99);
-- job B: INSERT INTO product (sku, stock) VALUES ('X', 42);
Step-by-step explanation. Both records share sku = X, so they land in the same bucket's LSM tree. On merge, partial-update folds them field by field: job A supplies price = 9.99 and a null stock, so stock is left for whoever provides it; job B supplies stock = 42 and a null price, so price is untouched. The result is a single row carrying both values — the null column from each record is a skip, never an overwrite.
Output.
| sku | price | stock |
|---|---|---|
| X | 9.99 | 42 |
Rule of thumb. Use partial-update when several pipelines each own a slice of one wide row; never write the whole row from one job with real nulls, because a genuine null you intend to store looks identical to "no update" and will be skipped.
Paimon interview question on real-time aggregation
Question. You need a real-time table of per-user_id running totals: total_spend should sum every incoming amount and last_seen should keep the latest timestamp, updated continuously by a Flink stream. Which merge engine and configuration, and how does a new event change the stored row?
Solution Using the aggregation merge engine
Code.
CREATE TABLE user_metrics (
user_id BIGINT,
total_spend DOUBLE,
last_seen TIMESTAMP(3),
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'merge-engine' = 'aggregation',
'fields.total_spend.aggregate-function' = 'sum',
'fields.last_seen.aggregate-function' = 'max',
'bucket' = '8'
);
-- stream: INSERT INTO user_metrics VALUES (7, 12.0, TIMESTAMP '2026-01-01 09:00:00');
Step-by-step trace.
| incoming event (user_id, amount, ts) | stored total_spend | stored last_seen |
|---|---|---|
| (7, 12.0, 09:00) | 12.0 | 09:00 |
| (7, 8.0, 09:05) | 20.0 | 09:05 |
| (7, 5.0, 08:50) late | 25.0 | 09:05 (max keeps 09:05) |
-
merge-engine = aggregationtells Paimon to fold each non-key field with its configured function instead of overwriting. -
fields.total_spend.aggregate-function = summeans every newamountis added to the running total, so three events for user 7 accumulate to25.0. -
fields.last_seen.aggregate-function = maxkeeps the newest timestamp; the late08:50event still contributes to the sum but cannot movelast_seenbackward. - The merge is applied on read and during compaction, so a batch reader always sees the fully-folded current value without a downstream
GROUP BY.
Output:
| user_id | total_spend | last_seen |
|---|---|---|
| 7 | 25.0 | 2026-01-01 09:05:00 |
Why this works — concept by concept:
-
Aggregation merge engine — moves the
GROUP BYinto the storage layer, so the table is the materialized aggregate and readers never re-scan history. -
Per-field functions —
sumon the amount andmaxon the timestamp co-exist on one row, each field folding by its own rule. -
Order independence for commutative aggregates —
sumandmaxare order-insensitive, so a late event lands correctly without a sequence field; non-commutative fields would needsequence.field. - Compaction applies the fold — because aggregation runs on compaction, read amplification stays bounded even as millions of events fold into a few keys.
- Cost — O(1) fold per event per field, versus an O(events) re-aggregation every query if you stored raw rows and grouped at read time.
Merge
Topic — merge
Upsert and merge-engine problems
4. Changelog producers & streaming CDC
A changelog producer is how Paimon emits a correct +I/-U/+U/-D stream — get it wrong and downstream streaming joins silently break
Paimon's superpower is that a downstream Flink job can read a primary-key table as a stream and receive a complete changelog — inserts, updates (as a retract -U plus an add +U), and deletes. But producing that changelog is not free: the merge that happens inside the LSM tree can hide whether a record was new or an update. The changelog-producer option decides how Paimon generates the change stream, and it is the single most misunderstood Paimon config in interviews.
Why a changelog is hard.
-
Streaming joins need retractions. A downstream job that aggregates or joins must be told the old value (
-U) before the new one (+U), or its state drifts. A stream of only "here is the current row" is not enough. - The LSM merge loses the before-image. When two records for a key merge during compaction, the pre-update value is gone unless Paimon deliberately captures it. The changelog producer is that capture strategy.
The four producers.
-
none(default) — no changelog file. Cheapest write; the streaming read emits only the merged records, so a downstream operator must normalize (dedup / compute retractions) itself. Fine when downstream does not need exact retractions. -
input— trust the input. Paimon double-writes the incoming records as the changelog. Correct only when the source is already a complete changelog (e.g. a Flink CDC stream that emits-U/+U/-D), because Paimon just passes the input through. -
lookup— produce a complete changelog via lookup. Before commit, Paimon looks up the previous value for each changed key and emits a correct-U/+U. Lower latency than full-compaction, at the cost of extra lookup I/O; works even when the input is insert-only. -
full-compaction— produce the changelog during full compaction. A complete changelog is derived when full compaction runs, so it is cheap on the write path but the changelog latency equals the full-compaction interval.
Choosing a producer.
-
Input already a complete CDC changelog →
input(cheapest correct option). -
Insert-only or upsert input, need correct retractions, want low latency →
lookup. -
Insert-only or upsert input, can tolerate higher latency, want low write cost →
full-compaction. -
Downstream does not need retractions →
none.
Streaming read knobs.
-
scan.modepicks the starting point:latest-full(full snapshot then continue),latest(only new changes),from-snapshot/from-timestampfor a specific point,compacted-fullto start from a compacted snapshot. -
consumer-idrecords read progress in the table so a restarted job resumes exactly where it left off and Paimon does not expire snapshots the consumer still needs.
Worked example — reading a primary-key table as a changelog stream
Detailed explanation. The everyday streaming pattern is a Flink CDC source writing into a Paimon primary-key table configured with changelog-producer = lookup, and a second streaming job reading that table. The reader receives a real changelog, so it can maintain a correct downstream aggregate without re-reading history.
Question. A orders table has changelog-producer = lookup. An order row updates from PENDING to PAID. What rows does a downstream streaming reader see for that change?
Input.
| event | order_id | status |
|---|---|---|
| insert | 100 | PENDING |
| update | 100 | PAID |
Code.
CREATE TABLE orders (
order_id BIGINT,
status STRING,
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'bucket' = '4',
'changelog-producer' = 'lookup'
);
-- downstream job reads the changelog as a stream
SELECT status, COUNT(*) FROM orders GROUP BY status;
Step-by-step explanation. The insert of order 100 emits a single +I (100, PENDING). When 100 updates to PAID, the lookup producer looks up the previous value and emits two changelog rows: -U (100, PENDING) then +U (100, PAID). The downstream GROUP BY status applies the retraction — decrementing PENDING and incrementing PAID — so the counts stay exact without rescanning the table.
Output.
| changelog rows the reader sees | downstream count effect |
|---|---|
+I (100, PENDING) |
PENDING: 1 |
-U (100, PENDING), +U (100, PAID)
|
PENDING: 0, PAID: 1 |
Rule of thumb. If any downstream job aggregates or joins on a Paimon primary-key table, it needs real retractions — use lookup (low latency) or full-compaction (low write cost); none will quietly corrupt the downstream state.
Paimon interview question on changelog-producer correctness
Question. Your Flink CDC source already emits a full -U/+U/-D changelog for a MySQL table, and you want the Paimon table to hand that exact change stream to three downstream jobs with minimum write cost. Which changelog-producer do you set, and why is lookup the wrong choice here?
Solution Using changelog-producer = input
Code.
CREATE TABLE cdc_mirror (
id BIGINT,
payload STRING,
op_ts TIMESTAMP(3),
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'bucket' = '4',
'changelog-producer' = 'input', -- pass the source changelog through
'sequence.field' = 'op_ts' -- order by CDC operation time
);
Step-by-step trace.
| source CDC row | producer = input action | changelog emitted |
|---|---|---|
+I (1, a) |
write + pass through | +I (1, a) |
-U (1, a) / +U (1, b)
|
write + pass through |
-U (1, a), +U (1, b)
|
-D (1, b) |
write + pass through | -D (1, b) |
- Because the Flink CDC source already produces a complete
-U/+U/-Dchangelog, Paimon does not need to reconstruct before-images. -
changelog-producer = inputdouble-writes the incoming records as the changelog file, so there is no lookup I/O and no wait for full compaction. -
sequence.field = op_tsorders records by CDC operation time, so out-of-order delivery still resolves to the correct latest state. - All three downstream jobs read the same pre-built changelog, so the cost is paid once on write, not per reader.
Output:
| table | write cost | changelog latency |
|---|---|---|
cdc_mirror |
lowest (no lookup, no forced full compaction) | equals commit interval |
Why this works — concept by concept:
- input producer — trusts a source that is already a complete changelog, so Paimon simply persists it; the cheapest correct option when the precondition holds.
-
Precondition matters —
inputis only correct because the CDC source emits retractions; on an insert-only source it would emit no-U, solookuporfull-compactionwould be required instead. - sequence.field — decouples correctness from arrival order, the standard safeguard for any CDC ingest.
- Write once, read many — the changelog is materialized on write, so N downstream readers share it without re-deriving change events.
-
Cost — O(input rows) extra write for the changelog file, versus the O(changed keys) lookup I/O of
lookupor the full-compaction latency offull-compaction.
Streaming
Topic — streaming
Streaming changelog and CDC problems
5. Snapshots, compaction, time travel & tags
Every commit is a snapshot — compaction merges files, time travel reads old snapshots, and tags pin them for batch
Paimon's table state is a sequence of snapshots, each an atomic commit that points to a set of manifest and data files. That one idea gives you ACID commits, time travel, and tags, and it is the layer an interviewer probes when asking "how does Paimon stay consistent and let me read the past?" Say it plainly: a snapshot is one commit; time travel reads an old snapshot; a tag is a snapshot you promise not to delete.
Snapshots — the commit unit.
- Each write commit creates a new numbered snapshot referencing the manifests that define the table at that instant; readers pin one snapshot for a consistent view.
-
Snapshots expire.
snapshot.num-retained.min/snapshot.num-retained.maxandsnapshot.time-retainedbound how many/how long snapshots live, so metadata does not grow unbounded — but expiration also limits how far back you can time-travel or resume a streaming read. -
Consumers protect snapshots. A streaming reader with a
consumer-idholds back expiration of snapshots it has not yet consumed, so restarts are safe.
Compaction — keeping reads fast.
- Minor compaction merges small L0 runs to bound the number of overlapping files; full compaction merges everything in a bucket to one sorted run, which also produces the full-compaction changelog.
-
Inline vs dedicated. By default the writer compacts inline, but for heavy ingest you run a dedicated compaction job (a separate Flink job /
compactaction) so writing and compacting scale independently. -
write-only. Setwrite-only = trueon the ingest job and delegate all compaction to the dedicated job, avoiding write/compact contention.
Time travel — reading the past.
-
Flink:
... FOR SYSTEM_TIME AS OF ..., orscan.snapshot-id/scan.timestamp-millisin the table hints, selects an old snapshot. -
Spark:
VERSION AS OF <snapshot-id>andTIMESTAMP AS OF <ts>read the same historical state. - Use it for reproducing a report, debugging "what did the table look like at 9am," or auditing a change.
Tags — long-lived named snapshots.
- A tag names a snapshot and exempts it from expiration, so it survives even after the snapshot would normally be cleaned up — the mechanism for stable batch reads over a streaming table.
-
Automatic tags.
tag.automatic-creation(e.g. daily) plustag.num-retained-maxcreate and retain periodic tags, giving you a daily "end-of-day" view for batch pipelines. -
Incremental between tags. You can read the delta between two tags (
incremental-between) to feed a batch job only what changed between two end-of-day points.
Engine integration.
- Flink is the streaming read/write engine (CDC ingest, changelog read, dedicated compaction actions).
- Spark reads and writes, does batch time travel and tag reads, and DDL/DML via the Paimon catalog.
- Query engines — Trino, Presto, StarRocks, Doris, Hive — read Paimon tables for interactive/batch analytics, so one table serves streaming and BI at once.
Worked example — reading yesterday's state with time travel
Detailed explanation. Because every commit is a numbered snapshot, reading an earlier state is just naming an older snapshot id or timestamp. The live table keeps advancing; the time-travel query pins a fixed snapshot, so it returns exactly what the table held at that commit, independent of writes happening now.
Question. The orders table is at snapshot 57. A report needs the table exactly as it was at snapshot 50. Show the Flink and Spark reads and what each returns.
Input.
| snapshot | order_id=100 status |
|---|---|
| 50 | PENDING |
| 57 (current) | PAID |
Code.
-- Flink: pin an old snapshot via table hint
SELECT * FROM orders /*+ OPTIONS('scan.snapshot-id' = '50') */
WHERE order_id = 100;
-- Spark: version-as-of time travel
SELECT * FROM orders VERSION AS OF 50 WHERE order_id = 100;
Step-by-step explanation. Both queries resolve the table to snapshot 50's manifest set instead of the current snapshot 57. The planner reads only the data files that snapshot 50 referenced, applying the merge engine over those files, and ignores every file committed afterward. So even though the live row is now PAID (snapshot 57), the time-travel read returns the PENDING value that was current at snapshot 50.
Output.
| query | order_id=100 status returned |
|---|---|
Flink scan.snapshot-id = 50
|
PENDING |
Spark VERSION AS OF 50
|
PENDING |
Rule of thumb. Time travel only reaches snapshots that have not expired — if you need a specific point to stay readable for weeks (month-end close, audit), create a tag, because tags are exempt from snapshot expiration.
Paimon interview question on stable batch reads over a streaming table
Question. A Flink job continuously upserts into a Paimon table, but a nightly Spark batch must always read a consistent "end of day" view — even though snapshot expiration would normally remove old snapshots within hours. How do you guarantee the batch always sees a stable daily state?
Solution Using automatic daily tags
Code.
CREATE TABLE dwd_orders (
order_id BIGINT,
status STRING,
dt STRING,
PRIMARY KEY (dt, order_id) NOT ENFORCED
) PARTITIONED BY (dt) WITH (
'bucket' = '8',
'snapshot.time-retained' = '2 h', -- snapshots expire fast
'tag.automatic-creation' = 'process-time',
'tag.creation-period' = 'daily', -- pin one snapshot per day
'tag.num-retained-max' = '90' -- keep 90 daily tags
);
-- nightly Spark batch reads a stable tagged snapshot
SELECT * FROM dwd_orders VERSION AS OF '2026-01-01';
Step-by-step trace.
| time | event | effect |
|---|---|---|
| all day | Flink upserts | new snapshots; old ones expire after 2h |
| 00:00 | daily tag created | snapshot pinned as tag 2026-01-01
|
| next day | more upserts + expiry | tag 2026-01-01 survives expiration |
| nightly | Spark reads tag | stable end-of-day view returned |
- Ordinary snapshots expire after
snapshot.time-retained = 2 h, so raw time travel cannot reach yesterday. -
tag.automatic-creation = process-timewithtag.creation-period = dailypins one snapshot per day as a named tag. - A tag is exempt from snapshot expiration, so tag
2026-01-01remains readable even after its underlying snapshot would have been cleaned up. - The Spark batch reads
VERSION AS OF '2026-01-01', resolving the tag to a fixed, consistent daily state regardless of concurrent streaming writes.
Output:
| reader | sees |
|---|---|
| nightly Spark batch | consistent end-of-day snapshot via the daily tag |
| live Flink stream | latest snapshot, unaffected |
Why this works — concept by concept:
- Snapshot = atomic commit — a stable point to read from exists after every write, and readers pin one for a consistent view.
- Tag pins a snapshot — naming a snapshot exempts it from expiration, converting a short-lived streaming snapshot into a durable batch checkpoint.
-
Automatic daily tags — periodic tag creation gives you an end-of-day series without a manual step, while
num-retained-maxbounds history. - Streaming and batch decoupled — the Flink writer keeps advancing and expiring snapshots aggressively, yet the batch reader is insulated because it reads a tag, not a live snapshot.
- Cost — a tag is metadata only (it keeps referenced files alive), so the cost is retained storage for tagged snapshots, not extra compute.
Partitioning
Topic — partitioning
Partitioned time-travel and tag problems
Cheat sheet — Paimon recipes
Primary-key table (Flink SQL).
CREATE TABLE t (
id BIGINT, v STRING,
PRIMARY KEY (id) NOT ENFORCED
) WITH ('bucket' = '4');
Append-only table with an explicit bucket key.
CREATE TABLE events (
id BIGINT, kind STRING
) WITH ('bucket' = '8', 'bucket-key' = 'id');
Merge engine + sequence field.
CREATE TABLE dim (
id BIGINT, attr STRING, ts TIMESTAMP(3),
PRIMARY KEY (id) NOT ENFORCED
) WITH ('merge-engine' = 'deduplicate', 'sequence.field' = 'ts');
Aggregation table (running totals).
CREATE TABLE m (
k BIGINT, total DOUBLE,
PRIMARY KEY (k) NOT ENFORCED
) WITH ('merge-engine' = 'aggregation',
'fields.total.aggregate-function' = 'sum');
Changelog producer for streaming reads.
ALTER TABLE t SET ('changelog-producer' = 'lookup');
-- 'input' if the source is already a full changelog; 'full-compaction' for cheap writes
Streaming read with a scan mode + consumer.
SELECT * FROM t
/*+ OPTIONS('scan.mode' = 'latest-full', 'consumer-id' = 'job1') */;
Time travel + tag read.
-- Flink
SELECT * FROM t /*+ OPTIONS('scan.snapshot-id' = '50') */;
-- Spark
SELECT * FROM t VERSION AS OF 50; -- by snapshot id
SELECT * FROM t VERSION AS OF '2026-01-01'; -- by tag name
Disposition / engine picker.
| Situation | Choice |
|---|---|
| Current state, CDC mirror | primary-key + deduplicate
|
| Immutable event log | append-only table |
| Wide row from many streams | partial-update |
| Real-time metrics | aggregation |
| Downstream streaming join | changelog-producer = lookup |
| Stable daily batch view | automatic daily tag
|
Frequently asked questions
What is Apache Paimon?
Apache Paimon is an open-source lake table format that stores each table as a log-structured merge (LSM) tree on object storage such as S3, HDFS, or OSS. It began as Flink Table Store and is now an Apache top-level project, and it is designed streaming-first: primary-key tables accept high-frequency upserts and CDC, background compaction merges the data, and downstream jobs can read the table as a bounded batch table or an unbounded changelog stream. It gives you ACID snapshots, time travel, and engine-agnostic reads while making real-time ingest the fast path.
How is Paimon different from Iceberg and Hudi?
Iceberg is snapshot-first and batch-optimized, handling row-level updates as merge-on-read delete files — great for large scans and occasional updates but costly for continuous CDC. Hudi is upsert-oriented but grew up Spark-centric around a record index and a copy-on-write / merge-on-read split. Paimon is streaming-first: an LSM tree makes continuous upserts cheap, merge engines define how records for a key combine, and changelog producers emit a correct +I/-U/+U/-D stream that neither Iceberg nor Hudi produces natively. Choose Paimon when your table is fed by a real-time stream and read as a changelog.
What is a primary-key table vs an append-only table in Paimon?
A primary-key table declares a PRIMARY KEY, so each bucket is a sorted LSM tree keyed by it; writes are upserts and a merge engine decides how two records for the same key combine, giving current-state semantics. An append-only table has no primary key, so every record is an immutable insert — it behaves like a scalable, queryable log, ideal for events and raw ODS ingestion. Merge engines and changelog producers apply only to primary-key tables.
What are Paimon merge engines?
A merge engine decides what happens when two records share a primary key. deduplicate (default) keeps the latest record; partial-update merges non-null columns from several narrow streams into one wide row; aggregation folds each field with a function such as sum or max; and first-row keeps the first record and ignores later updates. You set it with merge-engine at table create, and sequence.field controls which record is considered newest when events arrive out of order.
What is a changelog producer and why does streaming need it?
A changelog producer controls how Paimon emits the +I/-U/+U/-D change stream that a downstream streaming job needs to keep aggregates and joins correct. none writes no changelog (cheapest, downstream must normalize); input passes through a source that is already a full changelog; lookup reconstructs correct retractions with low latency; and full-compaction derives the changelog during full compaction at low write cost but higher latency. Pick lookup or full-compaction when the input is insert-only or upsert and downstream needs exact retractions.
Does Paimon support time travel and tags?
Yes. Every commit is a numbered snapshot, and you can time-travel to an older one with scan.snapshot-id / scan.timestamp-millis in Flink or VERSION AS OF / TIMESTAMP AS OF in Spark. Because snapshots expire, Paimon also supports tags — named, long-lived snapshots exempt from expiration — which you can create automatically (for example daily) to give batch jobs a stable end-of-day view over a continuously streaming table, and read the delta between two tags for incremental batch processing.
Practice on PipeCode
Pipecode.ai is Leetcode for Data Engineering — every Paimon idea above, from the primary-key LSM table to the partial-update merge engine, the lookup changelog producer, and the daily tag over a streaming table, maps to a hands-on practice room where you build the pipeline against real graded inputs. PipeCode pairs each reading with 450+ DE-focused problems and a real-time scoring engine, so your answer to "how would you keep a downstream streaming join correct?" holds up under a senior interviewer's depth probes.





Top comments (0)