A lake-stream unified table provides transparent access to real-time incremental data (Fluss) and historical full data (Paimon) through a unified metadata layer. Depending on your business needs, you can choose from real-time streaming reads, offline batch reads, or hybrid read modes.
Supported query engines
Engine compatibility depends on the access method.
|
Access method |
Description |
Supported engines |
|
Native Fluss access (including Union Read) |
Uses the lake-stream fusion hybrid read feature to unify stream and batch processing. |
Flink |
|
Underlying Paimon access ( |
Directly accesses the underlying Paimon-formatted storage, supporting multi-engine analytics. |
All Paimon-compatible engines, such as Flink, Spark, Trino, and StarRocks |
Access data with Flink
Consume data in real time
In stream processing scenarios, you often need to monitor the latest data changes. You can specify a consumption mode by configuring the startup mode in Flink SQL.
-- Continuously monitor changes to the orders table and read only the latest incoming data
SELECT *
FROM fluss_catalog.fluss.orders
/*+ OPTIONS('scan.startup.mode' = 'latest-offset') */;
For other startup mode configurations, such as reading from the earliest offset or a specific timestamp, refer to Consumption mode.
Batch queries
Fluss automatically synchronizes data in real time to the underlying Paimon-format storage. For offline tasks that need to analyze the full historical dataset or access a specific snapshot, you can directly query the underlying Paimon table.
In the Fluss Catalog, you can directly access the corresponding Paimon table by appending the $lake suffix to the table name.
-- Batch mode: Directly query all data from the underlying Paimon storage
SELECT * FROM fluss_catalog.fluss.orders$lake;
-- Batch mode: Query snapshot metadata from the underlying Paimon storage
SELECT * FROM fluss_catalog.fluss.orders$lake$snapshots;
A $lake table is essentially a standard Paimon table. It inherits all Paimon features, including Time Travel and support for multiple query engines.
Union Read
Union Read provides a unified view for lake-stream fusion. The system automatically merges data from the Fluss layer (hot data with second-level latency) with data from the Paimon layer (cold data with minute-level latency). This process presents the data as a single, logically complete table, abstracting the underlying physical storage.
In traditional architectures that separate stream and batch processing, accessing both real-time and offline data requires you to build complex data pipelines for data migration, ETL aggregation, and deduplication. This approach is costly to develop and maintain and makes it difficult to ensure data consistency. Union Read natively fuses stream and lake data at the storage engine level. A single SQL query provides transparent access to all hot and cold data, eliminating the need to manage the underlying data sharding and merge logic.
-- Automatically merge and read hot and cold data
SELECT * FROM fluss_catalog.fluss.orders;
The Flink engine uses the following execution logic.
|
Execution mode |
Logic |
|
Streaming mode |
When a job starts, it first reads all historical data from Paimon, then seamlessly switches to consuming real-time incremental data from Fluss with millisecond-level latency. |
|
Batch mode |
Reads the base data from Paimon and the unarchived incremental data from Fluss, merges them, and returns the latest global snapshot. |
For a primary key table, Flink must merge and deduplicate data from both sources by primary key when running in batch mode. This process incurs computational overhead compared to reading a log table.
For more information about accessing data with Flink, refer to Fluss Connector.
Access data with StarRocks
After you create and configure the Fluss Catalog (for details, refer to Engine Integration - StarRocks), you can query data in the lake-stream unified table using the following methods.
Query the data lake
StarRocks uses the $lake suffix to directly access all historical data in the underlying Paimon storage.
-- Query all data from the underlying Paimon storage
SELECT * FROM <catalog_name>.<database_name>.<table_name>$lake;
Query the stream
StarRocks uses the $rt suffix to directly query unarchived, real-time incremental data in Fluss.
-- Query real-time incremental data from the Fluss layer
SELECT * FROM <catalog_name>.<database_name>.<table_name>$rt;
Query all data (Union Read)
When no suffix is used, StarRocks automatically merges the Fluss layer (hot data) and the Paimon layer (cold data) to return a complete, up-to-date dataset.
-- Automatically merge hot and cold data to return a complete and up-to-date dataset
SELECT * FROM <catalog_name>.<database_name>.<table_name>;