Skip to main content
Version: 1.0

Table Overview

A Fluss table combines a typed schema with rules for storing, distributing, updating, and retaining rows. Those choices determine how applications write data, how readers access it, and how the workload uses cluster resources.

This overview introduces the logical and physical structure of tables, then explains the design decisions to make before creating one. For the servers and storage services that implement this model, see Fluss Architecture.

Start with the Access Pattern​

Begin with the data the application needs to preserve and the operations its consumers perform.

Design questionFluss mechanism
Must each event remain independently readable within a retention window?A log table stores appended records without replacing earlier rows.
Must applications update an entity and look up its current state?A primary key table maintains row state by key and records its changes.
Should data be managed or filtered in groups, such as by date?Partition columns organize rows into separately managed partitions.
How should writes and reads share cluster capacity?Bucket count and bucket keys control distribution and parallelism.
Which update should win when a key appears again?The primary key table's merge engine defines how writes combine with stored state.
How long are current state and change history needed?Log TTL, key-value TTL, and partition retention govern different data lifecycles.

These decisions are related. For example, partitioning a primary key table by date also makes that date part of the row's identity. Retaining current rows for a month does not automatically retain a month of their change history.

Database​

A database is a namespace containing tables. A table is identified by its database and table name, such as commerce.orders. Applications can use databases to group related datasets while configuring schemas, distribution, and retention independently for each table.

Native client APIs expose database and table administration. In Flink SQL, a Fluss catalog makes those objects available through SQL. The examples on this page assume that a Fluss catalog is already configured and selected; see Flink DDL for setup.

Flink SQL
CREATE DATABASE commerce;
USE commerce;

Table​

A table defines named, typed columns together with optional primary keys, partition keys, bucket keys, and table properties. The presence of a primary key determines the table type.

PropertyLog tablePrimary key table
Row identityEach append creates another record.The primary key identifies the row to insert, update, or delete.
WritesAppend only.Upsert and delete, subject to the selected merge engine's restrictions.
Repeated valuesRepeated records remain separate appends.Writes for an existing key are processed using the merge engine.
Typical readsStreaming consumption and replay by bucket offset.Key lookups, supported snapshot scans, and change consumption.
StorageOrdered log.Current key-value state and an ordered log of changes.
Typical usesEvents, activity streams, and ingestion logs.Entity state, reference data, and maintained results.

Log Tables​

Choose a log table when records represent events that should remain individually consumable. An order_created event and a later order_shipped event can both remain in the log even if they refer to the same order. Repeating an event_id does not, by itself, deduplicate records.

Log readers track progress through offsets within each bucket. Records in one bucket have an order, but there is no single ordering across all buckets. When related events need to reach the same bucket, define an appropriate hash bucket key.

See Log Tables for write behavior, log consumption, compression, and projected reads.

Primary Key Tables​

Choose a primary key table when a row represents the current state of an entity, such as an order or user profile. An upsert for an existing key updates its state according to the merge engine. The default engine retains the last written row; additional engines support other update rules.

A primary key table also exposes changes for downstream processing. Keeping one current row per key does not remove the need to plan changelog retention for consumers that may fall behind or need to replay data.

Full-key lookups retrieve the matching row. Supported clients and connectors also provide snapshot scans and prefix lookups, with constraints on key order and bucket keys. See Primary Key Tables and the Client Support Matrix for the available operations.

Schema and Key Design​

Select column types for the values applications exchange, including numeric precision, timestamp semantics, and nullability. For example, use an appropriate DECIMAL precision and scale for monetary values and define a consistent timestamp interpretation across writers. Consult Data Types and the client support matrix for type availability and mappings.

Three kinds of keys serve different purposes:

KeyPurposePrimary key table constraints
Primary keyDefines row identity and the target of updates and lookups.Includes every partition column.
Partition keySelects the partition containing the row.Must be a subset of the primary key.
Bucket keyHashes the row to a bucket within its table or partition.Must be a subset of the primary key, excluding partition columns.

