Fluss是一款具备亚秒级低延迟特性的流式存储系统,支持高效的实时数据读写。借助Lakehouse Storage,Fluss在Lakehouse架构之上构建了实时流数据服务,从而实现了数据流与数据湖仓的无缝集成与统一管理。

一、湖仓概述

1. 传统湖仓架构的挑战

数据湖仓(Lakehouse)作为一种创新的开放架构,融合了数据湖的灵活扩展性和经济高效性,以及数据仓库的稳定可靠性和高性能优势。在这一架构中,Apache Iceberg、Apache Paimon、Apache Hudi 和 Delta Lake 等领先的数据湖格式发挥了重要作用,它们通过统一的平台,实现了数据存储、可靠性与分析能力之间的高效协同与平衡。

然而,传统湖仓架构面临着实时性与分析能力平衡的挑战:

(1) 实时性与分析效率的矛盾

  • 如果要求低延迟,就需要频繁写入和提交,这会产生大量小型Parquet文件,导致读取效率低下
  • 如果要求读取效率,就需要积累数据直到能写入大型Parquet文件,但这会引入更高的延迟 

(2) 数据新鲜度限制

即使在最佳使用条件下,这些数据湖格式通常也只能在分钟级粒度内实现数据新鲜度。

二、流式湖仓一体化

1. 流与湖的统一

Fluss是一种支持亚秒级低延迟流式读写的流式存储系统。通过Lakehouse Storage,Fluss在Lakehouse之上提供实时流数据服务,实现了数据流和数据湖仓的统一。这不仅为数据湖仓带来了低延迟,还为数据流增加了强大的分析能力。 

为了构建流式湖仓,Fluss维护了一个分层服务,该服务将实时数据从Fluss集群压缩到存储在Lakehouse Storage中的数据湖格式。Fluss集群中的数据(流式Arrow格式)针对低延迟读写进行了优化,而Lakehouse中的压缩数据(带压缩的Parquet格式)针对强大的分析和长期数据存储进行了优化。因此,Fluss集群中的数据作为实时数据层,保留了具有亚秒级新鲜度的数据;而Lakehouse中的数据作为历史数据层,保留了具有分钟级新鲜度的数据。

2. 共享数据与共享元数据

流式湖仓的核心理念是流和湖仓之间共享数据和共享元数据,避免数据重复和元数据不一致。它提供了以下强大功能:

  • 统一元数据:Fluss为流和湖仓中的数据提供统一的表元数据。用户只需处理一个表,但可以访问实时流数据、历史数据或它们的联合。
  • 联合读取:计算引擎对表执行查询时将读取实时流数据和湖仓数据的联合。目前,只有Flink支持联合读取,但更多引擎正在规划中。
  • 实时湖仓:联合读取帮助湖仓从近实时分析发展到真正的实时分析,使企业能够从实时数据中获得更有价值的洞察。
  • 分析流:联合读取帮助数据流具备强大的分析能力,这减少了开发流应用程序的复杂性,简化了调试,并允许立即访问实时数据洞察。
  • 连接到湖仓生态系统:Fluss在将数据压缩到Lakehouse时保持表元数据与数据湖目录同步,这允许Spark、StarRocks、Flink、Trino等外部引擎通过连接到数据湖目录直接读取数据。 
3. 数据湖集成

(1) Paimon集成

Apache Paimon创新地结合了湖格式和LSM结构,将高效更新引入湖架构。 

要将Fluss与Paimon集成,必须启用lakehouse存储并将Paimon配置为lakehouse存储。具体步骤如下:

① 配置Lakehouse存储

首先,在server.yaml中配置lakehouse存储:

复制

datalake.format: paimon  


