Skip to main content

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.

Choose primary-key tables for merging incoming records by key, or append tables for independently stored rows; then select distribution, read mode, and optional AI features.

Quick Decision​

WorkloadStarting pointCheck before adopting it
CDC or upsert replicationPrimary-key table, default deduplicate merge engineSource ordering, partition changes, and downstream changelog requirements
Independent updates to columns of one entityPrimary-key table, partial-updatePer-source sequence groups and null/delete semantics
Accumulating metric contributionsPrimary-key table, aggregationInput replay and each function's retraction behavior
Keep the first record received for a keyPrimary-key table, first-rowArrival order differs from earliest event time
Batch ETL, logs, or partition overwriteUnaware-bucket append tablePartition granularity and file sizes
Equality filters or compatible bucketed joinsFixed-bucket append tableComplete bucket-key predicates, skew, and the query plan
Preserve append order within a partition and bucketOrdered bucketed append tableNo global or event-time ordering guarantee
BLOBs, feature backfills, or vector searchAI Data PipelinesStorage 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 contractProducer
Batch reads only, or consumers can handle upsertsLeave the default none.
Input already contains the complete required before/after changesUse input.
Input lacks before images, but consumers need complete changesConsider 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 modeConfigurationMain consideration
Dynamic (default)bucket = -1Assignment adapts to data; use one writing job.
FixedPositive bucketPlan for skew, parallelism, file sizes, and later rescaling.
Postponebucket = -2New 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​

ModeConfigurationWork to budget for
MOR (default)No extra settingReaders merge overlapping versions; compaction reduces that work.
COWfull-compaction.delta-commits = 1Frequent full compaction increases write amplification.
MOWdeletion-vectors.enabled = trueWriters 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​

RequirementUnaware-bucketFixed-bucket
Bucket-key pruning and compatible bucketed joinsNoYes, with compatible keys and query plans
Append order within a partition and bucketNot guaranteedYes, when bucket-append-ordered = true
Incremental clusteringSupportedRequires bucket-append-ordered = false
Data EvolutionSupported with its required optionsNot 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 taskWalkthrough
Store payloadsBLOBs and metadata
Search embeddingsVector storage and search
Update featuresColumn backfill
Train or analyze in PythonTraining 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.