This topic describes how to create and manage Fluss tables using Flink SQL.
Limitations
This feature is available only in Realtime Compute for Apache Flink VVR 11.2 and later.
Create a primary key table
-
Log in to the Realtime Compute for Apache Flink console.
-
In the Actions column of the target workspace, click Console.
-
In the left-side navigation pane, choose .
-
Write and run the SQL statement.
CREATE TABLE `my-catalog`.`my_db`.`my_pk_table` (
shop_id BIGINT,
user_id BIGINT,
num_orders INT,
total_amount INT,
PRIMARY KEY (shop_id, user_id) NOT ENFORCED
) WITH (
'bucket.num' = '4'
);
You must plan the number of buckets when you create a table.
Although you can adjust the bucketing for a log table, the number of buckets for a primary key table cannot be changed after creation. Therefore, you must determine the number of buckets based on your data volume when you create the table. Use the following formula:
Number of buckets = Total data volume in a single partition / Target bucket size
Create a log table
CREATE TABLE `my-catalog`.`my_db`.`my_log_table` (
order_id BIGINT,
item_id BIGINT,
amount INT,
address STRING
) WITH (
'bucket.num' = '8'
);
Create a partitioned table
-
Fluss currently supports only the STRING data type for partition keys.
-
For a partitioned primary key table, the partition key (the dt field in the following example) must be a subset of the primary key.
Partitioned primary key table
The partition key is dt, and the primary key is [dt, shop_id, user_id].
CREATE TABLE `my-catalog`.`my_db`.`my_part_pk_table` (
dt STRING,
shop_id BIGINT,
user_id BIGINT,
num_orders INT,
total_amount INT,
PRIMARY KEY (dt, shop_id, user_id) NOT ENFORCED
) PARTITIONED BY (dt) WITH (
'bucket.num' = '4'
);
Partitioned log table
The following example creates a partitioned log table with dt as the partition key.
CREATE TABLE `my-catalog`.`my_db`.`my_part_log_table` (
dt STRING,
shop_id BIGINT,
user_id BIGINT,
num_orders INT,
total_amount INT
) PARTITIONED BY (dt) WITH (
'bucket.num' = '4'
);
Add a partition
ALTER TABLE my-catalog.my_db.my_part_pk_table
ADD PARTITION (dt = '2025-03-05');
Create an auto-partitioned table
An auto-partitioned table extends a partitioned table with automatic partitioning. Fluss creates partitions in advance based on your auto-partitioning policy.
Auto-partitioned primary key table
The partition key is dt, and the auto-partition interval is day.
CREATE TABLE `my-catalog`.`my_db`.`my_auto_part_pk_table` (
dt STRING,
shop_id BIGINT,
user_id BIGINT,
num_orders INT,
total_amount INT,
PRIMARY KEY (dt, shop_id, user_id) NOT ENFORCED
) PARTITIONED BY (dt) WITH (
'bucket.num' = '4'
'table.auto-partition.enabled' = 'true',
'table.auto-partition.time-unit' = 'day'
);
Auto-partitioned log table
The partition key is dt, and the auto-partition interval is day.
CREATE TABLE `my-catalog`.`my_db`.`my_auto_part_log_table` (
dt STRING,
shop_id BIGINT,
user_id BIGINT,
num_orders INT,
total_amount INT
) PARTITIONED BY (dt) WITH (
'bucket.num' = '4'
'table.auto-partition.enabled' = 'true',
'table.auto-partition.time-unit' = 'day'
);
Update a table
Add columns
This feature requires Realtime Compute for Apache Flink VVR 11.5 or later and Fluss 0.8-ali-3.0 or later. To avoid disrupting read and write operations, first upgrade your Flink job to VVR 11.5 or later, and then upgrade your Fluss cluster to 0.8-ali-3.0.
Fluss allows you to evolve a table's schema by adding new columns. You can add columns of any data type, including complex types such as ROW, MAP, and ARRAY.
This is a lightweight, metadata-only operation with the following benefits:
-
Zero data rewrite: You do not need to rewrite or migrate existing data files when you add columns.
-
Instant execution: The operation completes in milliseconds, regardless of the table size.
-
High availability: The table remains online and fully available during schema evolution, with no disruption to running clients.
This feature currently has the following limitations:
-
Column order: New columns are always appended to the end of the existing column list.
-
Nullability: Only nullable columns can be added to a table to ensure compatibility with existing data.
-
Nested fields: Adding fields to an existing nested
ROWtype is not currently supported. This type of operation falls under "updating column types" and will be supported in a future release.
Use the ALTER TABLE statement to add a single column or multiple columns.
-- Add a single column at the end of the table
ALTER TABLE `my-catalog`.`my_db`.`my_pk_table` ADD user_email STRING COMMENT 'User email address';
-- Add multiple columns at the end of the table
ALTER TABLE `my-catalog`.`my_db`.`my_pk_table` ADD (
user_email STRING COMMENT 'User email address',
order_quantity INT
);
Modify table properties
Use the ALTER TABLE SET statement to modify one or more connector options and storage options. New values overwrite existing ones.
The Fluss cluster applies storage option changes dynamically. You do not need to rebuild the table.
Supported options
-
All Read Options, Write Options, Lookup Options, and Other Options, except for
bootstrap.servers. -
The following Storage Options:
-
table.datalake.enabled: Enables or disables data lake and stream integration. -
table.datalake.freshness: Sets the data freshness for data lake storage.
-
-- Enable data lake and stream integration
ALTER TABLE `my-catalog`.`my_db`.`my_pk_table` SET ('table.datalake.enabled' = 'true');
-- Set data lake freshness to 5 minutes
ALTER TABLE `my-catalog`.`my_db`.`my_pk_table` SET ('table.datalake.freshness' = '5min');
After data lake and stream integration (table.datalake.enabled) is enabled on a table, lake format prefix options such as paimon.* can no longer be modified.
Reset table properties
Use the ALTER TABLE RESET statement to restore one or more connector options and storage options to their default values.
-- Reset data lake and stream integration to its default (disabled)
ALTER TABLE `my-catalog`.`my_db`.`my_pk_table` RESET ('table.datalake.enabled');
Clone a table (schema cloning)
The CREATE TABLE LIKE syntax creates a new table with the same schema, partition information, and properties as an existing table.
You must run this operation in a temporary query script. To open a new script, choose .
-- Create a temporary table named datagen
CREATE TEMPORARY TABLE datagen (
user_id BIGINT,
item_id BIGINT,
behavior STRING,
dt STRING,
hh STRING
) WITH (
'connector' = 'datagen',
'rows-per-second' = '10'
);
-- Create a Fluss table based on the datagen table, excluding its options
CREATE TABLE my-catalog.my_db.my_like_db LIKE datagen (EXCLUDING OPTIONS);
Drop a table
DROP TABLE `my-catalog`.`my_db`.`my_pk_table`;