All Products
Search
Document Center

Realtime Compute for Apache Flink:Stream storage Fluss

Last Updated:Aug 20, 2026

Learn how to use the Stream Storage Fluss connector.

Background information

Fluss uses a lake-stream integration architecture that seamlessly integrates with data lake formats such as Apache Paimon to provide low-latency data processing and significantly reduce the cost of building and maintaining a real-time data warehouse. Fluss helps enterprises break down data silos, accelerate deep analytics, and efficiently share data within their data lake by providing a unified mechanism for data streams and lakehouse synergy.

Category

Description

Supported types

Source tables, dimension tables, and result tables, as well as data ingestion sinks.

Execution mode

Streaming mode and batch mode.

Data format

Not supported yet.

API type

SQL and data ingestion YAML.

Support for result table updates and deletes

Yes.

Create a custom catalog

Fluss evolves faster than Realtime Compute for Apache Flink. The latest custom Catalog package is provided here for you to try Fluss Catalog features that have not yet been released with VVR versions. For more information, see Register a Fluss Catalog.

SQL

The Fluss connector can be used as a source table, dimension table, or result table in SQL jobs. For more information, see Engine Integration - Flink.

Prerequisites

  • You have read and write permissions for the target Fluss table. For more information, see Grant access to a Fluss cluster.

  • The Flink workspace and the Fluss cluster must be in the same VPC.

  • You have created a Fluss Catalog in the Flink console. For details, see Create a Fluss Catalog.

Key features

The Fluss Connector provides the following key features:

Feature

Description

unified full and incremental consumption

Supports full, incremental, and unified full and incremental consumption.

data exploration

Fluss supports fast data exploration in Flink SQL's batch mode, which is useful for data development, debugging, data validation, and troubleshooting.

CDC log subscription and virtual table

Fluss primary key tables offer native support for Change Data Capture (CDC) log subscription, letting downstream Flink jobs consume a complete stream of change events (INSERT, UPDATE_BEFORE, UPDATE_AFTER, and DELETE) in real time.

lookup join (dimension table join)

Its support for efficient key-value (KV) point lookups and secondary index lookups makes Fluss ideal for low-latency dimension table join scenarios.

partial update

Updates only the data in modified columns, not the entire row.

delta join

This novel join solution preserves dual-stream join semantics by pushing data update operations to the Fluss table. This process significantly reduces Flink's resource consumption and improves job stability and execution efficiency.

merge engine

Fluss primary key tables let you configure different merge engines to control how records with the same primary key are merged.

lake-stream integration

Provides deep integration of real-time stream storage (Fluss) with a Paimon-based data lake where metadata is managed by Data Lake Formation (DLF). This "one data, two views" model allows a single logical table to offer both low-latency real-time data access and high-throughput historical data analysis.

partitioned table

Supports creating partitioned tables, reading from and writing to them, and automatically handling partition creation and cleanup.

add columns

Lets you dynamically add columns without interrupting ongoing read and write jobs.

Limitations

  • Log tables only support INSERT operations, not UPDATE or DELETE.

  • Primary key tables support INSERT, UPDATE, and DELETE operations.

  • Lookup joins are supported only on dimension tables that are also primary key tables. The join condition must include all primary key fields.

  • For a Delta join, both tables must be Fluss primary key tables, and the join key must include the bucket key.

Considerations

  • Fluss tables come in two types: a log table, which has no primary key and is similar to a Kafka topic, and a primary key table, which has a primary key and supports upsert semantics. Choose the table type that best fits your use case.

  • In a production environment, properly configure bucket.num to ensure optimal parallel read and write performance and balanced data distribution.

  • For partitioned tables, design your partitioning strategy based on data volume and query patterns. Choose a suitable time granularity (e.g., daily or monthly) and use the table.auto-partition.num-retention parameter to control how many historical partitions are retained, preventing metadata bloat. If a primary key table contains data with natural partitioning characteristics (such as by time or business dimension), use a partitioned table instead of storing all data in a single non-partitioned table. Partitioned tables simplify data lifecycle management and enable partition pruning for streaming and batch queries, which significantly reduces unnecessary data scans and improves query performance.

  • Before using Delta Join, configure the required runtime parameters in your Flink job to enable this optimization.

Features

Full and incremental consumption

Overview

Fluss primary key tables support full and incremental consumption. The system first reads a full historical snapshot of the table. After the snapshot read is complete, it automatically and seamlessly switches to incremental stream consumption. This process ensures data completeness and consistency.

Consumption modes

You can control the consumption starting point using the scan.startup.mode parameter.

Mode

Description

full

