Apache Fluss 升为顶级项目:Flink 实时分析快速上手

8 月 11 日,Apache 软件基金会宣布:由阿里云捐赠并开源的 Apache Fluss 正式从孵化器毕业,成为 Apache 顶级项目(TLP),委员会全票通过。对做 Agent 应用的开发者来说,这值得关注——智能体做决策需要的是秒级新鲜度的上下文,而传统 Lakehouse(Paimon/Iceberg/Hudi)擅长存历史,却接不住「两秒前发生了什么」。Fluss 就是在数据湖之上补上「实时」这一环的流式存储层。

本文带你跑一遍官方 quickstart:用 Docker Compose 起一个 Fluss + Flink 集群,写入流数据、用 lookup join 做维度关联、跑实时分析,全程 Flink SQL,不用写 Java。

Fluss 解决了什么

Fluss(德语「河流」)不是替代数据湖,而是盖在湖上面:热数据、高频更新的数据存在 Fluss,长周期历史留在 Paimon/Iceberg/Hudi,引擎通过 Union Read 拿到「从历史到现在」的统一视图。底层是基于 Apache Arrow 的列式流存储,服务端做列裁剪和条件下推,引擎只读自己需要的字节。

对 Agent 场景最有用的几个特性:

  • 主键表:高 QPS 点查、去重、部分更新、delta join
  • Changelog:状态/决策变更的 append-only 历史,天然是审计和可复现日志
  • 亚秒级新鲜度:数据落库即可查
  • 存算分离:Fluss 管状态和存储,Flink/Spark 只算

第 1 步:起集群

需要 Docker 和 Compose v2。建个目录,写入精简版 docker-compose.yml

services:
  # S3 兼容存储(演示用 RustFS,生产换云对象存储)
  rustfs:
    image: rustfs/rustfs:1.0.0-alpha.83
    ports: ["9000:9000", "9001:9001"]
    environment:
      - RUSTFS_ACCESS_KEY=rustfsadmin
      - RUSTFS_SECRET_KEY=rustfsadmin

  # Fluss 集群
  coordinator-server:
    image: apache/fluss:0.9.1-incubating
    command: coordinatorServer
    environment:
      - FLUSS_PROPERTIES=
          zookeeper.address: zookeeper:2181
          bind.listeners: FLUSS://coordinator-server:9123
          remote.data.dir: s3://fluss/remote-data
          s3.endpoint: http://rustfs:9000
          s3.access-key: rustfsadmin
          s3.secret-key: rustfsadmin

  tablet-server:
    image: apache/fluss:0.9.1-incubating
    command: tabletServer
    environment:
      - FLUSS_PROPERTIES=
          zookeeper.address: zookeeper:2181
          data.dir: /tmp/fluss/data
          remote.data.dir: s3://fluss/remote-data
          s3.endpoint: http://rustfs:9000
          s3.access-key: rustfsadmin
          s3.secret-key: rustfsadmin

  zookeeper:
    image: zookeeper:3.9.2

  # Flink 集群(镜像已内置 Fluss connector + flink-faker)
  jobmanager:
    image: apache/fluss-quickstart-flink:1.20-0.9.1-incubating
    ports: ["8083:8081"]
    command: jobmanager

  taskmanager:
    image: apache/fluss-quickstart-flink:1.20-0.9.1-incubating
    command: taskmanager

  sql-client:
    image: apache/fluss-quickstart-flink:1.20-0.9.1-incubating
    command: /opt/sql-client/sql-client
docker compose up -d
docker compose ps

Flink UI 在 http://localhost:8083;RustFS 控制台(S3 bucket fluss)在 http://localhost:9001,账号 rustfsadmin/rustfsadmin

第 2 步:建 catalog 和表

docker compose run sql-client
CREATE CATALOG fluss_catalog WITH (
  'type' = 'fluss',
  'bootstrap.servers' = 'coordinator-server:9123'
);
USE CATALOG fluss_catalog;

建主键表(点查性能的关键):

CREATE TABLE fluss_order (
  `order_key`   BIGINT,
  `cust_key`    INT NOT NULL,
  `total_price` DECIMAL(15, 2),
  `order_date`  DATE,
  `order_priority` STRING,
  `ptime` AS PROCTIME(),
  PRIMARY KEY (`order_key`) NOT ENFORCED
);

CREATE TABLE fluss_customer (
  `cust_key`  INT NOT NULL,
  `name`      STRING,
  `phone`     STRING,
  `nation_key` INT NOT NULL,
  `acctbal`   DECIMAL(15, 2),
  `mktsegment` STRING,
  PRIMARY KEY (`cust_key`) NOT ENFORCED
);

第 3 步:写入流数据

quickstart 镜像预置了 faker 源表生成演示数据,直接同步进 Fluss:

EXECUTE STATEMENT SET
BEGIN
  INSERT INTO fluss_nation   SELECT * FROM `default_catalog`.`default_database`.source_nation;
  INSERT INTO fluss_customer SELECT * FROM `default_catalog`.`default_database`.source_customer;
  INSERT INTO fluss_order    SELECT * FROM `default_catalog`.`default_database`.source_order;
END;

第 4 步:lookup join 做实时关联

这是 Fluss 的主场:订单流对主键维度表做关联是高 QPS 点查,不是全表扫。

INSERT INTO enriched_orders
SELECT o.order_key, o.cust_key, o.total_price, o.order_date, o.order_priority,
       c.name, c.phone, c.acctbal, c.mktsegment, n.name
FROM fluss_order o
LEFT JOIN fluss_customer FOR SYSTEM_TIME AS OF `o`.`ptime` AS c ON o.cust_key = c.cust_key
LEFT JOIN fluss_nation   FOR SYSTEM_TIME AS OF `o`.`ptime` AS n ON c.nation_key = n.nation_key;

第 5 步:实时分析

SET 'sql-client.execution.result-mode' = 'tableau';
SET 'execution.runtime-mode' = 'batch';
SET 'table.dml-sync' = 'true';

SELECT * FROM enriched_orders LIMIT 2;

-- COUNT(*) 很快:Fluss 维护表级统计,不扫全表
SELECT COUNT(*) FROM enriched_orders;

-- 主键点查
SELECT * FROM fluss_customer WHERE `cust_key` = 1;

多跑几次 COUNT(*),数字会随着 faker 持续生产而增长——Fluss 是连续写入的,结果实时变化。

第 6 步:更新与删除

UPDATE fluss_customer SET `name` = 'fluss_updated' WHERE `cust_key` = 1;
SELECT * FROM fluss_customer WHERE `cust_key` = 1;  -- name 变为 fluss_updated

DELETE FROM fluss_customer WHERE `cust_key` = 1;
SELECT * FROM fluss_customer WHERE `cust_key` = 1;  -- 空

收尾:quit 退出,然后 docker compose down -v

实践建议

  • 别用演示账号上线。 示例里的 rustfsadmin + AssumeRole STS 是给本地 RustFS 的,生产环境换成云厂商的凭证链。
  • 当 Agent 的实时上下文库用。 Fluss 原生支持去重、部分更新、delta join——「实时特征存储」的经典模式直接映射到一张主键表。
  • Changelog 就是免费的决策日志。 Agent 系统要审计状态变更轨迹,Fluss 开箱即生成 append-only changelog。
  • 生态现状: Flink、Spark connector 已可用,StarRocks 在路上;阿里、小红书、京东、蚂蚁已在千亿级流量下生产落地。

资源

滚动至顶部