For primary key tables, the default bucket key is the primary key with partition columns removed. Explicit bucket keys must not contain partition columns for either table type.

Choose a stable primary key that matches the application's identity model. If an order should have one current row throughout its lifetime, an unpartitioned table keyed by order_id may be appropriate. A table partitioned by order_date with primary key (order_date, order_id) instead treats the same order ID on different dates as different keys. Writing a new key does not remove a row stored under the previous key.

Schema Evolution​

Fluss supports adding nullable columns at the end of an existing schema without rewriting stored data. Existing rows have no stored value for the new columns and are read with null values for those fields. Schema changes have constraints, so plan keys and required fields at table creation and consult Adding Columns before evolving a table.

For lake-enabled tables, evolution must also be supported by the selected lake format. Apply supported schema changes through Fluss so the managed lake schema stays aligned with the Fluss table.

Table Data Organization​

The logical hierarchy is a database containing tables, with optional partitions and one or more buckets. Tablets are the physical storage structures used to serve each bucket.

Partition​

A partition groups rows by a distinct value or combination of values in the partition columns. For example, a table partitioned by event_date stores each date's records in a separate partition. Supported readers can use partition filters to limit the data they scan, and operators can manage retention in whole partitions.

Both table types support partitioning. Fluss provides manual partition management, automatic time-based partition management, and client-side dynamic creation when writes encounter new partition values. These mechanisms can coexist, but their configuration and client support differ.

Choose partition columns that match filtering and lifecycle requirements. A column with many distinct values can create a large number of partitions, each with its own buckets and associated storage overhead. Bucketing provides parallelism within each partition, so a partition does not need to represent one server or one writer.

See Partitioning for creation strategies, multi-column partitions, and historical partition access.

Bucket​

A bucket divides a table or partition into a unit of storage placement, replication, and processing parallelism. The bucket.num property controls the number of buckets. Multiple buckets can reside on the same TabletServer, and their replicas can be reassigned as the cluster changes.

For primary key tables, hash bucketing ensures that operations on a key route to the same bucket. Log tables use sticky assignment by default, which groups records into batches, and can also use round-robin or hash assignment.

Choose bucket keys with enough distinct values to distribute the workload. Hashing on a heavily skewed value can concentrate traffic in one bucket even when many buckets exist. Adding TabletServers does not split an individual hot key across servers.

Bucket count balances parallelism against the overhead of additional replicas, files, and storage instances. In a partitioned table, each partition captures the configured bucket count when it is created. Altering bucket.num affects future partitions only, including those created dynamically; existing and already pre-created partitions retain their original counts. This operation does not redistribute existing records.

See Bucketing for assignment strategies and the rules for changing bucket counts.

LogTablet​

Every bucket has log storage. A LogTablet contains the appended records of a log table or the changes of a primary key table. It is organized into log segments with data files and indexes that support locating records by offset and timestamp.

The bucket's log can be replicated across TabletServers. Older segments can also be uploaded to remote storage and remain readable within the configured retention period. For primary key tables, the log supports both change consumption and recovery of key-value state.

KvTablet​

A KvTablet maintains a primary key bucket's current state in an embedded RocksDB instance. The bucket leader hosts this state alongside its log and serves key-based operations.

Follower replicas replicate the log without maintaining equivalent live KvTablets. Recovery normally restores a remote key-value snapshot and replays subsequent log records. The architecture guide explains replication and recovery in more detail.

Update and Change Semantics​

A primary key defines which rows are related; the merge engine defines what to do when another write arrives for that key.

Merge engineBehaviorExample use
Default (LastRow)Applies the latest write; supports partial updates when configured.Current order or profile state.
FirstRowRetains the first record for a key.Deduplicating repeated event identifiers.
VersionedUses a version column to accept records whose version is at least the stored version.Handling out-of-order entity updates.
AggregationCombines values using configured per-column aggregate functions.Maintaining counters or aggregated results.

The engines differ in delete, partial-update, and changelog support. Select one using its documented semantics, especially when consuming upstream change events. The default last-write behavior follows write processing order; it does not automatically choose the largest event timestamp.