Default. Reads a full snapshot first, then consumes incremental data. This mode is suitable for scenarios that require complete historical data.

earliest

Consumes data from the earliest available log offset without reading a full snapshot.

latest

Consumes data from the latest offset. Only new data generated after the job starts is consumed.

timestamp

Consumes data from a specific timestamp. This mode requires you to also set the scan.startup.timestamp parameter.

SQL examples
-- Full and incremental consumption (default mode)
SELECT * FROM `fluss_catalog`.`fluss_db`.orders;

-- Consume from a specific timestamp by using a hint
SELECT * FROM `fluss_catalog`.`fluss_db`.orders
  /*+ OPTIONS('scan.startup.mode' = 'timestamp', 'scan.startup.timestamp' = '2025-01-01 00:00:00') */;

-- Consume from the latest offset
SELECT * FROM `fluss_catalog`.`fluss_db`.orders
  /*+ OPTIONS('scan.startup.mode' = 'latest') */;

CDC log subscription and virtual tables

Log subscription

Fluss primary key tables natively support CDC (Change Data Capture) log subscription. This allows downstream Flink jobs to consume a complete changelog stream (INSERT, UPDATE_BEFORE, UPDATE_AFTER, and DELETE events) in real time without requiring a separate binlog configuration.

This means:

  • All changes to a primary key table automatically generate CDC logs.

  • Downstream applications can directly subscribe to the primary key table to receive the complete changelog stream for real-time data synchronization or streaming ETL.

SQL example
-- Read the primary key table directly to get the CDC changelog stream
SELECT * FROM `fluss_catalog`.`fluss_db`.orders;
-- The output includes complete INSERT, UPDATE, and DELETE events.
Virtual tables

In addition to subscribing to the CDC stream directly, Fluss provides a virtual table mechanism. You can access changelogs with metadata by appending a specific suffix to the base table name, which eliminates the need for extra data storage. The system generates these virtual tables automatically, eliminating the need for manual creation.

Changelog table ($changelog)

This feature applies to primary key tables and append-only tables and provides the raw changelog stream with metadata.

Table structure (The system automatically adds the following metadata columns to the original table schema):

Metadata column

Type

Description

_change_type

STRING

The type of change. For a primary key table, possible values are insert, update_before, update_after, and delete. For an append-only table, the only possible value is insert.

_log_offset

BIGINT

The offset of the change in the log.

_commit_timestamp

TIMESTAMP_LTZ

The commit timestamp of the change.

SQL examples:

-- Read the $changelog virtual table for the orders table
SELECT * FROM `fluss_catalog`.`fluss_db`.orders$changelog;

-- Consume from the latest offset by using a hint
SELECT * FROM `fluss_catalog`.`fluss_db`.orders$changelog
  /*+ OPTIONS('scan.startup.mode' = 'latest') */;

-- Filter for delete events
SELECT * FROM `fluss_catalog`.`fluss_db`.orders$changelog
WHERE _change_type = 'delete';
Binlog table ($binlog)

This feature is available only for primary key tables and provides Binlog-formatted data that includes before/after images, which makes it easy to accurately retrieve the state of the data before and after each change.

Table structure (In addition to metadata columns, the full before and after row data is presented in a nested ROW type):

Column

Type

Description

_change_type

STRING

The type of change. Possible values are insert, update, and delete.

_log_offset

BIGINT

The log offset.

_commit_timestamp

TIMESTAMP_LTZ

The commit timestamp.

before

ROW

The row data before the change. The value is NULL for insert operations.

after

ROW

The row data after the change. The value is NULL for delete operations.

Change event mapping:

Change type

before

after

insert

NULL

New value

update

Old value

New value

delete

Old value

NULL

SQL examples:

-- Read the $binlog virtual table for the users table
SELECT * FROM `fluss_catalog`.`fluss_db`.users$binlog;

-- Get the before and after values of a field from an update event
SELECT
  `before`.user_name AS old_name,
  `after`.user_name  AS new_name,
  _commit_timestamp  AS change_time
FROM `fluss_catalog`.`fluss_db`.users$binlog
WHERE _change_type = 'update';
Limitations
  • The $changelog and $binlog tables currently do not support projection pushdown, partition pruning, or predicate pushdown.

  • The $binlog table is available only for primary key tables. It is not supported for append-only tables.

  • The supported startup modes for virtual tables are earliest (default), latest, and timestamp. The full mode is not supported.


Partial update

Overview

Fluss supports partial updates, which allow multiple data streams to independently write to their respective columns in the same wide primary key table. The Fluss engine automatically merges fields at the record level based on the primary key. This approach avoids the need to maintain complex multi-stream join state in Flink and prevents issues related to large state sizes.

