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-clientdocker compose up -d
docker compose psFlink UI 在 http://localhost:8083;RustFS 控制台(S3 bucket fluss)在 http://localhost:9001,账号 rustfsadmin/rustfsadmin。
第 2 步:建 catalog 和表
docker compose run sql-clientCREATE 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 在路上;阿里、小红书、京东、蚂蚁已在千亿级流量下生产落地。
