Scenario Guide
Choose a table by what incoming records mean. Then choose distribution, read behavior, and
streaming output. A need for SQL UPDATE or DELETE alone does not require a primary-key table:
append tables also support engine-specific row-level operations.
Quick Decision
| Workload | Starting point | Check before adopting it |
|---|---|---|
| CDC or upsert replication | Primary-key table, default deduplicate merge engine | Source ordering, partition changes, and downstream changelog requirements |
| Independent updates to columns of one entity | Primary-key table, partial-update | Per-source sequence groups and null/delete semantics |
| Accumulating metric contributions | Primary-key table, aggregation | Input replay and each function's retraction behavior |
| Keep the first record received for a key | Primary-key table, first-row | Arrival order differs from earliest event time |
| Batch ETL, logs, or partition overwrite | Unaware-bucket append table | Partition granularity and file sizes |
| Equality filters or compatible bucketed joins | Fixed-bucket append table | Complete bucket-key predicates, skew, and the query plan |
| Preserve append order within a partition and bucket | Ordered bucketed append table | No global or event-time ordering guarantee |
| BLOBs, feature backfills, or vector search | AI Data Pipelines | Storage requirements and index freshness are separate decisions |
The SQL examples below use Flink SQL in an existing Paimon catalog. Bucket counts are small illustrative values, not sizing recommendations. For Spark syntax, use the linked feature guides.
Primary Key Table
Use a primary-key table when incoming records should merge into one logical row per key. Pick the merge engine before tuning the storage mode.
CDC Real-Time Sync
For a source that supplies a monotonically increasing version per key:
CREATE TABLE orders (
order_id BIGINT,
amount DECIMAL(12, 2),
status STRING,
source_version BIGINT,
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'sequence.field' = 'source_version'
);
The default deduplicate engine keeps the row with the largest sequence value. Use a source
version that is non-null and expresses the required precedence; timestamps can tie. A sequence field alone does
not resolve equal versions deterministically. See Sequence Field.
Choose the changelog producer according to the records reaching Paimon and the consumer's needs:
| Consumer/input contract | Producer |
|---|---|
| Batch reads only, or consumers can handle upserts | Leave the default none. |
| Input already contains the complete required before/after changes | Use input. |
| Input lacks before images, but consumers need complete changes | Consider lookup; budget for lookup compaction. |
The transport does not determine changelog completeness: a Kafka stream can carry complete CDC, and a database connector can omit information required by a particular consumer. Partial-update and aggregation inputs also differ from their final merged rows.
If analytical reads dominate, evaluate MOW with deletion vectors at table creation. Check its compaction and visibility requirements instead of enabling it as a universal CDC default. For whole-database ingestion and schema evolution, see CDC Ingestion.
If you partition a primary-key table, the partition columns normally belong to its primary key. An entity moving between partitions needs an explicit design; see Cross-Partition Upsert.
Partial Column Updates
For two sources that update different parts of an order:
CREATE TABLE order_wide (
order_id BIGINT PRIMARY KEY NOT ENFORCED,
amount DECIMAL(12, 2),
order_version BIGINT,
delivery_status STRING,
delivery_version BIGINT
) WITH (
'merge-engine' = 'partial-update',
'fields.order_version.sequence-group' = 'amount',
'fields.delivery_version.sequence-group' = 'delivery_status',
'changelog-producer' = 'lookup'
);
Each source supplies its own version and leaves the other group's version null. A null sequence
skips that group; an accepted group update can set its value fields to null. Fields outside
sequence groups use the default non-null update behavior. lookup serves consumers that need
complete merged rows. See Partial Update.
Combine the input streams into one writing job when using the default dynamic buckets. Multiple sources do not imply that multiple independent writers can safely assign dynamic buckets.
Streaming Metrics
Use aggregation when each input is a contribution, such as an additional quantity sold:
CREATE TABLE product_metrics (
product_id BIGINT PRIMARY KEY NOT ENFORCED,
quantity BIGINT,
max_price DOUBLE
) WITH (
'merge-engine' = 'aggregation',
'fields.quantity.aggregate-function' = 'sum',
'fields.max_price.aggregate-function' = 'max',
'changelog-producer' = 'lookup'
);
Replaying a contribution can add it again. Do not send full replacement totals to a sum field
unless that is the intended calculation. Functions differ in support for retracts, deletes, and
input encoding; sketches require prepared sketch values. See
Aggregation.
Keep the First Record
CREATE TABLE first_login (
user_id BIGINT,
dt STRING,
login_time TIMESTAMP(3),
PRIMARY KEY (user_id, dt) NOT ENFORCED
) PARTITIONED BY (dt) WITH (
'merge-engine' = 'first-row',
'changelog-producer' = 'lookup'
);
This keeps the first row received for a key, not necessarily the row with the earliest
login_time. Lookup compaction produces an insert-only changelog of newly retained keys.
User-defined sequence fields and ordinary row deletes are not supported. Do not add deletion
vectors to this configuration; see First Row.
Choose Buckets and Read Behavior
| Bucket mode | Configuration | Main consideration |
|---|---|---|
| Dynamic (default) | bucket = -1 | Assignment adapts to data; use one writing job. |
| Fixed | Positive bucket | Plan for skew, parallelism, file sizes, and later rescaling. |
| Postpone | bucket = -2 | New data needs compaction before it becomes queryable. |
Use Data Distribution to size and validate the layout. A single bytes-per-bucket rule cannot capture update rate, partition count, or available memory.
Table Mode
| Mode | Configuration | Work to budget for |
|---|---|---|
| MOR (default) | No extra setting | Readers merge overlapping versions; compaction reduces that work. |
| COW | full-compaction.delta-commits = 1 | Frequent full compaction increases write amplification. |
| MOW | deletion-vectors.enabled = true | Writers resolve versions through lookup compaction; readers apply deletion vectors. |
These are read/write trade-offs, not interchangeable presets.
In particular, lookup changelog generation cannot be combined with full-compaction.delta-commits,
and deletion-vector tables do not support the full-compaction changelog producer.
Append Table
Use an append table when each input is stored independently, or when your engine manages batch replacements and row-level changes. A business identifier in the schema does not automatically need to be a declared primary key.
Batch ETL and Analytical Queries
CREATE TABLE events (
event_id BIGINT,
user_id BIGINT,
event_type STRING,
dt STRING
) PARTITIONED BY (dt);
With no primary key and no bucket setting, this uses the unaware-bucket layout. Start here for batch inserts and partition overwrite. Add partition filters and inspect file statistics before adding layout options. For selective queries, evaluate File Indexes or Incremental Clustering. Supported SQL mutations and their costs are described in Row-Level Operations.
Equality Lookups and Bucketed Joins
CREATE TABLE product_logs (
product_id BIGINT,
log_time TIMESTAMP(3),
message STRING
) WITH (
'bucket' = '8',
'bucket-key' = 'product_id'
);
SELECT * FROM product_logs WHERE product_id = 12345;
An equality or IN predicate on the complete bucket key can prune unrelated buckets. This does
not guarantee reading exactly one eighth of the bytes: bucket sizes can be skewed, and matching
buckets still contain other keys. A composite bucket key requires constraints on every key column.
Spark can exploit compatible bucket distributions in joins with V2 bucketing enabled. Confirm
shuffle removal with EXPLAIN; see Bucketed Join.
Ordered Streaming
The product_logs table above also preserves append order within each partition and bucket
because bucket-append-ordered defaults to true. It does not sort by log_time or define a
global order across buckets. See Bucketed Streaming.
Bucketed incremental clustering requires bucket-append-ordered = false, giving up that order
guarantee. Decide whether ordering or clustering is the requirement before combining the options.
Compare Append Layouts
| Requirement | Unaware-bucket | Fixed-bucket |
|---|---|---|
| Bucket-key pruning and compatible bucketed joins | No | Yes, with compatible keys and query plans |
| Append order within a partition and bucket | Not guaranteed | Yes, when bucket-append-ordered = true |
| Incremental clustering | Supported | Requires bucket-append-ordered = false |
| Data Evolution | Supported with its required options | Not supported |
Multimodal Data Lake
For payload storage, embedding search, feature backfills, and Python training readers, continue with AI Data Pipelines. That guide separates storage, column updates, and indexes, and provides a complete small vector-search example.
| Pipeline task | Walkthrough |
|---|---|
| Store payloads | BLOBs and metadata |
| Search embeddings | Vector storage and search |
| Update features | Column backfill |
| Train or analyze in Python | Training readers |
Distributed Feature Backfill with Ray
The Ray backfill walkthrough shows how to select records, compute a derived feature, and merge only that feature back into a BLOB table.
Validate Your Choice
Try a representative write and query, including duplicate input, out-of-order updates, and any required deletes. Verify both the batch result and downstream changelog. Then use Understand Files to inspect the committed layout and Streaming Writes and Small Files to assess file growth.