Use case

A typical use case is building real-time user profile wide tables. Different data sources, such as recommendation systems, impression logs, click logs, and shopping cart data, can update their respective fields in the user profile table. This process creates a complete 360-degree view of the user.

Limitations
  • The target table must be a primary key table.

  • Each data stream's write operation must include all primary key fields.

  • The computation topology must not contain multi-stream join nodes. Fluss handles the merge operation directly at the storage layer.

SQL examples
-- 1. Create a wide table
CREATE TABLE IF NOT EXISTS `fluss_catalog`.fluss_db.user_rec_wide (
    user_id   STRING,
    item_id   STRING,
    rec_score DOUBLE,
    imp_cnt   INT,
    click_cnt INT,
    PRIMARY KEY (user_id, item_id) NOT ENFORCED
) WITH ('bucket.num' = '3');

-- 2. Use multiple streams to write to their respective fields
-- Stream 1: Writes recommendation scores
INSERT INTO `fluss_catalog`.fluss_db.user_rec_wide (user_id, item_id, rec_score)
SELECT user_id, item_id, rec_score 
FROM `fluss_catalog`.fluss_db.recommendations;

-- Stream 2: Writes impression counts
INSERT INTO `fluss_catalog`.fluss_db.user_rec_wide (user_id, item_id, imp_cnt)
SELECT user_id, item_id, imp_cnt 
FROM `fluss_catalog`.fluss_db.impressions;

-- Stream 3: Writes click counts
INSERT INTO `fluss_catalog`.fluss_db.user_rec_wide (user_id, item_id, click_cnt)
SELECT user_id, item_id, click_cnt 
FROM `fluss_catalog`.fluss_db.clicks;

The three INSERT statements above can run as separate Flink jobs. The Fluss storage layer automatically merges the fields for the same primary key into complete rows.


Delta Join

Overview

Delta Join is a next-generation stream-stream join paradigm provided by Fluss. It offloads data update operations to the Fluss storage layer, leveraging primary key tables and secondary indexes for efficient lookups. This design eliminates the need to maintain join state in Flink.

Key advantages:

  • No join state: Eliminates redundant state storage and avoids garbage collection (GC) issues caused by large state sizes.

  • Low resource consumption: Reduces Flink memory and CPU consumption by over 86% in tests.

  • Greater stability and efficiency: Avoids performance bottlenecks caused by state bloat.

Comparison with traditional joins

Feature

Traditional Flink join (with Kafka)

Fluss Delta Join

State Storage

Requires caching all upstream data in Flink.

No join state; relies on Fluss indexes.

Resource Consumption

High, due to frequent GC caused by large state.

Low

Stability

Prone to issues from state bloat.

More stable

Mechanism

Updates from both streams drive the join.

Fluss offloads updates to its storage layer.

Use Cases

General-purpose stream-stream joins.

Efficient joins between Fluss primary key tables.

Limitations
  • Both the left and right tables must be Fluss primary key tables (partitioned tables are supported).

  • The bucket key must be a prefix of the primary key.

  • The join key must include the bucket key. For a partitioned table, it must also include the partition key.

SQL examples

1. Create tables

CREATE TABLE orders (
  user_id BIGINT,
  order_date DATE,
  order_id BIGINT,
  amount DECIMAL(10, 2),
  PRIMARY KEY (user_id, order_id, order_date) NOT ENFORCED
) WITH ('bucket.key' = 'user_id, order_id');

CREATE TABLE order_enhance (
  user_id BIGINT,
  order_id BIGINT,
  risk_level TINYINT,
  is_fraud BOOLEAN,
  update_time TIMESTAMP(3),
  PRIMARY KEY (user_id, order_id) NOT ENFORCED
) WITH ('bucket.key' = 'user_id, order_id');

2. Run a Delta Join

SELECT * FROM orders o
INNER JOIN order_enhance e
ON o.user_id = e.user_id AND o.order_id = e.order_id;

3. Optimize lookup performance with hints

SELECT * FROM orders o /*+ OPTIONS('client.lookup.queue-size'='2560') */
INNER JOIN order_enhance e /*+ OPTIONS('client.lookup.queue-size'='2560') */
ON o.user_id = e.user_id AND o.order_id = e.order_id;
Job parameter configuration

To use a Delta Join, add the following configurations to the runtime parameters of your Flink job:

table.exec.delta-join.cache-enabled: 'true'
table.exec.async-lookup.buffer-capacity: '3000'
table.optimizer.delta-join.strategy: FORCE
table.exec.delta-join.left.cache-size: '100'
table.exec.delta-join.right.cache-size: '1000'

Parameter

