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 |
|
Supports full, incremental, and unified full and incremental consumption. |
|
|
Fluss supports fast data exploration in Flink SQL's batch mode, which is useful for data development, debugging, data validation, and troubleshooting. |
|
|
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. |
|
|
Its support for efficient key-value (KV) point lookups and secondary index lookups makes Fluss ideal for low-latency dimension table join scenarios. |
|
|
Updates only the data in modified columns, not the entire row. |
|
|
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. |
|
|
Fluss primary key tables let you configure different merge engines to control how records with the same primary key are merged. |
|
|
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. |
|
|
Supports creating partitioned tables, reading from and writing to them, and automatically handling partition creation and cleanup. |
|
|
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.numto 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-retentionparameter 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 |
|
|
Default. Reads a full snapshot first, then consumes incremental data. This mode is suitable for scenarios that require complete historical data. |
|
|
Consumes data from the earliest available log offset without reading a full snapshot. |
|
|
Consumes data from the latest offset. Only new data generated after the job starts is consumed. |
|
|
Consumes data from a specific timestamp. This mode requires you to also set the |
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 |
|
|
STRING |
The type of change. For a primary key table, possible values are |
|
|
BIGINT |
The offset of the change in the log. |
|
|
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 |
|
|
STRING |
The type of change. Possible values are |
|
|
BIGINT |
The log offset. |
|
|
TIMESTAMP_LTZ |
The commit timestamp. |
|
|
ROW |
The row data before the change. The value is NULL for |
|
|
ROW |
The row data after the change. The value is NULL for |
Change event mapping:
|
Change type |
|
|
|
|
NULL |
New value |
|
|
Old value |
New value |
|
|
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, andtimestamp. Thefullmode 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 |
|
|
Set to |
|
|
The buffer capacity for asynchronous requests. We recommend setting this value in the thousands. |
|
|
Enables local caching to reduce requests to Fluss. |
|
|
The number of keys to cache for the left table. Adjust this value based on memory and hot data patterns. |
|
|
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) |
|
Provides high throughput. Recommended for most use cases. |
|
Synchronous |
|
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 |
|
|
Keeps the last written record. Later data overwrites earlier data. This engine is suitable for most upsert scenarios. |
|
|
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. |
|
|
Keeps the record with the highest version number from a specified version column. Requires the |
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 |
|
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 usingtable.auto-partition.key. Multi-field partitioning supports only automatic expiration and cleanup, not automatic pre-creation. Thenum-precreateparameter 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 |
|
|
Boolean |
false |
Specifies whether to enable automatic partition management. |
|
|
String |
None |
For multi-field partitioning, specifies the time-based key for automatic partitioning. This parameter is not needed for single-field partitioning. |
|
|
Enum |
DAY |
The time granularity for automatic partitioning. Valid values are |
|
|
Integer |
2 |
The number of future partitions to pre-create. Not supported for multi-field partitioning, where the value is forced to 0. |
|
|
Integer |
7 |
The number of historical partitions to retain. Expired partitions beyond this limit are automatically deleted. |
|
|
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
DAYorHOUR. For smaller data volumes, you can partition byMONTHorQUARTER. -
Set a reasonable retention policy: Use
num-retentionto 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.numparameter (default: 1000) limits the number of partitions per table. The total number of buckets in the cluster is limited bymax.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 |
|
|
Boolean |
false |
Enables lake storage tiering. |
|
|
Enum |
Cluster default |
The data lake format. Currently, only |
|
|
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 |
|
|
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 |
|
|
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. |
|
|
Duration |
7d |
No |
The retention period for log data. |
|
|
Integer |
Cluster default |
No |
The replication factor for the log table. This value cannot be greater than the number of tablet servers. |
|
|
Enum |
ARROW |
No |
The storage format for the log. Options include |
|
|
Enum |
ZSTD |
No |
The compression type for the ARROW format. Valid options are |
|
|
Integer |
3 |
No |
The ZSTD compression level (1-22). |
|
|
Enum |
COMPACTED |
No |
The KV storage format. Valid options are |
|
|
Enum |
last_row |
No |
The merge strategy for the primary key table. Valid options are |
|
|
String |
None |
Conditionally required |
Required if |
|
|
Boolean |
false |
No |
Enables or disables data lake storage tiering. |
|
|
Enum |
Cluster default |
No |
The data lake format. Currently, only |
|
|
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 |
|
|
Boolean |
false |
No |
Enables or disables automatic partition creation. |
|
|
Enum |
DAY |
No |
The time granularity for auto-partitioning. Valid options are |
|
|
Integer |
2 |
No |
The number of future partitions to pre-create during checks. |
|
|
Integer |
7 |
No |
The number of historical partitions to retain. Older partitions are automatically deleted. |
|
|
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 |
|
|
Enum |
full |
No |
The consumption startup mode. Valid options are |
|
|
Long/String |
None |
Conditionally required |
The startup timestamp, which can be a millisecond value or a string in the |
|
|
Duration |
1 minute |
No |
The interval for automatically discovering new partitions. Set to a negative value to disable automatic discovery. |
|
|
Boolean |
true |
No |
Specifies whether to perform a CRC32 check on messages to ensure data integrity. |
|
|
Integer |
500 |
No |
The maximum number of records that a single |
|
|
MemorySize |
16mb |
No |
The maximum number of bytes to fetch from the server per request. |
|
|
MemorySize |
1mb |
No |
The maximum number of bytes to fetch per bucket in a single request. |
|
|
MemorySize |
1b |
No |
The minimum data size the server should return for a fetch request. |
|
|
Duration |
500ms |
No |
The maximum time that the server waits when the minimum number of bytes is not met. |
|
|
String |
System temporary directory |
No |
The local directory for storing temporary files such as snapshots and log segments. |
|
|
Integer |
4 |
No |
The number of remote log segments to prefetch. |
|
|
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 |
|
|
Boolean |
false |
No |
Specifies whether to ignore |
|
|
Boolean |
true |
No |
Specifies whether to shuffle data by bucket ID before writing, which can improve write performance. |
|
|
MemorySize |
64mb |
No |
The total memory size of the internal buffer for row data. |
|
|
MemorySize |
2mb |
No |
The target batch size for records within the same bucket. |
|
|
Duration |
Long.MAX |
No |
The maximum time the writer blocks while waiting for an available segment. |
|
|
Duration |
100ms |
No |
The maximum time to buffer records before sending them in a batch. |
|
|
Enum |
STICKY |
No |
The bucket assignment strategy for tables without a primary key. Valid options are |
|
|
String |
all |
No |
The acknowledgment level for writes. |
|
|
MemorySize |
10mb |
No |
The maximum size of a single request in bytes. |
|
|
Integer |
MAX_VALUE |
No |
The number of times to retry a send that failed with a retriable error. |
|
|
Boolean |
true |
No |
Enables idempotent writes, which provide exactly-once semantics and message ordering guarantees. |
|
|
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 |
|
|
Boolean |
true |
No |
Enables asynchronous lookup, which can improve throughput for dimension table joins. |
|
|
Enum |
NONE |
No |
The cache policy. Valid options are |
|
|
Integer |
3 |
No |
The maximum number of retries if a lookup fails. |
|
|
Duration |
None |
No |
The cache expiration period after the last access. |
|
|
Duration |
None |
No |
The cache expiration period after the last write. |
|
|
Boolean |
true |
No |
Specifies whether to cache keys that do not exist in the dimension table (i.e., cache null values). |
|
|
Long |
None |
No |
The maximum number of rows in the cache. |
|
|
Integer |
25600 |
No |
The size of the queue for pending lookup operations. |
|
|
Integer |
128 |
No |
The maximum batch size for combining multiple lookup operations. |
|
|
Integer |
128 |
No |
The maximum number of in-flight lookup requests. |
|
|
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.
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: |
|
bucket.key |
Specifies the bucket key. |
No |
String |
None |
Specify the data distribution policy for each Fluss table. Tables are separated by Format: 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.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: If a table does not have a configured bucket number, the server-side configuration item The values of |
|
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) |