This topic covers the features, advantages, and usage of the Fluss fully-managed lake-stream integration service.
Overview
The Fluss lake-stream integration automatically synchronizes data from Stream Storage Fluss to a data lake built on Paimon. This approach provides a single data source with two views, enabling both millisecond-level real-time access and high-throughput historical analysis.
Data synchronization jobs for lake-stream integration now use a fully-managed mode. The system automatically handles the creation, execution, and maintenance of synchronization jobs, freeing you from monitoring their underlying status.
Comparison with the previous solution
|
Item |
Previous solution |
Fully-managed solution (New) |
|
Workspace |
You must manually select an existing workspace. |
Managed automatically by the system. No selection is required. |
|
Synchronization status visibility |
You can only view whether the synchronization job is running, but not the synchronization status of individual tables. |
Offers table-level metrics like synchronization latency and synchronization checkpoints, and supports alert notifications. |
|
Synchronization performance |
Based on standard Flink operators. |
It uses self-developed native operators that improve synchronization performance by over 30%. |
|
Maintenance |
Requires you to manually monitor and manage the running status of synchronization jobs. |
Fully managed by the system. No manual intervention is required. |
Core capabilities
Table-level synchronization monitoring
The fully-managed mode provides observability of the synchronization status at the table level.
-
Synchronization latency: Displays the delay for each table's data to be synchronized from the Fluss stream layer to the Paimon lake layer. This metric is used to evaluate data freshness.
-
Synchronization checkpoint: Displays the current synchronization progress for each table, reflecting the real-time status of data synchronization.
-
Alert notifications: Allows you to configure alert rules for key metrics such as synchronization latency. The system automatically sends notifications when these metrics are abnormal.
Self-developed native synchronization operator
Fully-managed synchronization jobs use a self-developed native operator instead of standard Flink operators. This approach deeply optimizes the data synchronization pipeline and improves performance by over 30% compared to the previous version.
Fully-managed O&M
-
Synchronization jobs are automatically created and managed by the system, requiring no manual configuration or maintenance.
-
The system automatically recovers from job exceptions to ensure the continuity of data synchronization.
-
The system automatically resumes from the last successful checkpoint to ensure exactly-once semantics.
Migration for existing users
For users of the previous solution, the system will perform an automated, universal migration:
-
After the migration, the system will automatically use fully-managed synchronization jobs for tables that have lake-stream integration enabled.
-
Your original, non-managed synchronization jobs will be automatically stopped and deleted.
-
The migration is seamless and does not affect your business or interrupt data synchronization. After the migration is complete, the non-managed synchronization method will no longer be supported.
-
The default limit for pay-as-you-go resources is 1,000 CUs. You can adjust this limit manually.
During the migration process, the system ensures a seamless transition of synchronization checkpoints to prevent data loss or duplication.
Billing
Fully-managed synchronization jobs consume computing resources during operation. The billing details are as follows:
-
Pay-as-you-go: Resources consumed by synchronization jobs are billed based on actual usage. This service is free of charge during the public preview period and will become a paid service upon official commercial release.
-
Provisioned resources: The system first deducts costs from your provisioned resources. You are billed on a pay-as-you-go basis only for usage that exceeds your provisioned amount.
Choose an appropriate amount of provisioned resources to reduce costs.
Prerequisites
Service Activation
-
You have activated the DLF service and created a new Catalog. This Catalog must be in the same region as Fluss, and both the VPC where Flink resides and the VPC where Fluss resides must be added to the trusted VPC list of DLF. For more information, see Authorize and activate DLF and Trusted VPC configuration.
Version restrictions
-
DLF-Legacy is not supported.
Permissions
-
Only the cluster owner can enable the fully-managed lake-stream integration service.
-
The account used to enable the service must have the following permissions:
The AliyunDLFFullAccess permission in DLF, which grants permissions to create databases and tables, and perform read and write operations on them. For more information, see User Authorization.
Enable lake-stream integration
Step 1: Enable lake-stream integration for the cluster
-
Log on to the Realtime Compute for Apache Flink console.
-
Navigate to the Streaming Storage Fluss tab and click Console in the Actions column to open the Fluss console.
-
On the Cluster Overview page, click the
toggle in the lower-right corner to enable the service and complete the service authorization.
Step 2: Enable lake-stream integration for a table
New table
In your Flink workspace, register a Fluss Catalog, and then create a table in .
CREATE TABLE `my-catalog`.`fluss`.`datalake_orders` (
shop_id BIGINT,
user_id BIGINT,
num_orders INT,
total_amount INT,
PRIMARY KEY (shop_id, user_id) NOT ENFORCED
) WITH (
'bucket.num' = '4',
'table.datalake.enabled' = 'true',
'table.datalake.freshness' = '30min'
);
Additional parameters for lake-stream integration
|
Parameter |
Description |
Default |
Remarks |
|
table.datalake.enabled |
Specifies whether to enable lake-stream integration. |
false |
If this parameter is not specified, you can enable lake-stream integration manually in the console. |
|
table.datalake.freshness |
Defines the maximum allowed latency (data freshness) for the Paimon table data relative to the source Fluss table. |
3min |
The Fluss service automatically performs data synchronization based on this configuration. A smaller value improves data freshness in the data lake. Conversely, if your business is not sensitive to latency, a larger value can reduce resource consumption. Important
This parameter cannot be modified after the table is created. |
|
paimon.* |
Any parameter prefixed with |
None |
Fluss lake-stream integration creates the underlying Paimon lake table using the default Paimon parameters. If you need to specify other parameters for the Paimon lake table, such as setting the file format to ORC, you can set For more parameter configurations, see Paimon Configuration. |
Creating Paimon tables with deletion vectors enabled is not currently supported. Do not pass the 'paimon.deletion-vectors.enabled' = 'true' parameter.
Existing table
-
Enable in the console
In the Fluss console, select Data from the left-side navigation pane. Click the target database. In the list of data tables on the right, click the
toggle to enable lake-stream integration. -
Using SQL: Use the
ALTER TABLEstatement.ALTER TABLE datalake_orders SET ('table.datalake.enabled' = 'true');
Query synchronized data
After you enable lake-stream integration for a data table, the system automatically synchronizes the table's metadata to DLF. In DLF, under the associated Catalog, you can find a database with the same name as the one in Fluss. This database contains a table with a structure and metadata identical to the source table, allowing for unified querying and management.
Before synchronization, ensure that no table with the same name exists in DLF, as this may cause the synchronization to fail or result in conflicts.
Data consistency in lake-stream integration
-
After you enable a lake-stream integration table, data from Fluss is continuously written to the Paimon-based data lake via a Flink synchronization job. The table's structural metadata is centrally managed by DLF and is retained indefinitely.
-
Deleting a lake-stream integration table in Fluss does not delete the synchronized lake table (Paimon table) in DLF. You must manually delete the table in DLF.
-
For a primary key table, data changes are handled using a changelog mode:
-
When a record is deleted, the system physically removes the historical data from the Paimon table after the next tiering process is complete. Simultaneously, a changelog entry with a type of DELETE is appended.
-
When you perform a Union Read query, the Flink engine automatically merges the real-time data in Fluss with the historical data in Paimon, and performs deduplication and state merging based on the primary keys and changelogs.
-
Logically deleted records will not appear in the final query result, ensuring the accuracy and consistency of query results.
-
Disable lake-stream integration
Disable lake-stream integration for a data table
-
From the console: In the Fluss console, select Data from the left-side navigation pane and disable lake-stream integration for each enabled table.
-
Using SQL: Use the
ALTER TABLEstatement.ALTER TABLE datalake_orders SET ('table.datalake.enabled' = 'false');
Disable lake-stream integration for a cluster
On the Cluster Overview page, click the
toggle in the lower-right corner to disable the lake-stream integration service.
-
You cannot disable the lake-stream integration service at the cluster level if the service is still enabled for any tables.
-
When you disable the lake-stream integration service, the fully-managed synchronization job also stops running.