Description

table.optimizer.delta-join.strategy

Set to FORCE to enforce Delta Join. The job fails if the conditions are not met.

table.exec.async-lookup.buffer-capacity

The buffer capacity for asynchronous requests. We recommend setting this value in the thousands.

table.exec.delta-join.cache-enabled

Enables local caching to reduce requests to Fluss.

table.exec.delta-join.left.cache-size

The number of keys to cache for the left table. Adjust this value based on memory and hot data patterns.

table.exec.delta-join.right.cache-size

The number of keys to cache for the right table. Adjust this value based on memory and hot data patterns.

After you enable Delta Join, the Delta Join node appears in the Flink job overview, indicating the optimization is active.


Lookup join (dimension table join)

Overview

You can use a Fluss primary key table as a Flink SQL dimension table. A lookup join allows you to enrich streaming data with dimensional data in real time. Fluss supports efficient KV point lookups and secondary index searches, making it suitable for low-latency dimension table join scenarios.

Limitations
  • Only a primary key table can be used as a dimension table.

  • The join condition must include all primary key fields of the dimension table.

Synchronous and asynchronous modes

Mode

Configuration

Description

Asynchronous (default)

'lookup.async' = 'true'

Provides high throughput. Recommended for most use cases.

Synchronous

'lookup.async' = 'false'

Queries records sequentially. Suitable for debugging or low-concurrency scenarios.

SQL examples
-- Configure lookup cache by using a hint to perform a dimension table join
SELECT
  o.order_id,
  o.amount,
  u.user_name,
  p.product_name
FROM `fluss_catalog`.`fluss_db`.orders AS o
LEFT JOIN `fluss_catalog`.`fluss_db`.user_dim
  /*+ OPTIONS('lookup.cache' = 'PARTIAL', 'lookup.partial-cache.expire-after-write' = '60s', 'lookup.partial-cache.max-rows' = '10000') */
  FOR SYSTEM_TIME AS OF o.proctime AS u
  ON o.user_id = u.user_id
LEFT JOIN `fluss_catalog`.`fluss_db`.product_dim
  FOR SYSTEM_TIME AS OF o.proctime AS p
  ON o.product_id = p.product_id;

-- Use synchronous lookup mode (suitable for debugging)
SELECT o.order_id, o.amount, d.user_name
FROM `fluss_catalog`.`fluss_db`.orders AS o
LEFT JOIN `fluss_catalog`.`fluss_db`.user_dim
  /*+ OPTIONS('lookup.async' = 'false') */
  FOR SYSTEM_TIME AS OF o.proctime AS d
  ON o.user_id = d.user_id;

Merge engine

Overview

Fluss primary key tables let you configure a merge engine to control how the system merges records with the same primary key. Use the table.merge-engine parameter to select one of the following strategies:

Merge engine

Description

last_row (default)

Keeps the last written record. Later data overwrites earlier data. This engine is suitable for most upsert scenarios.

first_row

Keeps the first written record, ignoring any subsequent data with the same primary key. This engine is suitable for deduplication scenarios, such as log deduplication.

versioned

Keeps the record with the highest version number from a specified version column. Requires the table.merge-engine.versioned.ver-column parameter to specify the version column. This engine is suitable for scenarios that require controlling the data update order based on versions.

SQL examples
-- Use first_row for deduplication
CREATE TABLE IF NOT EXISTS `fluss_catalog`.fluss_db.event_dedup (
    event_id STRING,
    event_type STRING,
    event_time TIMESTAMP(3),
    payload STRING,
    PRIMARY KEY (event_id) NOT ENFORCED
) WITH (
  'table.merge-engine' = 'first_row'
);

-- Use versioned to update based on a version column
CREATE TABLE IF NOT EXISTS `fluss_catalog`.fluss_db.user_profile (
    user_id BIGINT,
    user_name STRING,
    age INT,
    update_version BIGINT,
    PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
  'table.merge-engine' = 'versioned',
  'table.merge-engine.versioned.ver-column' = 'update_version'
);

Partitioned table

Overview

Fluss supports partitioned tables. You can use PARTITIONED BY to divide data into different physical partitions based on a partition key. This is ideal for managing data lifecycles and improving query performance based on time or business dimensions. Partitioned tables support single-field partitioning and multi-field partitioning, and offer three orthogonal partition management strategies: manual management, automatic management, and dynamic creation.

Basic rules
  • The data type of a partition key must be STRING.

  • For a primary key table, the partition key must be a subset of the primary key.

  • You can only configure automatic partitioning rules when creating the table; you cannot modify them later.

Partition management strategies

Strategy

Description

Manual management

