From 45e383285befc8230fbffbc4271473ba35fc4405 Mon Sep 17 00:00:00 2001 From: Giannis Polyzos Date: Mon, 5 Oct 2026 10:33:14 +0300 Subject: [PATCH] Update website concepts page (#4546) * [website] update concepts page * update fluss architecture * update table design (cherry picked from commit 64b3b46114882b9d40b492fa5d45e6c12604573f) --- website/docs/concepts/architecture.mdx | 205 +++++++++++++++++--- website/docs/quickstart/_category_.json | 2 +- website/docs/table-design/overview.mdx | 236 ++++++++++++++++++++---- 3 files changed, 383 insertions(+), 60 deletions(-) diff --git a/website/docs/concepts/architecture.mdx b/website/docs/concepts/architecture.mdx index 9e3db6f4420..a7316a0817b 100644 --- a/website/docs/concepts/architecture.mdx +++ b/website/docs/concepts/architecture.mdx @@ -2,10 +2,35 @@ title: "Architecture" sidebar_position: 1 --- + +{/* +Licensed to the Apache Software Foundation (ASF) under one or more +contributor license agreements. See the NOTICE file distributed with +this work for additional information regarding copyright ownership. +The ASF licenses this file to You under the Apache License, Version 2.0 +(the "License"); you may not use this file except in compliance with +the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/} + import ThemedImage from '@site/src/components/ThemedImage'; # Architecture -A Fluss cluster consists of two main processes: the **CoordinatorServer** and the **TabletServer**. + +Apache Fluss is a distributed streaming storage system that exposes data as tables. It combines an ordered log for continuous data consumption with a key-value store for updates and lookups, and integrates with remote filesystems and open lakehouse formats for recovery and historical access. + +This chapter explains how those components work together and what their responsibilities mean for deployment and operation. + +## Architecture Overview + +A Fluss cluster separates coordination from data serving. The **CoordinatorServer** manages metadata and placement, while **TabletServers** host table data and serve client requests. **ZooKeeper** provides shared metadata and coordination. Remote storage holds tiered log segments and key-value snapshots; an optional lakehouse integration maintains data in an open lake format. +| Component | Responsibility | +| --- | --- | +| CoordinatorServer | Manages databases, tables, partitions, TabletServer membership, and bucket replica placement and leadership. | +| TabletServer | Hosts bucket replicas, processes reads and writes, replicates logs, and manages local storage and background storage tasks. | +| ZooKeeper | Stores cluster metadata and supports server registration and coordinator leader election. | +| Remote storage | Stores Fluss log segments and key-value snapshots outside individual TabletServers. | +| Clients and connectors | Discover metadata, route requests, and expose table operations to applications and compute engines. | +| Lakehouse tiering service | Runs as a separate Flink job to maintain an open lakehouse representation of enabled tables. | + +The coordinator participates in metadata and administrative operations. Once a client has discovered the relevant servers, normal table reads and writes go directly to TabletServers. This separates the volume of data traffic from the coordination workload. + +## Tables, Partitions, and Buckets + +Fluss distributes a table through three related abstractions: + +- A **table** defines a schema, table type, and storage properties. +- A **partition**, when configured, groups rows by partition column values, such as a date or region. Each partition has its own buckets. +- A **bucket** is the unit of data distribution, log ordering, and replica assignment. A non-partitioned table has one set of buckets. + +The client selects a bucket using the table's distribution rules. Primary key tables use hash bucketing so operations on the same key reach the same bucket. Log tables can use hash bucketing or distribute records with sticky or round-robin assignment. See [Bucketing](../table-design/data-distribution/bucketing.md) and [Partitioning](../table-design/data-distribution/partitioning.md) for configuration and constraints. + +### Tablet / Bucket + +A bucket describes a logical slice of a table; a **tablet** is its storage representation on a TabletServer. A **LogTablet** stores the bucket's ordered log. For a primary key table, the bucket leader also maintains a **KvTablet** containing its key-value state. The leader's LogTablet and KvTablet are hosted together so updates, lookups, and log processing can be coordinated locally. + +A bucket's log can have multiple replicas on different TabletServers. Each bucket has a leader, which serves client writes, and follower replicas, which fetch its log. One TabletServer typically hosts replicas for many buckets and can be a leader for some and a follower for others. + +Ordering is defined within a bucket. Records in different buckets can be processed concurrently, so applications should not assume a single table-wide order. + ## CoordinatorServer -The **CoordinatorServer** serves as the central control and management component of the cluster. It is responsible for maintaining metadata, managing tablet allocation, listing nodes, and handling permissions. -Additionally, it coordinates critical operations such as: -- Rebalancing data during node scaling (upscaling or downscaling). -- Managing data migration and service node switching in the event of node failures. -- Overseeing table management tasks, including creating or deleting tables and updating bucket counts. +The CoordinatorServer manages the cluster's control plane: the metadata and decisions that allow clients and TabletServers to agree on where data belongs. + +Its responsibilities include: + +- **Metadata management:** Creating and altering databases, tables, and partitions, and maintaining their schemas and properties. +- **Server membership:** Tracking available TabletServers and responding to membership changes. +- **Replica placement and leadership:** Assigning bucket replicas and coordinating leader changes when servers become unavailable. +- **Cluster maintenance:** Coordinating replica reassignment and rebalancing as servers are added, maintained, or decommissioned. +- **Storage coordination:** Recording metadata such as completed key-value snapshots and lakehouse tiering progress. -As the **brain** of the cluster, the **CoordinatorServer** ensures efficient cluster operation and seamless management of resources. +Multiple CoordinatorServer instances can be deployed for high availability. They use ZooKeeper to elect one active coordinator; the other instances remain on standby. When leadership changes, a coordinator epoch identifies the new leader and fences stale coordinator operations. + +Coordinator failover temporarily affects administrative operations. Existing data reads and writes can continue against healthy TabletServers, but operations that require new metadata or a placement change depend on the coordinator becoming available. See [CoordinatorServer High Availability](../install-deploy/deploying-distributed-cluster.md#fluss-coordinatorserver-high-availability-ha-setup) for deployment instructions. ## TabletServer -The **TabletServer** is responsible for data storage, persistence, and providing I/O services directly to users. It comprises two key components: **LogStore** and **KvStore**. -- For **PrimaryKey Tables** which support updates, both **LogStore** and **KvStore** are activated. The KvStore is used to support updates and point lookup efficiently. LogStore is used to store the changelogs of the table. -- For **Log Tables** which only support appends, only the **LogStore** is activated, optimizing performance for write-heavy workloads. -This architecture ensures the **TabletServer** delivers tailored data handling capabilities based on table types. +TabletServers form the data-serving layer. They accept requests for the buckets they host, manage local log and key-value files, replicate logs, and perform background work such as snapshotting and remote log uploads. + +The storage components used by a bucket depend on its table type: +| Table type | LogStore | KvStore | Typical access | +| --- | --- | --- | --- | +| Log table | Stores appended records. | Not used. | Append events and consume or replay the log. | +| Primary key table | Stores changes used for streaming consumption and state recovery. | Maintains the current row state by primary key. | Upsert or delete rows, look up keys, and consume changes. | ### LogStore -The **LogStore** is designed to store log data, functioning similarly to a database binlog. -Messages can only be appended, not modified, ensuring data integrity. -Its primary purpose is to enable low-latency streaming reads and to serve as the write-ahead log (WAL) for restoring the **KvStore**. + +LogStore maintains an append-only sequence of records for each bucket. Records have offsets, which identify positions in that bucket's log and allow readers to track progress or resume consumption. The log is organized into segments on local disk, with older segments eligible for upload to remote storage. + +For log tables, the log contains the appended data. For primary key tables, it records changes generated by the server as writes are processed. This log also provides the recovery history for the key-value state. + +Fluss uses Apache Arrow as its default log format. Columnar record batches support compression and projected reads, allowing compatible clients to request the columns they need. See [Log Tables](../table-design/table-types/log-table.md) for the storage format, column pruning, and compression options. ### KvStore -The **KvStore** is used to store table data, functioning similarly to database tables. It supports data updates and deletions, enabling efficient querying and table management. Additionally, it generates comprehensive changelogs to track data modifications. -### Tablet / Bucket -Table data is divided into multiple buckets based on the defined bucketing policy. +KvStore maintains a primary key table's current row state using RocksDB. It supports key-based access and applies the table's merge behavior when new records arrive. Depending on the table configuration, writes can replace rows, update selected columns, or use another supported merge engine. + +The current state and its change history serve different access patterns: lookups retrieve rows by key, while log readers consume changes over time. The server coordinates updates to these representations so the log can also be used to recover the key-value state. + +The active KvTablet is maintained on the bucket leader. Followers replicate the log without maintaining equivalent live KvTablets. Fluss periodically uploads key-value snapshots to remote storage and uses snapshots together with log replay to restore state after a leadership change. This makes snapshot freshness and replay volume relevant to primary key table recovery time. + +See [Primary Key Tables](../table-design/table-types/pk-table.md) for merge engines, partial updates, and the change data feed. + +## Read and Write Paths + +### Discovery and Routing + +A native client starts with configured bootstrap servers, discovers cluster metadata, and resolves the servers responsible for the requested table and buckets. It uses this metadata to route operations and refreshes its view when leadership or server availability changes. -Data for the **LogStore** and **KvStore** is stored within tablets. Each tablet consists of a **LogTablet** and, optionally, a **KvTablet**, depending on whether the table supports updates. -Both **LogStore** and **KvStore** adhere to the same bucket-splitting and tablet allocation policies. As a result, **LogTablets** and **KvTablets** with the same `tablet_id` are always allocated to the same **TabletServer** for efficient data management. +Clients must be able to reach the server addresses advertised by the cluster. A reachable bootstrap endpoint alone is insufficient if the resulting TabletServer addresses are inaccessible from the application's network. See [Distributed Deployment](../install-deploy/deploying-distributed-cluster.md) for listener configuration. -The **LogTablet** supports multiple replicas based on the table's configured replication factor, ensuring high availability and fault tolerance. **Currently, replication is not supported for KvTablets**. +### Writing Data -## Zookeeper -Fluss currently utilizes **ZooKeeper** for cluster coordination, metadata storage, and cluster configuration management. -In upcoming releases, **ZooKeeper will be replaced** by **KvStore** for metadata storage and **Raft** for cluster coordination and ensuring consistency. This transition aims to streamline operations and enhance system reliability. See the [Roadmap](/roadmap) for more details. +For a **log table**, the client selects a bucket, batches records, and sends the append request to the bucket leader. The leader appends the batch to its local log, and followers fetch the new records. Completion depends on the writer's acknowledgment settings. +For a **primary key table**, the client routes an upsert or delete by key. The bucket leader evaluates the write using the table's merge rules, records the resulting change in the log, and updates the key-value state. The replicated log provides the history needed to reconstruct that state. Acknowledgment is governed by the configured replication requirements; periodic snapshot upload is a separate background operation. + +### Reading Data + +Fluss exposes several access paths with different purposes: + +| Access path | Behavior | +| --- | --- | +| Log scan | Reads a bucket's records from a chosen offset and can continue consuming new records. For primary key tables, this reads the change history. | +| Primary key lookup | Routes a key to its bucket and retrieves the row from the serving key-value state. | +| Snapshot scan | Reads a primary key snapshot through clients that support it, with an associated log position for subsequent change consumption. | +| Remote log read | Retrieves retained log data that has been offloaded from local storage. | +| Union Read | Uses a supported lakehouse integration to read committed lake data together with the relevant Fluss data. | + +Log consumers read up to the bucket's **high watermark**, the boundary of committed log data. This boundary advances with replication progress; for primary key buckets, advancement is also coordinated with key-value state updates. Separate buckets advance independently. + +Availability of each access path depends on the client and integration. Consult the [Client Support Matrix](../apis/client-support-matrix.md) before choosing an SDK. + +## Replication and Recovery + +### Log Replication and Write Durability + +Fluss replicates each bucket's log according to the table's replication factor. The leader tracks an **in-sync replica set (ISR)** containing replicas that meet the replication requirements. Followers that fall behind can leave this set and rejoin after catching up. + +Three settings work together to determine write durability and availability: + +| Setting | Role | +| --- | --- | +| `table.replication.factor` | Determines the number of log replicas assigned to a table's buckets. | +| `client.writer.acks` | Determines whether the writer waits for no acknowledgment, the leader, or all in-sync replicas. | +| `log.replica.min-in-sync-replicas-number` | Sets the minimum in-sync replica count required for writes using `acks=all`. | + +With `acks=all`, a write waits for the in-sync replicas, rather than every configured replica regardless of its health. If the minimum in-sync requirement cannot be met, the write fails. Choosing replication factor, acknowledgment policy, and minimum ISR together is therefore essential; replica count alone does not define the durability of an acknowledged write. + +Replicas should also be placed across the failure domains the deployment must tolerate. Copies on different processes offer limited protection if they share the same host or disk failure. See [Server Configuration](../maintenance/configuration.md) and [Rebalancing](../maintenance/operations/rebalance.md) for the relevant settings and placement controls. + +### TabletServer Failure + +When a TabletServer fails, the coordinator coordinates leadership and replica changes for affected buckets. Clients refresh their metadata and retry eligible operations against the current leader. + +For a primary key bucket, the new leader must also prepare its key-value state before serving it. The normal recovery path restores a completed remote snapshot and replays the log from the snapshot's recorded offset. When no snapshot is available, recovery relies on the required log history being available. Snapshot age, retained log volume, and storage bandwidth all influence how long recovery takes. + +## ZooKeeper + +ZooKeeper is part of the current Fluss deployment architecture. It stores shared metadata and supports coordinator leader election and server registration. Fluss processes use this information to discover membership and coordinate changes consistently. + +ZooKeeper holds coordination and metadata state; table records are stored by TabletServers and the configured storage layers. Its availability is nevertheless important for metadata operations and failure recovery. Deploy its ensemble as a shared cluster dependency and configure all Fluss servers to use the same ZooKeeper namespace. ## Remote Storage -**Remote Storage** serves two primary purposes: -- **Hierarchical Storage for LogStores:** By offloading LogStore data, it reduces storage costs and accelerates scaling operations. -- **Persistent Storage for KvStores:** It ensures durable storage for KvStore data and collaborates with LogStore to enable fault recovery. -Additionally, **Remote Storage** allows clients to perform bulk read operations on Log and Kv data, enhancing data analysis efficiency and reduce the overhead on Fluss servers. In the future, it will also support bulk write operations, optimizing data import workflows for greater scalability and performance. +Remote storage, such as S3, HDFS, or another supported filesystem, extends storage beyond the disks of individual TabletServers. It has two main roles: + +- **Tiered logs:** Fluss asynchronously uploads eligible log segments, allowing historical records to remain accessible after their local copies are removed according to retention policies. +- **Key-value snapshots:** Fluss stores completed snapshots and their metadata so primary key state can be restored and supported clients can scan snapshots. + +Recent log data remains local for low-latency access. Remote log retention, local retention, and snapshot retention have separate purposes and should be configured together with recovery and replay requirements. Asynchronous uploads also mean that remote storage does not replace replication for newly written records. + +These files use Fluss's storage representations. Uploading them to an object store does not itself create an Iceberg, Paimon, or other lakehouse table. See [Remote Storage](../maintenance/tiered-storage/remote-storage.md) and [Filesystem Configuration](../maintenance/tiered-storage/filesystems/overview.md) for setup and lifecycle details. + +## Lakehouse Integration + +For tables with lakehouse storage enabled, a separate **tiering service** runs as a Flink job. It incrementally reads Fluss data, writes the configured open lake format, commits the lake changes, and records the corresponding tiering progress in Fluss. + +That progress connects the lake representation to positions in the Fluss log. Compatible readers can use it for **Union Read**, combining committed lake data with the relevant streaming data. An engine reading the lake table directly sees the committed lake representation, whose freshness depends on tiering progress. + +Supported behavior varies by lake format, table type, engine, and execution mode. Shared table metadata does not imply a simultaneous snapshot across every table or engine. See [Lakestream Overview](../streaming-lakehouse/overview.mdx), [Tiering Service](../streaming-lakehouse/tiering-service.md), and [Union Read](../streaming-lakehouse/union-read.md) for the integration model and supported access paths. ## Client -Fluss clients/SDKs support streaming reads/writes, batch reads/writes, DDL and point queries. Currently, the main implementation of client is Flink Connector. Users can use Flink SQL to easily operate Fluss tables and data. \ No newline at end of file + +Applications can access Fluss through native clients, including Java, Rust, Python, and C++, or through compute engine integrations such as Flink and Spark. Clients expose table and administrative operations; the compute engines execute queries, joins, aggregations, and pipelines using those interfaces. + +The optional **Fluss Gateway** provides a stateless REST entry point for metadata, DDL, and schema-aware batch writes. It connects to the underlying cluster as a client and can be deployed separately from CoordinatorServers and TabletServers. Its preview API has a narrower scope than the native clients; see [Fluss Gateway](../gateway/index.md) for supported operations and deployment requirements. + +Start with the [Java Client](../apis/java/index.md), [Client Support Matrix](../apis/client-support-matrix.md), [Flink Integration](../engine-flink/getting-started.md), or [Spark Integration](../engine-spark/getting-started.md) for the access method appropriate to your application. + +## Deployment and Operational Considerations + +Architecture choices affect both capacity and recovery: + +- **Coordinator availability:** Run standby coordinators and provide a reliable ZooKeeper ensemble so metadata and placement operations can recover from a process failure. +- **TabletServer capacity:** Size local disk, memory, CPU, and network bandwidth for the combined workload of client traffic, replication, RocksDB maintenance, snapshots, and remote uploads. +- **Bucket distribution:** Use enough buckets for the required parallelism and watch for skew. Adding servers increases available capacity, while rebalancing distributes existing replicas across that capacity. +- **Scaling boundaries:** Moving bucket replicas changes placement. Changing a partitioned table's bucket count applies to future partitions; it does not redistribute existing partitions. +- **Storage and recovery:** Monitor replica health, snapshot completion, remote uploads, and lakehouse tiering progress. Retention and snapshot policies should support the application's replay and recovery requirements. + +For deployment procedures, continue with [Distributed Deployment](../install-deploy/deploying-distributed-cluster.md) or [Kubernetes Deployment](../install-deploy/deploying-with-helm.md). For ongoing operation, use the [Rebalancing Guide](../maintenance/operations/rebalance.md), [Monitoring Guide](../maintenance/observability/monitor-metrics.md), and [Security Overview](../security/overview.md). diff --git a/website/docs/quickstart/_category_.json b/website/docs/quickstart/_category_.json index d06bba85d6e..4cd65b74891 100644 --- a/website/docs/quickstart/_category_.json +++ b/website/docs/quickstart/_category_.json @@ -1,4 +1,4 @@ { - "label": "Quickstart", + "label": "Quickstarts", "position": 2 } diff --git a/website/docs/table-design/overview.mdx b/website/docs/table-design/overview.mdx index 6fdf15f405e..fe1a9634d6d 100644 --- a/website/docs/table-design/overview.mdx +++ b/website/docs/table-design/overview.mdx @@ -2,56 +2,230 @@ sidebar_label: Overview title: Table Overview sidebar_position: 1 +description: "Design Apache Fluss tables by choosing table types, schemas, primary keys, partitions, buckets, merge behavior, and retention policies." --- -import ThemedImage from '@site/src/components/ThemedImage'; + +{/* +Licensed to the Apache Software Foundation (ASF) under one or more +contributor license agreements. See the NOTICE file distributed with +this work for additional information regarding copyright ownership. +The ASF licenses this file to You under the Apache License, Version 2.0 +(the "License"); you may not use this file except in compliance with +the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/} # 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](../concepts/architecture.mdx). + +## Start with the Access Pattern + +Begin with the data the application needs to preserve and the operations its consumers perform. + +| Design question | Fluss 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 collection of Table objects. You can create/delete databases or create/modify/delete tables under a 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](../engine-flink/ddl.md) for setup. + +```sql title="Flink SQL" +CREATE DATABASE commerce; +USE commerce; +``` ## Table -In Fluss, a Table is the fundamental unit of user data storage, organized into rows and columns. Tables are stored within specific databases, adhering to a hierarchical structure (database -> table). -Tables are classified into two types based on the presence of a primary key: -- **Log Tables:** - - Designed for append-only scenarios. - - Support only INSERT operations. -- **Primary Key Tables:** - - Used for updating and managing data in business databases. - - Support INSERT, UPDATE, and DELETE operations based on the defined primary key. +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. + +| Property | Log table | Primary key table | +| --- | --- | --- | +| Row identity | Each append creates another record. | The primary key identifies the row to insert, update, or delete. | +| Writes | Append only. | Upsert and delete, subject to the selected merge engine's restrictions. | +| Repeated values | Repeated records remain separate appends. | Writes for an existing key are processed using the merge engine. | +| Typical reads | Streaming consumption and replay by bucket offset. | Key lookups, supported snapshot scans, and change consumption. | +| Storage | Ordered log. | Current key-value state and an ordered log of changes. | +| Typical uses | Events, 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](table-types/log-table.md) for write behavior, log consumption, compression, and projected reads. -A Table becomes a [Partitioned Table](/table-design/data-distribution/partitioning.md) when a partition column is defined. Data with the same partition value is stored in the same partition. Partition columns can be applied to both Log Tables and Primary Key Tables, but with specific considerations: -- **For Log Tables**, partitioning is commonly used for log data, typically based on date columns, to facilitate data separation and cleaning. -- **For Primary Key Tables**, the partition column must be a subset of the primary key to ensure uniqueness. +### Primary Key Tables -This design ensures efficient data organization, flexibility in handling different use cases, and adherence to data integrity constraints. +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](table-types/pk-table.md) and the [Client Support Matrix](../apis/client-support-matrix.md) 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](data-types.md) and the client support matrix for type availability and mappings. + +Three kinds of keys serve different purposes: + +| Key | Purpose | Primary key table constraints | +| --- | --- | --- | +| Primary key | Defines row identity and the target of updates and lookups. | Includes every partition column. | +| Partition key | Selects the partition containing the row. | Must be a subset of the primary key. | +| Bucket key | Hashes 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](../engine-flink/ddl.md#add-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** is a logical division of a table's data into smaller, more manageable subsets based on the values of one or more specified columns, known as partition columns. -Each unique value (or combination of values) in the partition column(s) defines a distinct 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](data-distribution/partitioning.md) for creation strategies, multi-column partitions, and historical partition access. ### Bucket -A **bucket** horizontally divides the data of a table/partition into `N` buckets according to the bucketing policy. -The number of buckets `N` can be configured per table. A bucket is the smallest unit of data migration and backup. -The data of a bucket consists of a LogTablet and a (optional) KvTablet. + +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](data-distribution/bucketing.md) for assignment strategies and the rules for changing bucket counts. ### LogTablet -A **LogTablet** needs to be generated for each bucket of Log and Primary Key Tables. -For Log Tables, the LogTablet is both the primary table data and the log data. For Primary Key Tables, the LogTablet acts -as the log data for the primary table data. -- **Segment:** The smallest unit of log storage in the **LogTablet**. A segment consists of an **.index** file and a **.log** data file. -- **.index:** An `offset sparse index` that maps message relative offsets to their corresponding physical byte addresses in the .log file. -- **.log:** Compact arrangement of log data. + +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 -Each bucket of the Primary Key Table needs to generate a KvTablet. Underlying, each KvTablet corresponds to an embedded RocksDB instance. RocksDB is an LSM (log structured merge) engine which helps KvTablet support high-performance updates and lookup queries. + +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](../concepts/architecture.mdx#replication-and-recovery) 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 engine | Behavior | Example use | +| --- | --- | --- | +| [Default (LastRow)](merge-engines/default.md) | Applies the latest write; supports partial updates when configured. | Current order or profile state. | +| [FirstRow](merge-engines/first-row.md) | Retains the first record for a key. | Deduplicating repeated event identifiers. | +| [Versioned](merge-engines/versioned.md) | Uses a version column to accept records whose version is at least the stored version. | Handling out-of-order entity updates. | +| [Aggregation](merge-engines/aggregation.md) | Combines 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](virtual-tables.md) 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. + +| Mechanism | Controls | Design implication | +| --- | --- | --- | +| `table.log.ttl` | Retention 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-segments` | Cleanup of local log copies after upload to remote storage. | Local retention determines the hot working set; remote logs can remain readable longer. | +| `table.kv.ttl` | Expiration of individual primary key rows. | Current state can have a different lifetime from its changelog. | +| Auto-partition retention | Expiration 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](data-distribution/ttl.md) 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](data-formats.md) 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](../streaming-lakehouse/overview.mdx) 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. + +```sql title="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. + +```sql title="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](../engine-flink/ddl.md), [Spark DDL](../engine-spark/ddl.md), or the [Java Client](../apis/java/index.md) to create and evolve tables. The detailed chapters on [Table Types](table-types/pk-table.md), [Partitioning](data-distribution/partitioning.md), [Bucketing](data-distribution/bucketing.md), [Merge Engines](merge-engines/index.md), and [Retention](data-distribution/ttl.md) provide the constraints behind each design choice.