认识Apache Fluss
你有没有被实时数仓的"实时"两个字骗过?
业务那边大促刚开场 3 秒,运营就在群里喊:"核心指标怎么还不刷新?"你盯着监控,发现链路是 Kafka 接消息 → Flink 跑 ETL → 落 Paimon/Iceberg → StarRocks 查——每一跳都差几分钟,串起来就是"分钟级实时"。你说这是实时,运营说这是"事后复盘"。
更扎心的是成本:为了撑住一个双流 Join,Flink 的 State 能涨到 几百 GB 甚至 TB 级,RocksDB compaction 卡死、Checkpoint 超时、任务频繁挂。你调了一周参数,结论是——延迟、一致性、成本这三角,怎么调都顾不全。
传统实时链路的痛,根子不在计算层,而在存储层:流(Kafka)和湖(Paimon/Iceberg)是两套东西,中间全靠 Flink 当胶水硬粘。
2026 年 8 月 6 日,由阿里云 Flink 团队开源、捐赠给 Apache 的 Apache Fluss 正式从孵化器毕业,成为 Apache 顶级项目(TLP),委员会全票通过。它想干的事很野:把"流"和"湖"在存储层直接焊成一条河,让 Lakehouse 真正拥有亚秒级的新鲜度。
什么是 Apache Fluss
Apache Fluss 的全称是 Flink Unified Streaming Storage——一条为实时分析与 AI 而生的"流式存储"河。它不是一个新消息队列,也不是又一个数据库,而是一层湖仓原生的实时流存储底座,把消息队列、在线 KV 存储、流处理状态后端、湖仓冷存储四件事,收敛进同一个存储内核。
一句话定位:让流和湖共享同一份数据、同一个 Schema、同一个抽象,Write Once,Read Anywhere。
它的出身很"实战派":最早由阿里云 Flink 团队为了治流批一体的成本病和链路复杂病而生,2024 年 11 月在 Flink Forward Asia 上海站宣布开源,2025 年 6 月进入 Apache 孵化器。毕业时社区已聚起 157 位贡献者、2000+ GitHub Stars、合并 1700+ PR,Apache 2.0 协议,根正苗红。
Fluss(德语"河流",读作 /flʊs/)的野心,是成为实时 Lakehouse 的开放数据底座——在 Agentic 时代,给智能体喂上"从过去到当下"都新鲜完整的实时上下文。
核心架构:两个进程撑起一条河
别被"分布式"吓到,Fluss 的集群就两类角色,很好记:
- CoordinatorServer(协调节点,集群的大脑):管全局元数据(库、表、Schema)、负责 Tablet 分配与负载均衡、扩缩容时的数据重分布、节点故障时的迁移切换,以及建表/删表/改 Bucket 这类表管理操作。
- TabletServer(数据节点,干活的):负责数据存储、持久化和 I/O,里面又拆成两个核心组件:
- LogStore:只追加的日志(很像数据库 binlog),支撑低延迟流式读取,同时也是 KvStore 的写前日志(WAL)。每个 Segment 由稀疏索引
.index和日志.log组成。 - KvStore:存表数据,支持更新和删除,主键表(PK Table)的点查就靠它,亚毫秒级响应。
- LogStore:只追加的日志(很像数据库 binlog),支撑低延迟流式读取,同时也是 KvStore 的写前日志(WAL)。每个 Segment 由稀疏索引
- 底层还依赖 ZooKeeper 做协调、远程存储层(S3/OSS 等)做持久化。
再加上一个关键角色 Tiering Service(分层服务):它会把 Fluss 里的实时数据周期性自动同步到 Paimon / Iceberg / Hudi / Lance 这些湖仓格式里,热数据在 Fluss、冷数据在湖,通过 Union Read 对外是同一张表视图。
核心特点
Fluss 的六大能力支柱,是冲着"替换一整套实时栈"去的:
- 亚秒级数据新鲜度:持续写入、立即可查,大规模下也能做低延迟分析和实时决策,大促开场几秒就能看到核心指标变化。
- 湖流一体(Lakestream 架构):流和批共享同一份数据和 Schema,不用复制、不用格式转换、不用额外的 ETL 搬运作业——这是它和传统"Kafka+Flink+Paimon 缝合怪"最本质的区别。
- 列式流式(基于 Apache Arrow):服务端列裁剪、谓词下推、分区下推层层叠加,引擎只读它真正要的数据,I/O 和网络开销砍掉一个数量级。
- 存算分离,计算无状态:Flink 只算纯计算,State 和存储交给 Fluss 的 Leader 来管。官方称相比 Kafka 系拓扑,资源成本能低到 85%。State 膨胀、Checkpoint 超时这些老毛病,从根上被卸掉了。
- AI / 向量就绪:同一张底层存储里同时放行存、列存、向量和多模态数据(通过 Lance 集成),实时特征、RAG 上下文、分析查询共用一张 PK Table 的不同视图。
- 内置变更日志与审计:自动生成 append-only 的 changelog,可回放,方便审计、复现和观测。
生态上也很开放,能通过 Flink、Spark、Trino、StarRocks、Doris、DuckDB 读写,冷层对接 Paimon、Iceberg、Hudi、Lance 等开放格式——没有厂商锁定。
如果说 Kafka 是"流式运输",那 Fluss 是"流式存储":它不是要取代 Kafka 做传输,而是要做所有实时分析、AI/ML、亚秒级湖仓背后那层共享的流式存储底座。
谁在用:真实案例,带硬数据
光说架构没意思,看真金白银的生产落地。Fluss 已在阿里、小红书、京东、蚂蚁集团、爱奇艺、Fresha 等公司的生产环境跑起来,主打日志采集、实时数仓、AI 搜索推荐、索引链路、实时特征服务。
- 淘宝闪购(淘宝即时零售):千亿级流量场景下的早期深度实践者。原来用
Kafka + Flink + Paimon缝合,一个双流 Join 要把上百亿条浏览/订单数据全塞进 Flink State,单作业 State 飙到 TB 级,Checkpoint 从几分钟恶化到 10–15 分钟,频繁失败。换成 Fluss 的 Delta Join 后,把全量双流数据从 Flink State 卸载到 Fluss KV 表,变成近乎无状态的双流 Join,配合前缀点查和列裁剪,转化率、漏斗指标做到 30 秒内刷新,并在 618 大促稳定支撑核心实时决策链路。 - 小红书:把 Fluss 的列式读写、冷热分层、湖流一体能力用在索引数据链路等核心场景,显著降低实时链路成本,提升构建与在线消费效率。
- 淘天集团:用 Fluss + Paimon + StarRocks 搭湖流一体链路——Fluss 扛秒级实时,Paimon 沉淀分钟级和历史数据,StarRocks 通过 Union Read 统一查询。相比传统 Kafka+Flink 链路,消费成本降低 80% 以上,开发效率提升 50%。
- 京东、蚂蚁集团:均在积极推进 Fluss 的技术实践与落地,覆盖实时特征、动态状态、持续上下文等 AI × Data 场景。
这些不是 PPT 案例,是千亿级流量、大促峰值、真实降本数据——对一个刚毕业的 Apache 项目来说,这种"出生即生产"的成色相当能打。
初体验:怎么把它跑起来
想亲手摸一下这条河?官方最省事的入口有两个:
方式一:Docker 快速体验镜像。Fluss 官方重新提供了 Quickstart 镜像,一条命令拉起本地最小集群,适合零基础上手、五分钟看到效果。
方式二:用 Flink SQL 直接挂上 Fluss(最像"生产感"的玩法)。在 Flink 里把 Fluss 注册成一个 Catalog,建表、插入、查询全文法:
-- 把 Apache Fluss 注册成 Flink 的 Catalog
CREATE CATALOG fluss_catalog WITH (
'type' = 'fluss',
'bootstrap.servers' = 'coordinator-server:9123'
);
USE CATALOG fluss_catalog;
-- 建一张主键表
CREATE TABLE pk_table (
shop_id BIGINT,
user_id BIGINT,
num_orders INT,
PRIMARY KEY (shop_id, user_id) NOT ENFORCED
) WITH ('bucket.num' = '4');
INSERT INTO pk_table VALUES (1234, 1234, 1);
SELECT * FROM pk_table WHERE shop_id = 1234;如果团队主力是 Spark 或 Trino,走对应的 fluss-spark / REST Catalog 接入即可,姿势类似。想深抠部署细节,直接看官网 Quickstart 文档和 GitHub apache/fluss,欢迎顺手点个 Star——毕竟 157 个贡献者里,下一个可能就是你。
入门别贪多:先把它当成"会存、会查、还能自动沉到湖里的 Kafka+KV 二合一"来用,等吃透了 Delta Join 和湖流分层,再谈替换整套链路。
进阶
更多开源技术干货和学习资料,关注公众号「遇码」,领取专属福利。