Manually create or delete partitions using ALTER TABLE. This strategy is suitable for scenarios with a limited and controllable number of partitions.

Automatic management

The system automatically creates new partitions and cleans up expired historical partitions based on time-based rules. This strategy is suitable for continuous, time-based data ingestion.

Dynamic creation

When data is written, the client automatically creates the target partition if it does not exist. This strategy is suitable for scenarios with unpredictable partition values.

These three strategies are orthogonal and can be used in combination. For example, you can enable both automatic management for time-based creation and cleanup, and dynamic creation as a fallback to handle partitions that were not pre-created.

SQL examples

1. Create a single-field partitioned table with automatic partition management

CREATE TABLE `fluss_catalog`.fluss_db.site_access (
  site_id BIGINT,
  user_name STRING,
  pv BIGINT,
  event_day STRING,
  PRIMARY KEY (site_id, event_day) NOT ENFORCED
) PARTITIONED BY (event_day) WITH (
  'table.auto-partition.enabled' = 'true',
  'table.auto-partition.time-unit' = 'DAY',
  'table.auto-partition.num-precreate' = '5',
  'table.auto-partition.num-retention' = '30',
  'table.auto-partition.time-zone' = 'Asia/Shanghai'
);

2. Create a multi-field partitioned table

CREATE TABLE `fluss_catalog`.fluss_db.region_access (
  access_id BIGINT,
  user_name STRING,
  event_date STRING,
  region STRING,
  PRIMARY KEY (access_id, event_date, region) NOT ENFORCED
) PARTITIONED BY (event_date, region) WITH (
  'table.auto-partition.enabled' = 'true',
  'table.auto-partition.key' = 'event_date',
  'table.auto-partition.time-unit' = 'DAY',
  'table.auto-partition.num-retention' = '7'
);
Note: For multi-field partitioning, you must specify the time-based key for automatic partitioning using table.auto-partition.key. Multi-field partitioning supports only automatic expiration and cleanup, not automatic pre-creation. The num-precreate parameter is forced to 0.

3. Manually manage partitions

-- View all partitions
SHOW PARTITIONS `fluss_catalog`.fluss_db.site_access;

-- Manually add a partition
ALTER TABLE `fluss_catalog`.fluss_db.site_access
  ADD PARTITION (event_day = '20250315');

-- Manually drop a partition
ALTER TABLE `fluss_catalog`.fluss_db.site_access
  DROP PARTITION (event_day = '20250101');

4. Dynamically create partitions

When you write data, if the target partition does not exist, the client automatically creates it. This feature is enabled by default.

-- Write data. If the partition '20250701' does not exist, it is automatically created.
INSERT INTO `fluss_catalog`.fluss_db.site_access
VALUES (1, 'hello', 100, '20250701');

The client.writer.dynamic-create-partition.enabled client parameter (default: true) controls dynamic creation. It is subject to the cluster-level limits of max.partition.num (default: 1000) and max.bucket.num (default: 128000).

Partition pruning

When you query a partitioned table in streaming or batch mode, Fluss supports partition pruning. By filtering on the partition key in a WHERE clause, the query reads only the matching partitions, which reduces unnecessary I/O overhead.

-- Read only the data from the partition for a specific date
SELECT * FROM `fluss_catalog`.`fluss_db`.site_access
WHERE event_day = '20250315';

-- Filter a multi-field partitioned table by region
SELECT * FROM `fluss_catalog`.`fluss_db`.region_access
WHERE region = 'US';
Automatic partitioning parameters

Parameter

Type

Default

Description

table.auto-partition.enabled

Boolean

false

Specifies whether to enable automatic partition management.

table.auto-partition.key

String

None

For multi-field partitioning, specifies the time-based key for automatic partitioning. This parameter is not needed for single-field partitioning.

table.auto-partition.time-unit

Enum

DAY

The time granularity for automatic partitioning. Valid values are HOUR, DAY, MONTH, QUARTER, and YEAR.

table.auto-partition.num-precreate

Integer

2

The number of future partitions to pre-create. Not supported for multi-field partitioning, where the value is forced to 0.

table.auto-partition.num-retention

Integer

7

The number of historical partitions to retain. Expired partitions beyond this limit are automatically deleted.

table.auto-partition.time-zone

String

System time zone

The time zone used for automatic partitioning.

