Apache Paimon is a lake table format designed for streaming updates. Like other table formats it keeps data as Parquet (by default) or ORC files on object storage and adds metadata that makes a set of files behave like a transactional table. What sets Paimon apart is its primary-key table: each bucket of the table is a log-structured merge tree, so a Flink job can upsert a high-rate stream of changes into the lake with minutes or seconds of latency, and another job can read the table back as a stream of changes.

This article explains the storage model from the files up, then the choices you make when you create a table: bucket mode, merge engine and changelog producer. It closes with compaction, streaming reads, a worked example and the failure modes that show up in production. Configuration keys and defaults follow the Paimon documentation as checked on 2026-10-01; check the docs for your version, because options do move between releases.

Advertisement

The problem it solves

Writing a change stream into a lake is hard because object storage files are immutable. An update to one row means either rewriting the file that holds it, which is expensive at high update rates, or writing the update somewhere else and merging at read time, which makes reads expensive unless something cleans up. Copy-on-write and merge-on-read are the two ends of that trade. Paimon takes the LSM approach familiar from databases such as RocksDB: new writes go into small sorted files, reads merge the overlapping files by key, and background compaction merges files into larger ones so reads stay cheap.

The result is a table that behaves well in three roles at once: a sink for a CDC or event stream, a source for downstream streaming jobs, and a normal batch table for Spark, Flink batch or other engines. Apache Hudi covers similar ground with a different design; the comparison is at the end.

Snapshots, manifests and files

Everything starts from a snapshot. A snapshot file records a committed table state: the schema id and the manifest lists that describe its files, plus optional changelog and index metadata. A manifest list names manifest files and carries partition statistics. A data manifest records additions and deletions of data files with their row counts and statistics; a changelog manifest does the same for changelog files. Data files hold the records; in a primary-key table each record also carries its row kind (insert, update or delete) and a sequence number.

On disk a table directory has snapshot/, schema/, manifest/ and index/ directories next to the partition directories, and each partition holds bucket directories such as bucket-0. Writers prepare files first and then publish a snapshot; publishing makes all of the changes visible together, so a reader of a snapshot sees a committed file set and never a partial one. That is what gives you atomic commits on storage that has no transactions.

Paimon primary-key table: Flink writers flush sorted runs per bucket and commit a snapshot at each checkpointCDC / Kafka sourceinserts, updates, deletesFlink writer taskshash key to bucketCommitteron checkpoint: new snapshotfile listBucket LSM treelevel 0: new sorted runslevel 1..n: compacted runsmerge engine applied on readand during compactionflushsnapshot-Nschema id, manifest listsmanifest list-> manifests -> fileswriteReadersbatch: one snapshotstream: each new snapshotchangelog files if producedconsumer-id pins progressObject storage: snapshot/ schema/ manifest/ index/ dt=.../bucket-k/ data and changelog filesA snapshot is the unit of visibility: readers see a committed file set, never a partially written one.
The write path for a primary-key table. Writers flush sorted runs into per-bucket LSM trees; the committer publishes a snapshot that references them through manifests.
Advertisement

Primary-key tables are LSM trees per bucket

Within a bucket, writers buffer incoming records in memory, sort them by key and flush them as new files. A sorted run is a set of files whose key ranges do not overlap, so within one run each key appears at most once. Different runs can overlap and hold different versions of the same key. A read merges the relevant runs and applies the table's merge engine to records with the same key, a merge-on-read scan. Compaction merges runs to reduce the number a reader has to combine.

Two settings bound that work. When the number of sorted runs reaches num-sorted-run.compaction-trigger (default 5), compaction starts; each level-0 file counts as one run and each occupied higher level counts as one. When the backlog reaches num-sorted-run.stop-trigger (default trigger plus 3), writers wait for compaction. That pause is backpressure, and it is the first thing to look at when a Paimon sink suddenly slows; see streaming backpressure.

Choosing a bucket mode

The bucket is the unit of parallel writing and the scope of one LSM tree. Paimon has three modes, set with the bucket option:

ModeSettingBehaviourUse when
Fixed'bucket' = 'N'Hash of the bucket key modulo N; predictable layout; changing N needs an explicit rescaleYou know the size and want stable, parallel writes
Dynamic'bucket' = '-1' (default)Keeps a key-to-bucket index and adds buckets as data grows, aiming at dynamic-bucket.target-row-numSize is unknown; one writer per partition is acceptable
Postpone'bucket' = '-2'Stages data first and assigns real buckets during compaction or batch writes, so partitions can have different bucket countsPartition sizes vary widely

