All Products
Search
Document Center

Realtime Compute for Apache Flink:Manage tables

Last Updated:Aug 10, 2026

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

  1. Log in to the Realtime Compute for Apache Flink console.

  2. In the Actions column of the target workspace, click Console.

  3. In the left-side navigation pane, choose Development > Scripts > New Temporary Script.

  4. 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'
);
Note

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

Note

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 ROW type 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');
Important

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 Development > Scripts > New Temporary Script.

-- 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`;