Best practices
  • Choose the right partition granularity: Select a time granularity based on your data volume and query patterns. For large data volumes and narrow query time ranges, we recommend partitioning by DAY or HOUR. For smaller data volumes, you can partition by MONTH or QUARTER.

  • Set a reasonable retention policy: Use num-retention to control the number of historical partitions. This helps avoid metadata bloat and reduces the cluster management burden from an excessive number of partitions.

  • Combine automatic management and dynamic creation: Use automatic management for routine time-based partition creation and cleanup, and use dynamic creation as a fallback to handle abnormal or delayed data.

  • Be aware of cluster limits: The max.partition.num parameter (default: 1000) limits the number of partitions per table. The total number of buckets in the cluster is limited by max.bucket.num (default: 128,000). The total number of buckets in a table (number of partitions × bucket.num) must not exceed this limit.


Streaming query pushdown

Fluss uses the Apache Arrow columnar storage format for log transmission. This supports pushing down filter and projection operations to the storage layer, which reduces data transfer volume and can improve query performance by up to 10x.

This feature is especially effective in the following scenarios:

  • Reading only a subset of columns from a wide table (projection pushdown).

  • Filtering a large amount of data based on specific conditions (predicate pushdown).

The log storage format is controlled by the table.log.format parameter, which defaults to ARROW.


Unified stream-lake architecture

Overview

Fluss supports a unified stream-lake architecture that can automatically synchronize real-time data from Fluss to a data lake, such as Apache Paimon. This achieves automatic tiered storage for hot and cold data, balancing the low latency of real-time queries with the low storage costs of historical data.

Key features:

  • Automatically writes streaming data from Fluss to Paimon lake storage.

  • Automatically tiers hot and cold data. Hot data is retained in Fluss stream storage, while cold data is moved to Paimon.

  • Provides a unified Flink SQL query interface to access data in both stream storage and lake storage simultaneously.

Configuration parameters

Parameter

Type

Default

Description

table.datalake.enabled

Boolean

false

Enables lake storage tiering.

table.datalake.format

Enum

Cluster default

The data lake format. Currently, only paimon is supported.

table.log.tiered.local-segments

Integer

2

The number of log segments to retain locally after tiered storage is enabled.


Data exploration

Overview

Fluss supports quick data exploration on tables in the Flink SQL batch mode, including LIMIT queries, primary key point lookups, and COUNT(*) aggregations. You can quickly view data in a table without starting a continuously running streaming job. This is useful for development, debugging, data validation, and troubleshooting.

LIMIT queries

You can run LIMIT queries on primary key tables and append-only tables to quickly retrieve a small sample of data.

-- View the first 10 records from a primary key table
SELECT * FROM `fluss_catalog`.`fluss_db`.orders LIMIT 10;

-- View the first 5 records from an append-only table
SELECT * FROM `fluss_catalog`.`fluss_db`.events LIMIT 5;

-- Combine with projection to view only specific columns
SELECT order_id, amount, order_time 
FROM `fluss_catalog`.`fluss_db`.orders LIMIT 20;
Primary key point lookups

For primary key tables, you can perform precise point lookups by specifying the primary key in a WHERE clause. This leverages the Fluss KV storage engine to achieve millisecond-level response times.

-- Perform a point lookup for a single record by its primary key
SELECT * FROM `fluss_catalog`.`fluss_db`.orders
WHERE order_id = 1001;

-- Perform a point lookup with a composite primary key
SELECT * FROM `fluss_catalog`.`fluss_db`.user_orders
WHERE user_id = 12345 AND order_id = 6789;
COUNT(*) aggregations

You can run a COUNT(*) query on a primary key table to quickly get the total number of records.

-- Count the total number of rows in a primary key table
SELECT COUNT(*) FROM `fluss_catalog`.`fluss_db`.orders;

WITH parameters

Storage parameters

The following parameters are persisted to the table's metadata.

Parameter

Type

Default value

Required

Description

bucket.num

Integer

None

No

The number of buckets for the Fluss table. This parameter affects data read and write parallelism. The default value is inherited from the cluster's default.bucket.number configuration.

bucket.key

String

None

No

The columns used for hash-based data distribution. These columns must be a subset of the primary key. Use a comma to separate multiple columns.

table.log.ttl

Duration

7d

No

The retention period for log data.

table.replication.factor

Integer

Cluster default

No

The replication factor for the log table. This value cannot be greater than the number of tablet servers.

table.log.format

Enum

ARROW

No

The storage format for the log. Options include ARROW (a columnar format that supports query pushdown) and INDEXED.

table.log.arrow.compression.type

Enum

ZSTD

No

The compression type for the ARROW format. Valid options are NONE, LZ4_FRAME, and ZSTD.

table.log.arrow.compression.zstd.level

Integer

3

No

The ZSTD compression level (1-22).

table.kv.format

Enum

COMPACTED

No

The KV storage format. Valid options are COMPACTED or INDEXED.

table.merge-engine

Enum