Two constraints are worth knowing before you design keys. In fixed mode the primary key must include the partition columns, because a key can only be found in its own partition. Dynamic bucket mode can support cross-partition upserts by maintaining a key-to-partition mapping on disk, but that index costs memory, disk and bootstrap time, and with deduplicate the row moves to the new partition while partial-update and aggregation apply the change in the original one. A bucket-key can be set to a subset of the primary key, excluding partition columns, to control distribution.

As a starting point for fixed buckets, aim for about 1 GB per bucket, the same target Paimon's postpone mode uses by default (postpone.target-size-per-bucket), and make N a multiple of the writer parallelism you expect, because parallelism above N leaves tasks idle.

Merge engines decide what a key's value is

When several records share a primary key, the merge engine decides the stored result. It is set with merge-engine:

  • deduplicate (default): keep the latest row; a later delete removes it. Note that a null in the latest row replaces an earlier value, which is correct for full-row CDC and wrong for sparse updates. Set sequence.field to an event-time or version column if arrival order is not the true order.
  • partial-update: update only non-null fields, so several streams can each fill in different columns of one wide row. Sequence groups let independent column groups be ordered separately.
  • aggregation: combine values with a per-field function set by fields.<name>.aggregate-function, for example sum or max, turning the table into a continuously updated aggregate.
  • first-row: keep the first row for each key and ignore later ones, which is a cheap way to deduplicate event or log streams.

The documentation makes a distinction that is easy to miss: the merge engine defines the stored table state, and the changelog producer defines what streaming readers receive. Getting the first right does not guarantee the second.

Changelog producers decide what streaming readers see

A downstream streaming job usually wants a proper changelog: for an update, the old value (update-before) and the new value (update-after). Whether Paimon can provide that depends on changelog-producer:

ProducerHow the changelog is madeCost and latency
none (default)No separate changelog; readers see new records without complete before-images. Flink can add a normalize operator that keeps state to rebuild themCheapest writes; expensive stateful normalize downstream
inputStores incoming records as changelog filesCorrect only if the source already supplies full before and after images, such as database CDC
lookupLooks up existing values during compaction and emits the resulting changesWriters wait for lookup compaction before commit; uses local lookup caches
full-compactionDiffs successive full compaction resultsIntermediate updates collapse; latency is the full compaction interval

lookup is the usual choice when the source does not provide before-images and consumers need fresh, accurate changes; size lookup.cache-max-memory-size (default 256 MB) and the disk cache to fit. lookup and full-compaction.delta-commits cannot be combined.

Commits, compaction and deletion vectors

In Flink, writers flush files continuously but a snapshot is committed when a checkpoint completes, so checkpoint interval sets data freshness and the number of small files. A one-minute checkpoint gives roughly one-minute freshness. Because the commit is tied to the checkpoint, Paimon integrates with Flink's exactly-once guarantees; see exactly-once semantics for the general mechanism.

By default each writer also compacts its own buckets. With several jobs writing to one table, or to give compaction its own resources, set write-only to true on the writers and run a single dedicated compaction job; two compaction jobs on the same partition will produce commit conflicts. Deletion vectors (deletion-vectors.enabled) take the merge work away from readers: compaction marks superseded rows in a bitmap so a reader can skip them without merging runs, which makes reads from engines with weak merge support much faster at the cost of more work at write time. For batch readers that only need data as of the last full compaction, the $ro read-optimized system table avoids merging altogether.

Reading the table as a stream

A streaming read starts from a position given by scan.mode and then follows each new snapshot. Two maintenance settings interact with it. Snapshots expire under snapshot.time-retained (default 1 h) and snapshot.num-retained.min (default 10), and an expired snapshot's files can be deleted. A streaming reader that falls behind past expiry fails. Setting a consumer-id on the reader records its progress in the table, protects the snapshots it still needs from expiry and lets a restarted job resume from where it stopped. The cost is that an abandoned consumer id holds storage forever, so remove stale consumers.

Worked example: orders into the lake

An orders service writes change events (CDC) for about 50 million live orders, roughly 3,000 updates per second at peak. The team wants a lake table that analysts can query and a streaming job that updates revenue per customer. Orders are partitioned by order date, and an update to an order rarely crosses partitions.