# Paimon目录配置,假设使用Filesystem目录  
datalake.paimon.metastore: filesystem  
datalake.paimon.warehouse: /tmp/paimon_data_warehouse

    ② 启动数据湖分层服务

    然后,启动数据湖分层服务,将Fluss的数据压缩到lakehouse存储:

    复制

    # 切换到Fluss目录  
    cd $FLUSS_HOME  
    
    
    # 启动分层服务,假设rest端点是localhost:8081  
    ./bin/lakehouse.sh -D flink.rest.address=localhost -D flink.rest.port=8081

    ③ 为表启用Lakehouse存储

    要为表启用lakehouse存储,必须在创建表时使用选项'table.datalake.enabled' = 'true':

    复制

    CREATE TABLE datalake_enriched_orders (  
        `order_key` BIGINT,  
        `cust_key` INT NOT NULL,  
        `total_price` DECIMAL(15, 2),  
        `order_date` DATE,  
        `order_priority` STRING,  
        `clerk` STRING,  
        `cust_name` STRING,  
        `cust_phone` STRING,  
        `cust_acctbal` DECIMAL(15, 2),  
        `cust_mktsegment` STRING,  
        `nation_name` STRING,  
        PRIMARY KEY (`order_key`) NOT ENFORCED  
    ) WITH ('table.datalake.enabled' = 'true');

      当在Fluss中创建或修改带有选项'table.datalake.enabled' = 'true'的表时,Fluss将创建一个具有相同表路径的相应Paimon表。Paimon表的模式与Fluss表的模式相同,只是在最后附加了两个额外的列__offset和__timestamp。这两列用于帮助Fluss客户端以流式方式消费Paimon中的数据,例如按偏移量/时间戳查找等。 

      然后,数据湖分层服务会持续将数据从Fluss压缩到Paimon。对于主键表,它还将生成Paimon格式的变更日志,使您能够以Paimon方式流式消费它。 

      (2) Iceberg集成规划

      目前,Fluss支持Paimon作为Lakehouse存储,更多种类的数据湖格式正在规划中。 

      Fluss社区正在积极努力增强流和湖仓统一功能,重点关注以下关键领域:

      ① 扩展联合读取生态系统

      目前,联合读取功能已与Apache Flink集成,实现了实时和历史数据的无缝查询。未来,社区计划扩展此功能以支持其他查询引擎,如Apache Spark和StarRocks,进一步扩大其生态系统兼容性和采用。

      ② 多样化湖存储格式

      目前,Fluss支持Apache Paimon作为其主要湖存储。为了满足多样化的用户需求,社区计划添加对更多湖格式的支持,包括Apache Iceberg和Apache Hudi,从而提供更大的灵活性和与更广泛的湖仓生态系统的互操作性。

      4. 联合读取(Union Read)

      (1) 实时数据与历史数据的无缝结合

      对于带有选项'table.datalake.enabled' = 'true'的表,有两部分数据:保留在Fluss中的数据和已经在Paimon中的数据。现在,您有两种表视图:一种是具有分钟级延迟的Paimon数据视图,一种是联合Fluss和Paimon数据的完整数据视图,具有秒级延迟。 

      Flink使您能够决定选择哪种视图:

      • 仅Paimon意味着更好的分析性能,但数据新鲜度较差
      • 结合Fluss和Paimon意味着更好的数据新鲜度,但分析性能下降

      ① 仅读取Paimon中的数据

      要指定读取Paimon中的数据,必须使用$lake后缀指定表,以下SQL显示了如何执行此操作:

      复制

      -- 假设我们有一个名为`orders`的表  
      
      
      -- 从paimon读取  
      SELECT COUNT(*) FROM orders$lake;  
      
      
      -- 我们还可以查询系统表  
      SELECT * FROM orders$lake$snapshots;

        当在查询中使用$lake后缀指定表时,它就像一个普通的Paimon表,因此它继承了Paimon表的所有能力。您可以享受Flink在Paimon上支持/优化的所有功能,如查询系统表、时间旅行等。 

        ② 联合读取Fluss和Paimon中的数据

        要指定读取联合Fluss和Paimon的完整数据,只需像查询普通表一样查询它,无需任何后缀或其他内容,以下SQL显示了如何执行此操作:

        复制

        -- 查询将联合Fluss和Paimon的数据  
        SELECT SUM(order_count) as total_orders FROM ads_nation_purchase_power;

          查询可能看起来比只查询Paimon中的数据慢,但它查询了完整数据,这意味着更好的数据新鲜度。您可以多次运行查询,由于数据持续写入表中,每次运行都应该得到不同的结果。 

          (2) 查询引擎支持

          ① Flink支持

          Fluss向Flink用户公开统一的API,允许他们选择是使用联合读取还是仅在Lakehouse上进行读取,使用以下SQL:

          SELECT * FROM orders

            这将读取orders表的完整数据,Flink将联合读取Fluss和Lakehouse中的数据。如果用户只需要读取数据湖上的数据,可以在要读取的表后添加$lake后缀。SQL如下:

            -- 分析查询  
            SELECT COUNT(*), MAX(t), SUM(amount)   
            FROM orders$lake  
            
            
            -- 查询系统表  
            SELECT * FROM orders$lake$snapshots
              5. 实时湖仓与分析流

              (1) 为湖仓带来真正的实时性

              传统的数据湖仓架构通常只能提供分钟级的数据新鲜度,这对于许多实时分析场景来说是不够的。Fluss的流式湖仓一体化架构通过联合读取(Union Read)机制,为湖仓带来了真正的实时性:

              • 秒级数据新鲜度:通过将Fluss的实时流数据与Paimon的历史数据联合起来,查询可以访问最新写入的数据,实现秒级数据新鲜度。
              • 实时分析能力增强:企业可以在保持湖仓强大分析能力的同时,获得实时数据洞察,使决策更加及时和准确。
              • 无需数据复制:不需要将数据复制到单独的实时系统中,减少了存储成本和数据一致性问题。

              (2) 为数据流增加强大的分析能力

              传统的流处理系统通常缺乏强大的分析能力,而Fluss通过与湖仓的集成,为数据流增加了这些能力:

              • 复杂分析查询支持:流数据可以与历史数据结合,支持更复杂的分析查询,如聚合、窗口分析等。
              • 开发简化:减少了开发流应用程序的复杂性,简化了调试过程。 
              • 生态系统集成:可以利用现有的湖仓生态系统工具和查询引擎,如Spark、StarRocks、Trino等。
              6. 流式湖仓一体化实施流程

              下面是实施Fluss流式湖仓一体化的详细步骤和流程图:

              (1) 准备环境

              确保已安装Fluss和Flink,并且两者都正常运行。Fluss需要一个运行中的Flink集群来执行数据湖分层服务。

              (2) 配置Lakehouse存储

              在Fluss的server.yaml配置文件中设置Lakehouse存储:

              # 指定数据湖格式为Paimon  
              datalake.format: paimon  
              
              
              # Paimon目录配置  
              datalake.paimon.metastore: filesystem  
              datalake.paimon.warehouse: /path/to/paimon/warehouse

                (3) 启动数据湖分层服务

                数据湖分层服务是一个Flink作业,负责将数据从Fluss压缩到Paimon:

                复制

                # 切换到Fluss安装目录  
                cd $FLUSS_HOME  
                
                
                # 启动分层服务,指定Flink REST端点  
                ./bin/lakehouse.sh -D flink.rest.address=localhost -D flink.rest.port=8081

                  您还可以设置其他Flink配置参数,例如:

                  # 设置检查点间隔为10秒  
                  ./bin/lakehouse.sh -D flink.rest.address=localhost -D flink.rest.port=8081 -D flink.execution.checkpointing.interval=10s  
                  
                  
                  # 仅同步特定数据库中的表  
                  ./bin/lakehouse.sh -D flink.rest.address=localhost -D flink.rest.port=8081 -D database=fluss_\\w+

                  (4) 创建启用Lakehouse的表

                  使用Flink SQL创建启用了Lakehouse的表,通过设置'table.datalake.enabled' = 'true'选项:

                  复制

                  -- 创建一个启用了Lakehouse的日志表  
                  CREATE TABLE log_orders (  
                      order_id BIGINT,  
                      customer_id BIGINT,  
                      order_time TIMESTAMP,  
                      amount DECIMAL(10, 2),  
                      status STRING  
                  ) WITH (  
                      'bucket.num' = '4',  
                      'table.datalake.enabled' = 'true'  
                  );  
                  
                  
                  -- 创建一个启用了Lakehouse的主键表  
                  CREATE TABLE pk_orders (  
                      order_id BIGINT,  
                      customer_id BIGINT,  
                      order_time TIMESTAMP,  
                      amount DECIMAL(10, 2),  
                      status STRING,  
                      PRIMARY KEY (order_id) NOT ENFORCED  
                  ) WITH (  
                      'bucket.num' = '4',  
                      'table.datalake.enabled' = 'true'  
                  );

                    当创建启用了Lakehouse的表时,Fluss会自动创建一个对应的Paimon表,并在表末尾添加两个额外的列__offset和__timestamp,用于支持流式消费。

                    (5) 数据写入

                    使用标准的Flink SQL向表中写入数据:

                    复制

                    -- 插入数据到日志表  
                    INSERT INTO log_orders  
                    VALUES (1001, 101, TIMESTAMP '2025-05-18 10:30:00', 199.99, 'PENDING');  
                    
                    
                    -- 插入数据到主键表  
                    INSERT INTO pk_orders  
                    VALUES (1001, 101, TIMESTAMP '2025-05-18 10:30:00', 199.99, 'PENDING');  
                    
                    
                    -- 或者从其他表插入数据  
                    INSERT INTO pk_orders  
                    SELECT order_id, customer_id, order_time, amount, status  
                    FROM source_table;

                    数据写入后,会首先存储在Fluss的实时层中,然后由数据湖分层服务定期压缩到Paimon中。对于主键表,还会生成Paimon格式的变更日志。

                    (6) 数据查询

                    ① 联合读取(实时+历史数据)

                    要查询完整的数据(Fluss中的实时数据和Paimon中的历史数据),直接查询表名,无需任何后缀:

                    复制

                    -- 联合读取实时和历史数据  
                    SELECT COUNT(*) FROM pk_orders;  
                    
                    
                    -- 流式查询,会持续获取最新数据  
                    SELECT * FROM pk_orders WHERE amount > 100;

                    联合读取提供秒级数据新鲜度,但可能会比仅查询Paimon数据慢一些。由于数据持续写入,多次运行相同的查询可能会得到不同的结果。

                    ② 仅读取Lakehouse数据

                    要仅查询Paimon中的历史数据,在表名后添加$lake后缀:

                    -- 仅查询Paimon中的历史数据  
                    SELECT COUNT(*) FROM pk_orders$lake;  
                    
                    
                    -- 查询Paimon系统表  
                    SELECT * FROM pk_orders$lake$snapshots;  
                    
                    
                    -- 使用Paimon的时间旅行功能  
                    SELECT * FROM pk_orders$lake FOR TIMESTAMP AS OF '2025-05-17 00:00:00';

                      仅读取Paimon数据提供更好的分析性能,但数据新鲜度较差(分钟级)。您可以使用Paimon的所有功能,如查询系统表、时间旅行等。

                      7. 流式湖仓一体化架构图

                      三、实际应用场景示例

                      场景1:实时仪表盘

                      一个电子商务平台需要一个实时仪表盘,显示最新的销售数据和趋势。

                      实现方式:

                      • 创建启用Lakehouse的订单表
                      • 使用Flink SQL进行联合读取,获取最新的销售数据
                      • 将结果推送到仪表盘应用

                      -- 创建启用Lakehouse的订单表  
                      CREATE TABLE sales_orders (  
                          order_id BIGINT,  
                          product_id BIGINT,  
                          customer_id BIGINT,  
                          order_time TIMESTAMP,  
                          amount DECIMAL(10, 2),  
                          PRIMARY KEY (order_id) NOT ENFORCED  
                      ) WITH (  
                          'bucket.num' = '8',  
                          'table.datalake.enabled' = 'true'  
                      );  
                      
                      
                      -- 实时计算最近1小时的销售额  
                      SELECT   
                          TUMBLE_START(order_time, INTERVAL '5' MINUTE) AS window_start,  
                          SUM(amount) AS total_sales  
                      FROM sales_orders  
                      WHERE order_time > CURRENT_TIMESTAMP - INTERVAL '1' HOUR  
                      GROUP BY TUMBLE(order_time, INTERVAL '5' MINUTE);
                      场景2:历史数据分析与实时数据结合

                      一个金融机构需要分析客户的历史交易模式,并将其与实时交易数据结合,以检测潜在的欺诈行为。

                      实现方式:

                      • 创建启用Lakehouse的交易表
                      • 使用Paimon查询历史交易模式
                      • 使用联合读取将历史模式与实时交易结合分析

                      -- 创建启用Lakehouse的交易表  
                      CREATE TABLE transactions (  
                          transaction_id BIGINT,  
                          account_id BIGINT,  
                          transaction_time TIMESTAMP,  
                          amount DECIMAL(15, 2),  
                          merchant_id BIGINT,  
                          location STRING,  
                          PRIMARY KEY (transaction_id) NOT ENFORCED  
                      ) WITH (  
                          'bucket.num' = '16',  
                          'table.datalake.enabled' = 'true'  
                      );  
                      
                      
                      -- 查询历史交易模式(仅Paimon数据)  
                      SELECT   
                          account_id,  
                          AVG(amount) AS avg_amount,  
                          STDDEV(amount) AS stddev_amount  
                      FROM transactions$lake  
                      WHERE transaction_time BETWEEN TIMESTAMP '2025-04-01 00:00:00' AND TIMESTAMP '2025-05-01 00:00:00'  
                      GROUP BY account_id;
                        -- 将历史模式与实时交易结合分析(联合读取)  
                        WITH historical_patterns AS (  
                            SELECT   
                                account_id,  
                                AVG(amount) AS avg_amount,  
                                STDDEV(amount) AS stddev_amount  
                            FROM transactions$lake  
                            WHERE transaction_time BETWEEN TIMESTAMP '2025-04-01 00:00:00' AND TIMESTAMP '2025-05-01 00:00:00'  
                            GROUP BY account_id  
                        )  
                        SELECT   
                            t.transaction_id,  
                            t.account_id,  
                            t.transaction_time,  
                            t.amount,  
                            t.merchant_id,  
                            t.location,  
                            h.avg_amount,  
                            h.stddev_amount,  
                            CASE   
                                WHEN ABS(t.amount - h.avg_amount) > 3 * h.stddev_amount THEN 'SUSPICIOUS'  
                                ELSE 'NORMAL'  
                            END AS risk_level  
                        FROM transactions t  
                        JOIN historical_patterns h ON t.account_id = h.account_id  
                        WHERE t.transaction_time > CURRENT_TIMESTAMP - INTERVAL '1' HOUR;
                        场景3:实时数据湖分层架构

                        一个大型零售企业需要构建一个数据湖分层架构(Bronze、Silver、Gold),同时保持各层的数据新鲜度。

                        实现方式:

                        • 创建各层的启用Lakehouse的表
                        • 使用Flink SQL将数据从一层转换到下一层
                        • 利用联合读取确保各层都具有秒级数据新鲜度

                        复制

                        -- Bronze层:创建原始销售数据表  
                        CREATE TABLE sales_raw (  
                            event_time TIMESTAMP,  
                            store_id BIGINT,  
                            product_id BIGINT,  
                            quantity INT,  
                            amount DECIMAL(10, 2),  
                            payment_method STRING,  
                            source STRING  
                        ) WITH (  
                            'bucket.num' = '16',  
                            'table.datalake.enabled' = 'true'  
                        );  
                        
                        
                        -- Silver层:创建清洗后的销售数据表  
                        CREATE TABLE sales_cleaned (  
                            event_time TIMESTAMP,  
                            store_id BIGINT,  
                            product_id BIGINT,  
                            quantity INT,  
                            amount DECIMAL(10, 2),  
                            payment_method STRING,  
                            source STRING,  
                            PRIMARY KEY (store_id, product_id, event_time) NOT ENFORCED  
                        ) WITH (  
                            'bucket.num' = '16',  
                            'table.datalake.enabled' = 'true'  
                        );  
                        
                        
                        -- Gold层:创建聚合的销售指标表  
                        CREATE TABLE sales_aggregated (  
                            window_start TIMESTAMP,  
                            window_end TIMESTAMP,  
                            store_id BIGINT,  
                            product_id BIGINT,  
                            total_quantity BIGINT,  
                            total_amount DECIMAL(15, 2),  
                            PRIMARY KEY (window_start, window_end, store_id, product_id) NOT ENFORCED  
                        ) WITH (  
                            'bucket.num' = '8',  
                            'table.datalake.enabled' = 'true'  
                        );  
                        
                        
                        -- Silver层ETL:清洗和转换数据  
                        INSERT INTO sales_cleaned  
                        SELECT   
                            event_time,  
                            store_id,  
                            product_id,  
                            CASE WHEN quantity <= 0 THEN 1 ELSE quantity END AS quantity,  
                            amount,  
                            payment_method,  
                            source  
                        FROM sales_raw  
                        WHERE amount > 0;  
                        
                        
                        -- Gold层ETL:聚合数据  
                        INSERT INTO sales_aggregated  
                        SELECT   
                            TUMBLE_START(event_time, INTERVAL '1' HOUR) AS window_start,  
                            TUMBLE_END(event_time, INTERVAL '1' HOUR) AS window_end,  
                            store_id,  
                            product_id,  
                            SUM(quantity) AS total_quantity,  
                            SUM(amount) AS total_amount  
                        FROM sales_cleaned  
                        GROUP BY TUMBLE(event_time, INTERVAL '1' HOUR), store_id, product_id;

                          四、流式湖仓一体化的技术优势

                          1. 数据分布一致性

                          Fluss和Paimon之间的数据分布是严格对齐的。Fluss支持分区表和桶,其分桶算法与Paimon完全一致。这确保了给定的数据始终分配到两个系统中的相同桶,创建了Fluss桶和Paimon桶之间的一对一对应关系。

                          这种数据分布的一致性提供了两个显著优势:

                          (1) 消除分层过程中的Shuffle开销

                          当将数据从Fluss分层到Paimon格式时:

                          • Fluss桶(例如bucket1)可以直接分层到相应的Paimon桶(bucket1)
                          • 无需读取Fluss桶中的数据,计算每条数据属于哪个Paimon桶,然后将其写入适当的Paimon桶

                          通过绕过这个中间重新分配步骤,架构避免了昂贵的shuffle开销,显著提高了压缩效率。

                          (2) 防止数据不一致

                          通过在Fluss和Paimon中使用相同的分桶算法,确保了数据一致性。该算法计算每条数据的桶分配如下:

                          bucket_id = hash(row) % bucket_num
                            2. 更高效的数据追踪

                            在Fluss架构中,历史数据存储在Lakehouse存储中,而实时数据保留在Fluss中。在流式读取期间,这种架构实现了历史和实时数据访问的无缝结合:

                            • 历史数据访问:Fluss直接从Lakehouse存储检索历史数据,利用其固有优势,如:
                            • 高效的过滤下推:使查询引擎能够在存储层应用过滤条件,减少读取的数据量并提高性能
                            • 列裁剪:允许仅检索必要的列,优化数据传输和查询效率
                            • 高压缩比:在保持快速检索速度的同时最小化存储开销
                            • 实时数据访问:Fluss同时从自己的存储中读取最新的实时数据,确保毫秒级的新鲜度
                            3. 一致的数据新鲜度

                            在构建数据仓库时,通常按层组织和管理数据,如Bronze、Silver和Gold层。当数据流经这些层时,维护数据新鲜度变得至关重要。

                            当Paimon作为每一层的唯一存储解决方案时,数据可见性取决于Flink检查点间隔,这会引入累积延迟:

                            • 给定层的变更日志仅在Flink检查点完成后才可见
                            • 当这个变更日志传播到后续层时,数据新鲜度延迟随着每个检查点间隔而增加

                            例如,使用1分钟的Flink检查点间隔:

                            • Bronze层经历1分钟延迟
                            • Silver层增加另一个1分钟延迟,总计2分钟
                            • Gold层再增加1分钟延迟,累积为3分钟延迟

                            而使用Fluss和Paimon,我们获得:

                            • 即时数据可见性:Fluss中的数据在摄入后立即可见,无需等待Flink检查点完成。变更日志立即传输到下一层
                            • 一致的数据新鲜度:所有层的数据新鲜度是一致的,以秒为单位,消除了累积延迟
                            4. 完整实施示例:电子商务实时分析平台

                            下面是一个完整的电子商务实时分析平台实施示例,展示了如何使用Fluss的流式湖仓一体化架构构建端到端解决方案。

                            步骤1:环境准备

                            首先,确保已安装并配置好Fluss和Flink环境:

                            # 下载并解压Fluss  
                            wget https://github.com/alibaba/fluss/releases/download/v0.6.0/fluss-0.6.0.tgz  
                            tar -xzf fluss-0.6.0.tgz  
                            cd fluss-0.6.0  
                            
                            
                            # 启动本地Fluss集群  
                            ./bin/start-local.sh  
                            
                            
                            # 启动本地Flink集群(用于数据湖分层服务)  
                            cd /path/to/flink  
                            ./bin/start-cluster.sh

                              步骤2:配置Lakehouse存储

                              在Fluss的conf/server.yaml中配置Paimon作为Lakehouse存储:

                              datalake.format: paimon  
                              datalake.paimon.metastore: filesystem  
                              datalake.paimon.warehouse: /path/to/paimon/warehouse

                                步骤3:启动数据湖分层服务

                                启动数据湖分层服务,将数据从Fluss压缩到Paimon:

                                cd /path/to/fluss  
                                ./bin/lakehouse.sh -D flink.rest.address=localhost -D flink.rest.port=8081 -D flink.execution.checkpointing.interval=10s

                                步骤4:创建数据模型

                                使用Flink SQL创建电子商务分析所需的表结构:

                                -- 创建订单表(主键表)  
                                CREATE TABLE orders (  
                                    order_id BIGINT,  
                                    customer_id BIGINT,  
                                    order_time TIMESTAMP,  
                                    total_amount DECIMAL(10, 2),  
                                    status STRING,  
                                    PRIMARY KEY (order_id) NOT ENFORCED  
                                ) WITH (  
                                    'bucket.num' = '8',  
                                    'table.datalake.enabled' = 'true'  
                                );  
                                
                                
                                -- 创建订单明细表(日志表)  
                                CREATE TABLE order_items (  
                                    item_id BIGINT,  
                                    order_id BIGINT,  
                                    product_id BIGINT,  
                                    quantity INT,  
                                    unit_price DECIMAL(10, 2),  
                                    discount DECIMAL(5, 2)  
                                ) WITH (  
                                    'bucket.num' = '8',  
                                    'table.datalake.enabled' = 'true'  
                                );  
                                
                                
                                -- 创建产品表(主键表)  
                                CREATE TABLE products (  
                                    product_id BIGINT,  
                                    name STRING,  
                                    category STRING,  
                                    brand STRING,  
                                    price DECIMAL(10, 2),  
                                    PRIMARY KEY (product_id) NOT ENFORCED  
                                ) WITH (  
                                    'bucket.num' = '4',  
                                    'table.datalake.enabled' = 'true'  
                                );  
                                
                                
                                -- 创建客户表(主键表)  
                                CREATE TABLE customers (  
                                    customer_id BIGINT,  
                                    name STRING,  
                                    email STRING,  
                                    registration_time TIMESTAMP,  
                                    vip_level INT,  
                                    PRIMARY KEY (customer_id) NOT ENFORCED  
                                ) WITH (  
                                    'bucket.num' = '4',  
                                    'table.datalake.enabled' = 'true'  
                                );  
                                
                                
                                -- 创建实时销售指标表(主键表)  
                                CREATE TABLE sales_metrics (  
                                    window_start TIMESTAMP,  
                                    window_end TIMESTAMP,  
                                    category STRING,  
                                    total_orders BIGINT,  
                                    total_sales DECIMAL(15, 2),  
                                    PRIMARY KEY (window_start, window_end, category) NOT ENFORCED  
                                ) WITH (  
                                    'bucket.num' = '4',  
                                    'table.datalake.enabled' = 'true'  
                                );

                                  步骤5:数据处理流程

                                  使用Flink SQL实现数据处理流程:

                                  -- 1. 计算实时销售指标(每5分钟更新一次)  
                                  INSERT INTO sales_metrics  
                                  SELECT   
                                      TUMBLE_START(o.order_time, INTERVAL '5' MINUTE) AS window_start,  
                                      TUMBLE_END(o.order_time, INTERVAL '5' MINUTE) AS window_end,  
                                      p.category,  
                                      COUNT(DISTINCT o.order_id) AS total_orders,  
                                      SUM(i.quantity * i.unit_price * (1 - i.discount/100)) AS total_sales  
                                  FROM orders o  
                                  JOIN order_items i ON o.order_id = i.order_id  
                                  JOIN products p ON i.product_id = p.product_id  
                                  WHERE o.status = 'COMPLETED'  
                                  GROUP BY   
                                      TUMBLE(o.order_time, INTERVAL '5' MINUTE),  
                                      p.category;  
                                  
                                  
                                  -- 2. 创建VIP客户实时推荐视图  
                                  CREATE VIEW vip_recommendations AS  
                                  SELECT   
                                      c.customer_id,  
                                      c.name,  
                                      p.category,  
                                      p.brand,  
                                      COUNT(*) AS purchase_count  
                                  FROM customers c  
                                  JOIN orders o ON c.customer_id = o.customer_id  
                                  JOIN order_items i ON o.order_id = i.order_id  
                                  JOIN products p ON i.product_id = p.product_id  
                                  WHERE   
                                      c.vip_level >= 3 AND  
                                      o.order_time > CURRENT_TIMESTAMP - INTERVAL '30' DAY  
                                  GROUP BY   
                                      c.customer_id, c.name, p.category, p.brand  
                                  HAVING COUNT(*) >= 2;

                                  步骤6:实时仪表盘查询

                                  使用联合读取查询实时销售指标:

                                  -- 实时销售趋势(联合读取)  
                                  SELECT   
                                      window_start,  
                                      category,  
                                      total_sales  
                                  FROM sales_metrics  
                                  WHERE   
                                      window_start >= CURRENT_TIMESTAMP - INTERVAL '24' HOUR  
                                  ORDER BY   
                                      window_start DESC, total_sales DESC;

                                  步骤7:历史分析查询

                                  使用Paimon查询历史数据分析:

                                  -- 历史销售分析(仅Paimon数据)  
                                  SELECT   
                                      p.category,  
                                      MONTH(o.order_time) AS month,  
                                      SUM(i.quantity * i.unit_price * (1 - i.discount/100)) AS total_sales  
                                  FROM orders$lake o  
                                  JOIN order_items$lake i ON o.order_id = i.order_id  
                                  JOIN products$lake p ON i.product_id = p.product_id  
                                  WHERE   
                                      o.order_time >= TIMESTAMP '2025-01-01 00:00:00' AND  
                                      o.order_time < TIMESTAMP '2025-05-01 00:00:00'  
                                  GROUP BY   
                                      p.category, MONTH(o.order_time)  
                                  ORDER BY   
                                      p.category, month;

                                  这种查询只访问Paimon中的历史数据,提供更好的查询性能,适合大规模历史数据分析。

                                  5. 流式湖仓一体化的关键技术点

                                  (1) 统一元数据管理

                                  在传统架构中,流存储系统(如Kafka)和湖仓存储解决方案(如Paimon)作为独立实体运行,各自维护自己的元数据。这给计算引擎(如Flink)带来了两个主要挑战:

                                  双重目录:用户需要创建和管理两个独立的目录——一个用于流存储,另一个用于湖存储

                                  手动切换:访问数据需要手动在目录之间切换,以确定是查询流存储还是湖存储,导致操作复杂性和效率低下

                                  而在Fluss中,虽然Fluss和Paimon仍然维护独立的元数据,但它们向计算引擎(如Flink)公开统一的目录和单一的表抽象。这种统一方法提供了几个关键优势:

                                  • 简化数据访问:用户可以通过单一目录无缝访问湖仓存储(Paimon)和流存储(Fluss),无需管理或在独立目录之间切换
                                  • 集成查询:统一表抽象允许直接访问Fluss中的实时数据和Paimon中的历史数据
                                  • 操作效率:通过提供一致的接口,该架构降低了操作复杂性,使用户更容易在单一工作流中处理实时和历史数据

                                  (2) 数据分布对齐

                                  Fluss和Paimon之间的数据分布是严格对齐的,这是通过使用相同的分桶算法实现的:

                                  bucket_id = hash(row) % bucket_num

                                    这种对齐确保了:

                                    • 分层效率:Fluss桶可以直接分层到对应的Paimon桶,无需重新分配数据
                                    • 数据一致性:相同的数据在两个系统中分配到相同的桶,防止数据不一致

                                    (3) 检查点间隔与数据新鲜度

                                    在传统的Paimon架构中,数据新鲜度受Flink检查点间隔的限制。例如,使用1分钟的检查点间隔时:

                                    而在Fluss流式湖仓架构中,数据在写入后立即可见:

                                    6. 流式湖仓一体化的最佳实践

                                    (1) 表设计最佳实践

                                    • 合理设置桶数量:根据数据量和查询模式设置适当的桶数量,通常为2的幂次方(4、8、16等)
                                    • 选择合适的主键:对于主键表,选择具有良好分布特性的主键,避免热点
                                    • 考虑分区策略:对于大型表,使用分区来提高查询性能和管理数据生命周期

                                    (2) 数据湖分层服务配置

                                    • 检查点间隔:根据实时性需求和系统负载调整检查点间隔,通常在10-60秒之间
                                    • 资源分配:为数据湖分层服务分配足够的资源,确保能够及时处理数据
                                    • 监控:定期监控分层服务的性能和延迟,确保数据及时压缩到Paimon

                                    ( 3)查询优化

                                    ①选择合适的查询模式

                                    • 对于需要最新数据的查询,使用联合读取
                                    • 对于大规模历史数据分析,使用仅Paimon查询

                                    ②过滤下推:在查询中尽可能早地应用过滤条件,利用Paimon的过滤下推能力

                                    ③列裁剪:只选择需要的列,减少数据传输和处理量

                                    Logo

                                    腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

                                    更多推荐