last_row

No

The merge strategy for the primary key table. Valid options are first_row, versioned, and last_row.

table.merge-engine.versioned.ver-column

String

None

Conditionally required

Required if table.merge-engine is set to versioned. Specifies the name of the version column.

table.datalake.enabled

Boolean

false

No

Enables or disables data lake storage tiering.

table.datalake.format

Enum

Cluster default

No

The data lake format. Currently, only paimon is supported.

table.log.tiered.local-segments

Integer

2

No

The number of log segments to retain locally after tiered storage is enabled.

Auto-partitioning parameters

Parameter

Type

Default value

Required

Description

table.auto-partition.enabled

Boolean

false

No

Enables or disables automatic partition creation.

table.auto-partition.time-unit

Enum

DAY

No

The time granularity for auto-partitioning. Valid options are YEAR, QUARTER, MONTH, DAY, and HOUR.

table.auto-partition.num-precreate

Integer

2

No

The number of future partitions to pre-create during checks.

table.auto-partition.num-retention

Integer

7

No

The number of historical partitions to retain. Older partitions are automatically deleted.

table.auto-partition.time-zone

String

System time zone

No

The time zone used for auto-partitioning.

Read parameters

The following parameters can be configured via the WITH clause or SQL hints.

Parameter

Type

Default value

Required

Description

scan.startup.mode

Enum

full

No

The consumption startup mode. Valid options are full, earliest, latest, and timestamp.

scan.startup.timestamp

Long/String

None

Conditionally required

The startup timestamp, which can be a millisecond value or a string in the yyyy-MM-dd HH:mm:ss format. This parameter applies only when scan.startup.mode is set to timestamp.

scan.partition.discovery.interval

Duration

1 minute

No

The interval for automatically discovering new partitions. Set to a negative value to disable automatic discovery.

client.scanner.log.check-crc

Boolean

true

No

Specifies whether to perform a CRC32 check on messages to ensure data integrity.

client.scanner.log.max-poll-records

Integer

500

No

The maximum number of records that a single poll() call can return.

client.scanner.log.fetch.max-bytes

MemorySize

16mb

No

The maximum number of bytes to fetch from the server per request.

client.scanner.log.fetch.max-bytes-for-bucket

MemorySize

1mb

No

The maximum number of bytes to fetch per bucket in a single request.

client.scanner.log.fetch.min-bytes

MemorySize

1b

No

The minimum data size the server should return for a fetch request.

client.scanner.log.fetch.wait-max-time

Duration

500ms

No

The maximum time that the server waits when the minimum number of bytes is not met.

client.scanner.io.tmpdir

String

System temporary directory

No

The local directory for storing temporary files such as snapshots and log segments.

client.scanner.remote-log.prefetch-num

Integer

4

No

The number of remote log segments to prefetch.

client.remote-file.download-thread-num

Integer

3

No

The number of threads used to download remote files.

Write parameters

The following parameters can be configured via the WITH clause or SQL hints.

Parameter

Type

Default value

Required

Description

sink.ignore-delete

Boolean

false

No

Specifies whether to ignore DELETE and UPDATE_BEFORE messages from the source.

sink.bucket-shuffle

Boolean

true

No

Specifies whether to shuffle data by bucket ID before writing, which can improve write performance.

client.writer.buffer.memory-size

MemorySize

64mb

No

The total memory size of the internal buffer for row data.

client.writer.batch-size

MemorySize

2mb

No

The target batch size for records within the same bucket.

client.writer.buffer.wait-timeout

Duration

Long.MAX

No

The maximum time the writer blocks while waiting for an available segment.

client.writer.batch-timeout

Duration

100ms

No

The maximum time to buffer records before sending them in a batch.

client.writer.bucket.no-key-assigner

Enum

STICKY

No

The bucket assignment strategy for tables without a primary key. Valid options are STICKY or ROUND_ROBIN.

client.writer.acks

String

all

No

The acknowledgment level for writes. all (or -1): Waits for all replicas to acknowledge. 1: Waits only for the leader to acknowledge. 0: Does not wait for any acknowledgment. The all setting is recommended for maximum data durability.

client.writer.request-max-size

MemorySize

10mb

No

The maximum size of a single request in bytes.

client.writer.retries

Integer

MAX_VALUE

No

The number of times to retry a send that failed with a retriable error.

client.writer.enable-idempotence

Boolean

true

No

Enables idempotent writes, which provide exactly-once semantics and message ordering guarantees.

client.writer.max-inflight-requests-per-bucket

Integer

5

No

The maximum number of in-flight requests per bucket. This parameter applies only when idempotent writes are enabled.

Dimension table parameters