SET 'execution.runtime-mode' = 'streaming';
SET 'execution.checkpointing.interval' = '1 min';

CREATE CATALOG lake WITH ('type' = 'paimon', 'warehouse' = 's3://lake/warehouse');
USE CATALOG lake;

CREATE TABLE orders (
  order_id    BIGINT,
  dt          STRING,
  customer_id BIGINT,
  status      STRING,
  amount      DECIMAL(12, 2),
  updated_at  TIMESTAMP(3),
  PRIMARY KEY (order_id, dt) NOT ENFORCED
) PARTITIONED BY (dt) WITH (
  'bucket' = '8',
  'merge-engine' = 'deduplicate',
  'sequence.field' = 'updated_at',
  'changelog-producer' = 'lookup'
);

INSERT INTO orders SELECT order_id, dt, customer_id, status, amount, updated_at FROM orders_cdc;

CREATE TABLE revenue_by_customer (
  customer_id BIGINT,
  revenue     DECIMAL(18, 2),
  PRIMARY KEY (customer_id) NOT ENFORCED
) WITH ('bucket' = '4', 'merge-engine' = 'aggregation',
        'fields.revenue.aggregate-function' = 'sum');

INSERT INTO revenue_by_customer
SELECT customer_id, amount
FROM orders /*+ OPTIONS('consumer-id' = 'revenue', 'scan.mode' = 'latest') */;

Why these choices: daily partitions are small enough for eight fixed buckets each, and the primary key includes dt as fixed mode requires. sequence.field protects against CDC events arriving out of order. The lookup producer gives the downstream job proper update-before and update-after records, so when an order's amount changes from 100 to 120 the aggregation sees minus 100 and plus 120 rather than a second plus 120. The aggregation table relies on that: without a correct changelog, a summing sink double counts every update. To check health, query orders$snapshots for commit cadence and orders$files for file counts per bucket.

Failure modes

  • Writers stalled at the stop trigger: compaction cannot keep up, so sorted runs pile up and writers wait. Give compaction more resources, use a dedicated compaction job or reduce small commits.
  • Small-file explosion: very frequent checkpoints with many buckets create many tiny files. Lengthen the checkpoint interval or reduce buckets.
  • Double counting downstream: a summing consumer reading a table whose changelog lacks before-images. Use lookup, input with full CDC, or a normalize step.
  • Streaming job fails after downtime: snapshots it needed expired. Use a consumer-id and longer retention.
  • Storage that never shrinks: a forgotten consumer id pins old snapshots. Audit consumers regularly.
  • Commit conflicts: two jobs compacting the same partition. Keep exactly one compactor per table.
  • Nulls wiping values: sparse updates written to a deduplicate table. Use partial-update.

Trade-offs against other formats

Paimon is strongest when the table is updated by a stream at high rate and read as a stream as well, with Flink at the centre. Hudi offers a similar merge-on-read design with a timeline and a broader set of table services. Formats such as Iceberg and Delta Lake are widely supported by query engines and are excellent for batch-oriented updates, but row-level upserts at streaming rates are not their primary design. The LSM approach costs background compaction work and makes the bucket layout a decision you must get roughly right up front. Check that the engines your readers use support the Paimon features you enable, especially deletion vectors and changelog reads.

What to do next

  1. Create a test table with Flink SQL against local storage, write a small CDC stream and inspect the snapshot/, manifest/ and bucket directories it produces.
  2. Choose the bucket mode from your partition sizes; for fixed buckets, include partition columns in the primary key and size buckets to about 1 GB.
  3. Pick the merge engine from the shape of your input: full rows, sparse columns, contributions to sum, or events to deduplicate.
  4. Pick the changelog producer from what downstream consumers need, and use lookup when they need correct before-images.
  5. Set the checkpoint interval from your freshness target, then watch file counts and sorted-run backlog.
  6. Give every production streaming reader a consumer-id, set snapshot retention deliberately and audit consumers monthly.
Key takeaway: Paimon turns a directory of immutable files into a streaming-updatable table by publishing snapshots that reference LSM sorted runs per bucket. Its behaviour comes from three choices made at table creation: the bucket mode sets the layout, the merge engine sets the stored value of a key, and the changelog producer sets what streaming readers receive. Size checkpoints and buckets to avoid small files, keep one compactor per table, and pin streaming readers with consumer ids.