Apache Fluss 深度拆解:当流存储把 Kafka、Redis、Flink State 与 Iceberg 揉进一个底座——湖仓原生实时架构全链路实战
选题来源:GitHub Trending / Apache 顶级项目动态(2026-08-11 Apache Fluss 从孵化器毕业成为顶级项目)
适用读者:做过实时数仓、被"Kafka + Flink + Redis + Iceberg"四件套折磨过的后端 / 数据工程师
一、背景:实时数仓的"多系统税"
如果你在一家稍微有点规模的公司做过实时业务,大概率见过下面这套"标准答案":
- 消息队列用 Kafka 做事件传输;
- 流处理用 Flink(或 Spark Structured Streaming)做清洗、聚合、打宽;
- 在线 KV用 Redis(或 DynamoDB、HBase)承接亚毫秒级点查,给 API 和特征服务喂数据;
- 离线湖用 Iceberg / Hudi / Delta 落在 S3/OSS 上,喂训练、喂 BI、喂历史回放;
- 同步层再叠一堆自研管道 + 新鲜度监控,把上面四个系统的数据"对齐"。
这套架构本身没毛病,业界跑了快十年。但它有一个被刻意忽视的隐性成本——每一条系统边界,都是数据分叉点。
Kafka 里的订单流、Redis 里的订单快照、Iceberg 里的订单历史,三份数据的"此刻真值"永远差那么几秒到几分钟。你以为你在做实时,其实你在做"勉强不崩的准实时"。一旦某个同步管道悄悄挂了,Redis 里的库存和 Iceberg 里的库存就会出现肉眼不可见的偏差,直到一次对账失败或一次大促超卖才炸出来。
Apache 基金会在 2026-08-11 宣布:由阿里云捐赠的流式存储项目 Fluss 正式从孵化器毕业,成为 Apache 全球顶级项目(TLP),提案在委员会投票中全票通过。官方给 Fluss 的定位一句话概括——"Streaming Storage for Real-Time Analytics & AI",它要做的不是又一个消息队列,而是把上面那套"消息代理 + 在线 KV + 流处理状态后端 + 湖仓冷存"坍缩成同一个底座,让 Lakehouse 真正变成实时的。
这篇文章不吹概念,我们从架构、读写路径、存算分离、列式流分析、湖仓统一,一直拆到能跑的代码和性能调优清单。
二、核心概念:Fluss 到底是什么
2.1 一句话定义与六个能力支柱
Fluss 官方的自我定位是 lakehouse-native streaming storage(湖仓原生的流存储)。它和 Kafka 最大的区别在于词:Kafka 是 streaming transport(传输),Fluss 是 streaming storage(存储)。传输只管把字节从 A 搬到 B 且可重放;存储要回答"我现在这份数据的真值是什么、点查多快、分析怎么扫、冷了往哪沉"。
Fluss 把被拆散的能力重新收拢成六个支柱:
- Unified Architecture(统一架构):一套系统同时承担消息传输、点查、分析查询。架构基础是 PK Table 的"双表示"——同一份数据同时有 Log Store 和 KV Store 两种形态。
- Stream & Lakehouse Unification(流湖统一):实时层和批式层共享同一份 schema、同一份数据,靠 Tiering Service 下沉冷数据 + Union Read 跨冷热统一查询。
- Compute / Storage Separation(存算分离):计算无状态、秒级恢复,状态常驻在 Fluss leader 而不是 Flink slot 里,官方宣称相比 Kafka 拓扑最高可便宜 85%。
- Columnar Streaming Analytics(列式流分析):Log 用 ARROW 列式格式,TabletServer 端做投影、谓词下推、分区裁剪,IO 与网络成本呈数量级下降。
- Feature & Context Stores(特征与上下文存储):行、列、向量三种数据形态同存一个 substrate,在线特征、RAG 上下文、分析查询都是同一张 PK Table 的不同视图。
- Ecosystem Openness(生态开放):Flink、Spark、Trino、StarRocks、Doris、DuckDB 都能读;冷层是 Iceberg / Paimon / Lance 等开放格式,无厂商锁定。
2.2 PK Table 的"双表示"——这是整篇的题眼
传统消息队列里,一条消息进了 Kafka 就是一个 append-only 的 offset 记录,你没法"按主键秒级拿到用户 12345 的当前资料"。要做点查,你得另起一个 Redis 把同一份数据再写一遍。
Fluss 的 Primary-Key Table 把这两件事合并了:落盘时同时维护两份结构——
- Log Store:顺序、可重放、带 offset 的流,给流式消费(等价于 Kafka 的 topic 能力);
- KV Store:按主键组织的点查结构,给亚毫秒级
get(key)(等价于 Redis 的能力)。
写入 (Append / Upsert)
│
▼
┌─────────────────┐
│ PK Table 写入 │
└────────┬────────┘
│ 同一份数据,双写
┌────────┴────────┐
▼ ▼
┌──────────────┐ ┌──────────────┐
│ Log Store │ │ KV Store │
│ (顺序/可重放) │ │ (主键/点查) │
│ offset 有序 │ │ 亚毫秒 get │
└──────┬───────┘ └──────┬───────┘
│ │
流式消费/回放 点查/维表 Join
这意味着:你不再需要在 Kafka 和 Redis 之间写同步管道。一份数据进来,既能当流回放,也能当 KV 点查,新鲜度天然一致,因为源是同一个。
2.3 它和 Kafka 的本质区别
一句话总结对比(官方也专门做了 Kafka 对比页):
| 维度 | Apache Kafka | Apache Fluss |
|---|---|---|
| 定位 | 流式传输(transport) | 流存储(storage) |
| 点查 | 不支持,需外挂 KV | PK Table 内建亚毫秒点查 |
| 列式分析 | 无,按字节流消费 | ARROW 列式 + 服务端裁剪 |
| 湖仓统一 | 需外部 Connect / 自研 | Tiering + Union Read 内建 |
| 状态后端 | 状态在 Flink 侧 | 状态外置到 Fluss leader |
| 向量/多模态 | 无 | Lance 集成,支持 RAG 上下文 |
选型判断:如果你只需要"大规模把事件从 A 搬到 B 再消费",Kafka 依旧是更轻、更成熟的选择;如果你要的是"实时分析 + 在线点查 + 湖仓一体 + AI 上下文",Fluss 是那块被缺了很久的共享底座。
三、架构分析:一个底座怎么把四件套坍缩掉
3.1 组件拓扑
一个最小 Fluss 集群由三类角色组成(外加你熟悉的 ZooKeeper 做协调):
- CoordinatorServer:集群大脑,管元数据、路由、bucket 分配。客户端
bootstrap.servers指向的就是它(默认端口9123)。 - TabletServer:真正扛读写的数据节点,做服务端投影/下推/裁剪的就是它;本地盘放热数据,冷数据落到远端对象存储。
- ZooKeeper:集群协调与选主(和很多 Hadoop 系组件一样)。
数据来源(CDC、事件日志、IoT、点击流)先灌进 Fluss 热层(Coordinator + TabletServer),再通过 Tiering Service 把冷数据沉到 Iceberg / Paimon / Lance。
3.2 写入路径:一次写入,双结构落盘
当一条 UPSERT 进 PK Table:
- 先写 WAL / Log Store(顺序追加,保证持久化与可重放);
- 同步更新内存中的 KV Store(按主键定位,覆盖旧值);
- KV 的变更也通过 WAL 保护,故障后可由 Log 重建;
- 后台把"足够冷"的 KV 分区 / Log 段按 Tiering 规则上传到对象存储(Iceberg/Paimon 文件)。
对客户端来说,这就是一次写入;对系统来说,它同时喂饱了"流"和"KV"两种消费需求。
3.3 读取路径:三种读法,各走各的最快路径
- 点查(PK Lookup):
get(shop_id, user_id)直接走 KV Store,亚毫秒返回,给 API / 特征服务 / 维表 Join 用。 - 流读(Streaming Log):消费者从某个 offset 顺序读 Log Store,等价于 Kafka Consumer,可重放、可回溯。
- 分析读(Columnar Analytics):下游 Flink/Spark/Trino/DuckDB 扫整表时,TabletServer 在 ARROW 列式格式上做服务端投影(只返回要的列)、谓词下推(只返回满足条件的行)、分区裁剪(跳过无关 bucket/分区)。三道裁剪叠加,IO 和网络是数量级下降,而不是线性下降。
3.4 存算分离:状态从 Flink slot 迁到 Fluss leader
这是 Fluss 在成本上"敢说便宜 85%"的底层原因。传统 Flink 作业的状态(Join 状态、聚合状态、维表缓存)全在 TaskManager 的内存/堆外里,作业一挂,状态恢复要靠 checkpoint 重放,又慢又占资源。
Fluss 把"流处理状态后端"也收进自己:状态外置到 Fluss leader,Flink 计算节点变成无状态、可秒级拉起。计算资源按峰谷弹性伸缩,存储独立扩缩,二者解耦。对云上按量付费的场景,这意味着你不用为了"保住状态"而常驻一大堆昂贵的 Flink slot。
3.5 湖仓统一:Tiering + Union Read
热数据在 TabletServer 本地盘,冷数据在对象存储上的 Iceberg/Paimon/Lance 文件。两层共享同一 schema,对查询引擎暴露成一张表:
- Tiering Service 负责把冷数据自动下沉、对历史分区做 compaction;
- Union Read 让一次查询同时命中热层(实时)和冷层(历史),对外是单一数据源。
结果:你既不用为"实时性"单独维护 Kafka→Iceberg 同步管道,也不用为"历史回放"再起一套批管道。实时和批式读的是同一份真值。
3.6 多模态:Lance 向量与 RAG 上下文
第五根支柱很多人会忽略,但恰恰契合 2026 年的"实时 AI"叙事。Fluss 通过 Lance 集成把向量也变成同一 substrate 上的一种数据形态,行、列、向量同存。在线特征、RAG 检索上下文、流式分析都是同一张 PK Table 的不同视图——这正是它敢喊"Agentic Lake 最强底座"的技术底气:企业级 Agent 做实时决策时,所需的"秒级更新的上下文"可以直接从 Fluss 一个点拿到,而不是在五个系统间拼凑。
四、代码实战
下面所有可直接跑的示例都来自 Fluss 官方 0.9 文档(quickstart-flink),我做了连贯化整理。最后我会用一段从零实现的 Python 迷你版把"双表示 + 列式裁剪"讲透——注意那段是我自己的教学代码,不是 Fluss 客户端 API,目的是让你看清内核原理。
4.1 一键拉起:Docker Compose
最小可玩环境包含:RustFS(S3 兼容对象存储,做分层存储)、Fluss(CoordinatorServer + TabletServer + ZooKeeper)、Flink(JobManager + TaskManager + SQL Client)。
services:
# RustFS: S3 兼容存储,承接 Fluss 的分层(tiered)冷数据
rustfs:
image: rustfs/rustfs:1.0.0-alpha.83
ports:
- "9000:9000"
- "9001:9001"
environment:
- RUSTFS_ACCESS_KEY=rustfsadmin
- RUSTFS_SECRET_KEY=rustfsadmin
- RUSTFS_CONSOLE_ENABLE=true
volumes:
- rustfs-data:/data
command: /data
rustfs-init:
image: minio/mc
depends_on:
- rustfs
entrypoint: >
/bin/sh -c "
until mc alias set rustfs http://rustfs:9000 rustfsadmin rustfsadmin; do
echo 'Waiting for RustFS...'; sleep 1;
done;
mc mb --ignore-existing rustfs/fluss;
"
# Fluss 集群:CoordinatorServer + TabletServer + ZooKeeper
coordinator-server:
image: apache/fluss:0.9.1-incubating
command: coordinatorServer
depends_on:
- zookeeper
- rustfs-init
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
s3.region: us-east-1
s3.path-style-access: true
s3.assumed.role.arn: arn:aws:iam::000000000000:role/rustfsadmin
s3.assumed.role.sts.endpoint: http://rustfs:9000
tablet-server:
image: apache/fluss:0.9.1-incubating
command: tabletServer
depends_on:
- coordinator-server
environment:
- |
FLUSS_PROPERTIES=
zookeeper.address: zookeeper:2181
bind.listeners: FLUSS://tablet-server:9123
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
s3.region: us-east-1
s3.path-style-access: true
s3.assumed.role.arn: arn:aws:iam::000000000000:role/rustfsadmin
s3.assumed.role.sts.endpoint: http://rustfs:9000
zookeeper:
restart: always
image: zookeeper:3.9.2
# Flink 集群(镜像里已打包 Fluss connector + flink-faker + S3 支持)
jobmanager:
image: apache/fluss-quickstart-flink:1.20-0.9.1-incubating
ports:
- "8083:8081"
command: jobmanager
environment:
- |
FLINK_PROPERTIES=
jobmanager.rpc.address: jobmanager
taskmanager:
image: apache/fluss-quickstart-flink:1.20-0.9.1-incubating
depends_on:
- jobmanager
command: taskmanager
environment:
- |
FLINK_PROPERTIES=
jobmanager.rpc.address: jobmanager
taskmanager.numberOfTaskSlots: 10
taskmanager.memory.process.size: 2048m
taskmanager.memory.framework.off-heap.size: 256m
sql-client:
image: apache/fluss-quickstart-flink:1.20-0.9.1-incubating
depends_on:
- jobmanager
command: /opt/sql-client/sql-client
environment:
- |
FLINK_PROPERTIES=
jobmanager.rpc.address: jobmanager
rest.address: jobmanager
volumes:
rustfs-data:
启动:
docker compose up -d
docker compose ps
# 打开 http://localhost:8083 看 Flink Web UI
# 打开 http://localhost:9001 (rustfsadmin/rustfsadmin) 看分层存储桶
4.2 Flink SQL:建 Catalog、建 PK Table、流式写入
进入 SQL Client:docker compose run sql-client。
-- 1) 注册 Fluss Catalog
CREATE CATALOG fluss_catalog WITH (
'type' = 'fluss',
'bootstrap.servers' = 'coordinator-server:9123'
);
USE CATALOG fluss_catalog;
-- 2) 建主键表(注意 bucket.num 控制分桶并行度)
CREATE TABLE fluss_order (
`order_key` BIGINT,
`cust_key` INT NOT NULL,
`total_price` DECIMAL(15, 2),
`order_date` DATE,
`order_priority` STRING,
`clerk` 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) 把 CDC / 事件源持续灌进 Fluss(EXECUTE STATEMENT SET 批量提交)
EXECUTE STATEMENT SET
BEGIN
INSERT INTO fluss_order SELECT * FROM `default_catalog`.`default_database`.source_order;
INSERT INTO fluss_customer SELECT * FROM `default_catalog`.`default_database`.source_customer;
END;
这里 fluss_order / fluss_customer 就是 PK Table:进来的数据同时进入 Log Store(可被流式消费)和 KV Store(可被点查)。bucket.num 决定分桶数,直接影响点查并行度和热点分布,是后面性能优化的核心旋钮。
4.3 点查 + Lookup Join:一套数据两种吃法
正是因为 PK Table 内建 KV,Flink 的 Lookup Join(维表 Join)可以直接打在 Fluss 上,亚毫秒拿到客户 / 国家维度,把订单流实时打宽:
CREATE TABLE 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
);
INSERT INTO enriched_orders
SELECT
o.order_key, o.cust_key, o.total_price, o.order_date,
o.order_priority, o.clerk,
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;
注意:在传统架构里,fluss_customer 这份维表你本该在 Redis 里再存一份;在 Fluss 里,它就是 PK Table 的 KV 视图,零额外同步。
4.4 从零实现一个迷你"双表示"流存储(教学版)
下面这段 Python 没有用任何外部依赖,目的是把 Fluss 的"双表示 + 列式服务端裁剪 + 分层"内核讲透。它不是 Fluss 客户端 API,请勿直接当生产代码用。
"""
mini_fluss.py —— 教学用迷你流存储,讲透 Fluss 的核心思想:
1) PK Table 双表示:Log Store(顺序可重放) + KV Store(主键点查)
2) 列式存储 + 服务端投影/谓词下推/分区裁剪
3) 分层:热层本地内存,冷层对象存储(模拟)
作者:三哥 / 程序员茄子(仅用于原理演示)
"""
from typing import Any, Dict, List, Optional
class LogStore:
"""顺序、带 offset、可重放的流(对应 Kafka 的 topic 能力)。"""
def __init__(self) -> None:
self._segments: List[Dict[str, Any]] = []
def append(self, record: Dict[str, Any]) -> int:
offset = len(self._segments)
self._segments.append({"offset": offset, **record})
return offset
def read_from(self, offset: int) -> List[Dict[str, Any]]:
"""从指定 offset 顺序消费(流读 / 回放)。"""
return self._segments[offset:]
class KvStore:
"""按主键组织的点查结构(对应 Redis 的点查能力)。"""
def __init__(self) -> None:
self._data: Dict[Any, Dict[str, Any]] = {}
def upsert(self, key: Any, value: Dict[str, Any]) -> None:
self._data[key] = value
def get(self, key: Any) -> Optional[Dict[str, Any]]:
return self._data.get(key)
class ColumnarTable:
"""极简列式存储:每列一个 list,演示服务端裁剪。"""
def __init__(self, columns: List[str]) -> None:
self.columns = columns
self._cols: Dict[str, List[Any]] = {c: [] for c in columns}
def insert_row(self, row: Dict[str, Any]) -> None:
for c in self.columns:
self._cols[c].append(row.get(c))
def scan(self, project: List[str], predicate=None) -> List[Dict[str, Any]]:
"""
服务端裁剪演示:
- project:只返回需要的列(投影下推,省网络/IO)
- predicate:只在服务端过滤满足条件的行(谓词下推)
"""
n = len(self._cols[self.columns[0]])
out: List[Dict[str, Any]] = []
for i in range(n):
row = {c: self._cols[c][i] for c in self.columns}
if predicate and not predicate(row):
continue
out.append({c: row[c] for c in project})
return out
class MiniFlussPkTable:
"""PK Table:一次写入,双表示落盘。"""
def __init__(self, primary_key: str) -> None:
self.pk = primary_key
self.log = LogStore()
self.kv = KvStore()
self.cold: List[Dict[str, Any]] = [] # 模拟对象存储冷层
def upsert(self, record: Dict[str, Any]) -> int:
offset = self.log.append(record)
self.kv.upsert(record[self.pk], record)
return offset
def lookup(self, key: Any) -> Optional[Dict[str, Any]]:
"""点查 —— 走 KV Store,亚毫秒。"""
return self.kv.get(key)
def stream(self, offset: int = 0) -> List[Dict[str, Any]]:
"""流读 —— 走 Log Store。"""
return self.log.read_from(offset)
def tier_to_cold(self) -> None:
"""分层:把冷数据批量下沉到对象存储(模拟 Tiering Service)。"""
for rec in self.log._segments:
self.cold.append(rec)
# ---------- 跑一个端到端示例 ----------
if __name__ == "__main__":
orders = MiniFlussPkTable(primary_key="order_key")
orders.upsert({"order_key": 1, "cust_key": 100, "total_price": 99.5, "mkt": "AUTO"})
orders.upsert({"order_key": 2, "cust_key": 200, "total_price": 12.0, "mkt": "BUILDING"})
orders.upsert({"order_key": 3, "cust_key": 100, "total_price": 250.0, "mkt": "AUTO"})
# 点查:等价于 Redis GET,无需外挂 KV
print("点查 order_key=2 ->", orders.lookup(2))
# 流读:等价于 Kafka Consumer,从 offset 0 重放
print("流读全部 ->", orders.stream(0))
# 列式分析:只投影 total_price,只保留 mkt=AUTO(服务端裁剪)
ct = ColumnarTable(["order_key", "total_price", "mkt"])
for r in orders.stream(0):
ct.insert_row(r)
cheap_auto = ct.scan(
project=["order_key", "total_price"],
predicate=lambda row: row["mkt"] == "AUTO",
)
print("列式裁剪结果(仅 AUTO 且只投影两列) ->", cheap_auto)
# 分层:下沉冷数据
orders.tier_to_cold()
print("冷层记录数 ->", len(orders.cold))
运行输出会清晰展示:同一次 upsert 既进了 Log 也进了 KV;点查走 KV、流读取 Log、分析走列式裁剪——三个消费需求来自同一份真值。这就是 Fluss 把 Kafka + Redis 收进一个底座的直觉来源。
4.5 与 Kafka 的迁移对照
如果你现在用 Kafka + Redis + Flink,向 Fluss 收敛时的心智变化:
| 你现在做的 | Fluss 里怎么表达 |
|---|---|
| Kafka topic 传事件 | Fluss Log Store(PK Table 的流视图) |
| Redis 存维表点查 | Fluss KV Store(同一 PK Table 的点查视图) |
| Flink 状态后端 | 外置到 Fluss leader,Flink 变无状态 |
| Kafka→Iceberg 同步管道 | Tiering Service 自动下沉 + Union Read |
| 自研新鲜度监控 | 单一数据源,不再需要跨系统对账 |
五、性能优化清单
Fluss 不是"装上就快",要把它喂好,下面几条是实战里最值钱的旋钮。
5.1 bucket.num 是点查与并行的第一旋钮
PK Table 按主键哈希分桶。桶数太少 → 单桶过热、点查热点、并行度上不去;桶数太多 → 元数据和小文件膨胀、KV 维护成本上升。经验值:单桶目标承载 百万级 ~ 千万级主键行,再按 QPS 反推。热点主键(比如大 V 用户)要靠主键设计打散(加 salt / 复合主键)而非单纯加桶。
5.2 服务端裁剪:把计算推到存储侧
这是 Fluss 相对"裸 Kafka + 你自己在 Flink 里全量拉"的最大收益点。确保查询走的是列裁剪 + 谓词下推路径:
- 只
SELECT需要的列,别SELECT *; - 把过滤条件下推到存储(Flink/Trino/DuckDB 的 planner 会自动做,但你要避免在 UDF 里把整行拉出来再过滤);
- 利用分区/桶裁剪,让无关 bucket 根本不被扫描。
5.3 点查优化:PK Table + 必要二级索引
亚毫秒点查依赖 KV Store 的主键结构。如果你的点查维度不是主键,需要在上游把"查询键→主键"的映射建模进表设计(例如用 (shop_id, user_id) 复合主键直接覆盖最常见的点查路径),而不是在查询时再做全表过滤。
5.4 分层存储调优:热层本地盘,冷层对象存储
- 热层(TabletServer 本地盘):用高 IOPS 本地 NVMe 放活跃 KV 与近期 Log,点查和最新流读都在这。
- 冷层(S3/OSS/Iceberg/Paimon):调 Tiering 的触发阈值(数据年龄 / 大小),让"够冷"的分区自动下沉,并把 compaction 放到低峰期,避免和实时读写抢 IO。
- Union Read 让你对冷热无感,但冷层访问延迟天然更高,实时看板要显式约束查询时间窗,别让热查询意外扫到冷分区。
5.5 Flink 集成优化
- Lookup Join 缓存:维表 Join 打开 Flink 的 lookup cache,减少重复点查;但缓存 TTL 要和数据新鲜度要求匹配,别为了性能牺牲一致性。
- State Backend 外置:把 Join / 聚合状态放到 Fluss,而不是堆在 TaskManager,作业失败恢复从"重放 checkpoint"变成"无状态拉起 + 读 Fluss 状态",恢复时间从分钟级降到秒级。
- 并行度对齐桶数:Flink 消费 Fluss 的并行度尽量与
bucket.num/ TabletServer 数对齐,避免个别 subtask 成为瓶颈。
5.6 成本:存算分离的真红利
把"状态常驻在昂贵 Flink slot"改成"状态在 Fluss、Flink 无状态弹性伸缩",是 Fluss 宣称比 Kafka 拓扑便宜最高 85% 的来由。真实收益取决于你的峰谷比:峰谷越悬殊,把计算弹性出去省得越多;如果是常驻平稳流量,收益会打折。这点和所有存算分离架构一致——它不是免费的午餐,但是峰谷明显的实时业务的一剂良药。
六、总结与展望
6.1 它解决的真问题
Fluss 不是在造一个新轮子,而是在填一个被行业习惯性忽略的坑:实时数据栈的"分叉税"。当一份数据的真值被拆在 Kafka、Redis、Iceberg 三处,你花在"对齐"上的工程力,往往超过花在"业务"上的。Fluss 用 PK Table 的双表示 + 湖仓统一 + 状态外置,把这几个系统重新收拢成一个 source of truth。
6.2 什么时候该用,什么时候别用
适合:实时分析 + 在线点查 + 湖仓一体叠加的场景;Agentic / 实时 AI 需要"秒级更新的统一上下文"的场景;峰谷明显的实时业务(吃存算分离红利)。
暂不建议:你只需要单纯的高吞吐事件传输、团队对 Kafka 生态极其熟练且不需要点查/湖仓统一——Kafka 更轻更稳;或者你需要的是严格 Exactly-Once 跨多系统的分布式事务(那是另一类问题,Fluss 不替你做全局 2PC)。
6.3 生态与里程碑
截至 2026-08-11,Fluss 已从 Apache 孵化器毕业成为顶级项目(TLP),已在阿里、小红书、爱奇艺、蚂蚁等大规模落地。开放生态是其关键卖点:热层由 Flink/Spark/Trino/StarRocks/Doris/DuckDB 直接读,冷层是 Iceberg/Paimon/Lance 开放格式,治理归 ASF,无厂商锁定。
6.4 结语
流式系统过去十年在"传输"上卷到了极致(Kafka 几乎成了事实标准),但"存储"这一层一直被将就——我们用一堆系统拼出了实时,却始终在为它们的边界买单。Fluss 的意义,是第一次把"流"和"存储"、"实时"和"湖仓"、"点查"和"分析"当成同一个问题的不同侧面来设计。Lakestream(湖流一体) 会不会成为下一个十年的默认范式,现在下结论还早;但至少,被 Kafka + Redis + Flink + Iceberg 四件套折磨过的工程师,多了一个值得认真评估的选择。
本文基于 Apache Fluss 0.9 官方文档与 2026-08-11 TLP 公告整理,代码示例除标注"教学版"的 Python 外均来自官方 quickstart。版本演进快,落地前请以你所用工版本的官方文档为准。