Parameter

Type

Default value

Required

Description

lookup.async

Boolean

true

No

Enables asynchronous lookup, which can improve throughput for dimension table joins.

lookup.cache

Enum

NONE

No

The cache policy. Valid options are NONE (no caching) or PARTIAL (partial caching).

lookup.max-retries

Integer

3

No

The maximum number of retries if a lookup fails.

lookup.partial-cache.expire-after-access

Duration

None

No

The cache expiration period after the last access.

lookup.partial-cache.expire-after-write

Duration

None

No

The cache expiration period after the last write.

lookup.partial-cache.cache-missing-key

Boolean

true

No

Specifies whether to cache keys that do not exist in the dimension table (i.e., cache null values).

lookup.partial-cache.max-rows

Long

None

No

The maximum number of rows in the cache.

client.lookup.queue-size

Integer

25600

No

The size of the queue for pending lookup operations.

client.lookup.max-batch-size

Integer

128

No

The maximum batch size for combining multiple lookup operations.

client.lookup.max-inflight-requests

Integer

128

No

The maximum number of in-flight lookup requests.

client.lookup.batch-timeout

Duration

100ms

No

The maximum time to wait for a lookup batch to fill before it is sent.

Data ingestion

The Fluss connector acts as a data sink in a data ingestion YAML job.

Syntax

source:
   type: xxx

sink:
  type: fluss
  name: Fluss Sink
  bootstrap.servers: localhost:9123
  properties.client.security.protocol: sasl
  properties.client.security.sasl.mechanism: PLAIN
  properties.client.security.sasl.username: developer
  properties.client.security.sasl.password: developer-pass
  

pipeline:
  schema.change.behavior: IGNORE

Schema change

Fluss Server supports adding new columns but does not support other column operations (such as deleting columns, renaming columns, or changing column types). Therefore, we recommend using the default Flink CDC schema.change.behavior = LENIENT setting:

  • When a column is renamed, the connector sends two events. The connector does not delete the original column. Instead, it changes the column's data type to nullable and adds a new column with the new name and a nullable data type.

  • When a column is dropped, the connector sends a column type change event, changing the column's data type to nullable.

  • When a column is added, the connector sends an add column event but changes the column's data type to nullable.

Note

This feature requires VVR 11.6 or a later version.

Parameters

Parameter

Description

Required

Type

Default

Remarks

type

Specifies the sink type.

Yes

String

None

Fixed value: fluss.

name

Specifies the sink name.

No

String

None

None.

bootstrap.servers

Specifies the address of the Fluss Server.

Yes

String

None

Format: host:port,host:port,host:port. Use commas (,) to separate multiple addresses.

bucket.key

Specifies the bucket key.

No

String

None

Specify the data distribution policy for each Fluss table. Tables are separated by ';', and bucket keys are separated by ','.

Format: database1.table1:key1,key2;database1.table2:key3.

Data is distributed to buckets based on the hash value of the bucket key. The bucket key must be a subset of the primary key and cannot include the partition key of a primary key table. If a table has a primary key but no bucket key is specified, the bucket key defaults to the primary key (excluding the partition key). If a table has no primary key and no bucket key is specified, data is randomly distributed to the buckets.

The values of bucket.key and bucket.num cannot be changed after the table is created.

bucket.num

Specifies the number of buckets for Fluss tables.

No

String

None

The bucket number for each Fluss table. The numbers are separated by a ';'.

Format: database1.table1:4;database1.table2:8.

If a table does not have a configured bucket number, the server-side configuration item default.bucket.number is used by default.

The values of bucket.key and bucket.num cannot be changed after the table is created.

properties.table.*

Fluss table properties.

No

String

None

Passes parameters supported by Fluss tables to the pipeline. For more information, see Fluss table options.

properties.client.*

Fluss client parameters.

No

String

None

Passes parameters supported by the Fluss client to the pipeline. For more information, see Fluss client options.

Type mapping

The following table shows the data type mapping for data ingestion.

CDC field type

Fluss field type

TINYINT

TINYINT

SMALLINT

SMALLINT

INT

INT

BIGINT

BIGINT

FLOAT

FLOAT

DOUBLE

DOUBLE

DECIMAL(p, s)

DECIMAL(p, s)

BOOLEAN

BOOLEAN

DATE

DATE

TIME

TIME

TIMESTAMP

TIMESTAMP

TIMESTAMP_LTZ

TIMESTAMP_LTZ

CHAR(n)

CHAR(n)

VARCHAR(n)

STRING

ARRAY

ARRAY

MAP

MAP

ROW

ROW

BINARY(n)

Unsupported

VARBINARY(n)