The changelog image is a separate choice from the merge engine and storage format. The default FULL image includes before and after records for updates. WAL omits update-before records, which affects consumers that need previous values. Check these requirements before tuning a table for lookup or ingestion performance.

Virtual Tables expose change records through names such as orders$changelog and, for primary key tables, orders$binlog. They read the retained Fluss log; a lake snapshot does not reconstruct expired change history for these interfaces.

Retention and Storage Layout​

Plan current state, replay history, and partition lifetime separately.

MechanismControlsDesign implication
table.log.ttlRetention of log-table records or primary key changelogs.Cover consumer lag, restart time, and required replay history.
table.log.local-ttl and table.log.tiered.local-segmentsCleanup of local log copies after upload to remote storage.Local retention determines the hot working set; remote logs can remain readable longer.
table.kv.ttlExpiration of individual primary key rows.Current state can have a different lifetime from its changelog.
Auto-partition retentionExpiration of entire time partitions.All data in an expired partition is subject to that lifecycle, regardless of longer row or log TTLs.

Expiration and cleanup run asynchronously. Key-value TTL cleanup does not emit delete events to the changelog, so downstream consumers should not rely on it to receive row-deletion notifications. Local log cleanup also has independent time-based and segment-count policies; a local TTL alone does not guarantee that every record remains local for that duration. See Data Retention and TTL for the complete behavior and configuration constraints.

The log format is another independent setting. ARROW, the default, uses a columnar layout that supports selective column reads. COMPACTED uses a row-oriented encoding. A format controls representation, while the table type and merge engine control row semantics. In particular, selecting COMPACTED does not make a log table deduplicate records by key. See Data Formats for the trade-offs and supported combinations.

Lakehouse tiering can maintain an additional representation for historical access. Its retention, freshness, and schema rules depend on the lake integration. Review Lakestream Overview when designing a table for both streaming and lakehouse readers.

Example Table Designs​

These examples use Flink SQL in the commerce database created above. The bucket counts illustrate the configuration; size them for the workload and cluster.

Partitioned Order Events​

An append-only table preserves individual order events, groups them by date for lifecycle management, and hashes by order ID to route related events within each date to the same bucket.

Flink SQL
CREATE TABLE order_events (
event_date STRING,
event_id STRING,
order_id BIGINT,
event_type STRING,
event_time TIMESTAMP(3)
) PARTITIONED BY (event_date)
WITH (
'bucket.num' = '8',
'bucket.key' = 'order_id',
'table.log.ttl' = '7 d'
);

ALTER TABLE order_events ADD PARTITION (event_date = '2026-10-05');

This table has no primary key, so two appends with the same event_id remain separate records. Hashing by order_id gives bucket-local ordering within a date partition; it does not provide ordering across dates. The explicit partition creation shown here does not enable automatic partition expiration.

Current Order State​

An unpartitioned primary key table maintains one current row per order, supports lookups by order ID, and retains a separate window of change history.

Flink SQL
CREATE TABLE orders (
order_id BIGINT,
customer_id BIGINT,
status STRING,
total_amount DECIMAL(18, 2),
updated_at TIMESTAMP(3),
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'bucket.num' = '8',
'table.log.ttl' = '7 d'
);

The default bucket key is order_id, and the default merge engine applies the latest write for each order. The NOT ENFORCED clause is part of Flink's primary key declaration syntax; Fluss still uses the key to maintain row state. Log expiration removes old change history without expiring current rows. The updated_at column does not control merge order unless a suitable merge engine is explicitly configured.

Continue Designing Your Table​

Before creating production tables, confirm the row identity and update rules, estimate partition and bucket counts, and define retention for both current state and replay history. Verify that the chosen client or engine supports the required reads, writes, data types, and lakehouse features.

Use Flink DDL, Spark DDL, or the Java Client to create and evolve tables. The detailed chapters on Table Types, Partitioning, Bucketing, Merge Engines, and Retention provide the constraints behind each design choice.