图解 Fluss(一):一张图看清整体架构,两张图理解核心服务
2026/8/31 4:47:37 网站建设 项目流程

图解 Fluss(一):一张图看清整体架构,两张图理解核心服务

阅读本文你将了解:Fluss 集群由哪几类节点构成、Coordinator 与 TabletServer 的职责边界为什么这样切分、元数据和数据分别存在哪里、以及 Coordinator 的"两阶段启动 + ZK Fence"设计如何防止脑裂。

配套图表:component-01-architectureclass-01-server-hierarchy

难度:⭐ | 适合人群:第一次接触 Fluss 的架构师与开发


一、从一个真实的选型场景说起

假设你是某电商公司的数据平台负责人。当前的实时链路是这样的:

业务库 → Canal → Kafka → Flink → Redis(维表) + HBase(结果) → 应用 ↓ Kafka → Flink → Iceberg(离线)

这套架构跑了两年,痛点越来越明显:

痛点具体表现
链路长同一份数据被写 3 遍(Kafka、HBase、Iceberg),存储成本翻三倍
维表关联慢维表存 Redis,每次关联一次网络往返,大促时 RT 从 5ms 涨到 80ms
状态太重Flink 双流 Join 的状态 8TB,Checkpoint 要 12 分钟,一次故障恢复半小时
查不了历史Kafka 只保留 7 天,想查一个月前的数据得去 Iceberg,两套语法

老板问你:能不能简化?

Fluss 给出的答案是:让 Kafka 具备 KV 点查能力、列裁剪能力,并原生对接湖格式。但要判断这个答案是否成立,你得先看懂它的架构。


二、图 1:整体架构组件图

2.1 四层结构

整张图从上到下可以切成四层:

┌─────────────────────────────────────────────────┐ │ ① 计算引擎层 Flink / Spark / Trino / StarRocks │ ← 无状态,可随意扩缩 ├─────────────────────────────────────────────────┤ │ ② Fluss 集群 CoordinatorServer (3 节点 HA) │ ← 大脑:元数据 + 调度 │ TabletServer (N 节点) │ ← 手脚:数据存储 + 读写 ├─────────────────────────────────────────────────┤ │ ③ 协调层 ZooKeeper │ ← 元数据存储 + 选主 ├─────────────────────────────────────────────────┤ │ ④ 存储层 S3 对象存储 / Iceberg / Paimon │ ← 冷数据归档 └─────────────────────────────────────────────────┘

2.2 Coordinator 内部:五个组件各司其职

图中 Coordinator 节点里画了五个组件,它们的分工是:

组件职责类比 Kafka
MetadataManagerDatabase / Table / Schema 的 CRUD,持久化到 ZKKafka 的AdminManager+ ZK 元数据
CoordinatorEventProcessor单线程事件循环,串行处理所有协调事件类似 Kafka Controller 的事件队列
AutoPartitionManager自动创建分区(比如按天自动建分区)Kafka 无对应组件
LakeTableTieringManager调度热数据 → 冷数据的分层任务Kafka 无对应组件
RebalanceManager负载不均衡时生成迁移计划类似 Kafka 的分区重分配

关键设计:Coordinator 不存任何业务数据。

这一点和 Kafka 的 Controller 有本质区别。Kafka 的 Controller 是"Broker 兼职",Controller 节点本身也存分区数据;而 Fluss 的 Coordinator 是一个纯粹的协调者,它只持有元数据视图和调度逻辑。

这带来两个直接好处:

  1. Coordinator 可以很轻量:4C8G 的机器就能扛住几万 Tablet 的元数据管理。
  2. 扩缩容不影响数据:加 TabletServer 时,Coordinator 只需要重新分配 Tablet,自己不用迁移任何东西——这就是存算分离。

2.3 TabletServer 内部:两套存储引擎

// TabletServer 持有的核心组件(简化)publicclassTabletServerextendsServerBase{privateLogManagerlogManager;// LogStore:日志存储privateKvManagerkvManager;// KvStore:键值存储privateReplicaManagerreplicaManager;// 副本管理privateRpcServerrpcServer;// RPC 服务}

图中 TabletServer 下方有一条关键的 note,它解释了 Fluss 两类表的存储差异:

PK 表:LogStore (WAL) + KvStore (RocksDB) Log 表:仅 LogStore 两者按 Bucket 切分,分布到不同 TabletServer

这就是 Fluss 相对 Kafka 最核心的增量:同一套集群,既可以当"消息队列"用(Log 表),也可以当"KV 数据库"用(PK 表)。

Log 表PK 表
存储只有LogTabletLogTablet(WAL) +KvTablet(RocksDB)
能力追加写、流式读追加写、流式读、按主键点查按主键更新/删除
典型场景埋点日志、CDC 流维表、宽表、实时聚合结果
存储成本1x约 2.5x(多一份 RocksDB)

2.4 连接关系逐条拆解

图里的箭头看似多,其实只有四组:

① 计算引擎 → 集群

Flink --> Coordinator : Fluss Catalog DDL/DML Flink --> TabletServer : 流式读写 / PK Lookup Spark --> TabletServer : 批量读取

注意 Flink 是双通道的:DDL(建表改表)走 Coordinator,实际数据读写直连 TabletServer。Coordinator 不在数据路径上,这是高吞吐的前提。

② 集群 → ZooKeeper

Coordinator --> ZK : 元数据持久化 + Leader 选举 TabletServer --> ZK : 注册/心跳 + Tablet 元数据

③ Coordinator → TabletServer(控制流)

Coordinator --> TabletServer : Tablet 分配 / Rebalance 指令

④ 数据分层(冷热分离)

TabletServer --> S3 : Tiering (Arrow → Parquet) S3 --> Iceberg : Compaction 提交

这条链路是 Fluss 的 Streaming Lakehouse 能力,第 5 篇会展开。

2.5 回到选型场景

现在可以回答开头的问题了:

原痛点Fluss 的解法
链路长、存三份PK 表一份数据同时支持流式读和点查;冷热分层自动归档到 Iceberg
维表关联慢维表直接放 Fluss PK 表,Lookup 走 LRU + Bloom + RocksDB,亚毫秒级
状态太重Delta Join 把 Join 状态外部化(第 3 篇会讲)
查不了历史Union Read 统一查询热层 + 冷层,一套 SQL

三、图 2:核心服务类图

看完了宏观架构,我们下沉到代码层面,看看这两个服务是怎么实现的。

3.1 ServerBase:公共基类

publicabstractclassServerBase{protectedConfigurationconf;protectedPluginManagerpluginManager;publicabstractvoidstartServices()throwsException;publicabstractvoidcloseAsync(@NullableCompletableFuture<Void>closeResult);protectedConfigurationloadConfiguration(String[]args){/* ... */}protectedvoidapplyServerDefaultConfigurations(Configurationconfiguration){/* ... */}}

ServerBase抽象了两类服务的共性:配置加载、插件管理、生命周期。子类只需实现startServices()closeAsync()

这是很标准的模板方法模式,但它有个值得注意的设计:closeAsync返回的是CompletableFuture而不是void

为什么?因为 Fluss 的资源释放是有依赖顺序的。Coordinator 失去 Leader 身份时,必须按逆序清理:先停EventProcessor(不再处理新事件)→ 再关ChannelManager(断开与 TabletServer 的连接)→ 最后清理RpcClient。异步返回让调用方能编排这个顺序。

3.2 CoordinatorServer:两阶段启动

这是整张图最值得关注的设计。看图右侧的 note:

两阶段启动: 1. initCoordinatorStandby() 基础设施 2. initCoordinatorLeader() 竞选成功后加载协调逻辑 防脑裂: ZK Fence (epoch 递增)

为什么不能一次性启动完?

因为一个 Fluss 集群通常部署 3 个 Coordinator,但同一时刻只有 1 个是 Leader。如果三个节点启动时就把EventProcessorChannelManager这些"Leader 专属资源"全部初始化了,会出两个问题:

  1. 浪费资源(Standby 节点白占内存和线程池)
  2. 更严重:万一发生网络分区,两个节点都以为自己是 Leader,就会同时下发冲突的 Tablet 分配指令

所以 Fluss 的做法是:

// 阶段一:所有节点都执行,只起基础设施privatevoidinitCoordinatorStandby(){// RPC 服务(仅健康检查端口)// ZooKeeper 连接// MetadataManager(只读模式)// DynamicConfigManager(监听配置变更)}// 阶段二:只有竞选成功的节点执行privatevoidinitCoordinatorLeader()throwsException{// 创建 CoordinatorEventProcessor(单线程事件循环)// 启动 AutoPartitionManager// 创建 CoordinatorChannelManager(主动连接所有 TabletServer)// 创建默认数据库 "fluss"}

3.3 防脑裂:ZK Fence

两阶段启动解决的是"资源按需分配",但还需要解决"旧 Leader 诈尸"。

考虑这个场景:

1. Coordinator A 是 Leader,epoch = 5 2. A 发生长时间 GC,ZK 会话超时,A 被摘除 3. Coordinator B 竞选成功,成为新 Leader,epoch = 6 4. A 的 GC 结束,恢复运行,它以为自己还是 Leader 5. A 和 B 同时下发指令 → 集群状态混乱

Fluss 的解法是fenceBecomeCoordinatorLeader()

classZooKeeperClient{/** * 递增 epoch 并返回新的 ZkEpoch。 * 若返回 null,说明存在更新的 epoch,本次竞选失败。 */@NullablepublicZkEpochfenceBecomeCoordinatorLeader(){/* ... */}}

原理和 Kafka 的 Controller epoch、HDFS 的 fencing token 一脉相承:每次竞选 Leader 时把 ZK 上的 epoch 计数器 +1,所有下发给 TabletServer 的指令都携带当前 epoch。TabletServer 只接受 epoch 大于已见最大值的指令,旧 Leader 的指令因为 epoch 过期而被直接丢弃。

3.4 组合关系:CoordinatorServer 的五个核心依赖

CoordinatorServer *-- CoordinatorService // 处理 RPC 请求 CoordinatorServer *-- CoordinatorEventProcessor // 事件循环 CoordinatorServer *-- MetadataManager // 元数据 CoordinatorServer *-- CoordinatorLeaderElection // 选主 CoordinatorServer *-- ZooKeeperClient // ZK 客户端

注意这里用的是组合(*--)而不是聚合:这些组件的生命周期完全由CoordinatorServer掌管,随它创建、随它销毁。

其中CoordinatorService是 RPC 请求的第一站:

classCoordinatorService{publicvoidhandleCreateTable(...){/* ... */}publicvoidhandleAlterTable(...){/* ... */}publicvoidhandleDropTable(...){/* ... */}publicvoidhandleFetchRequest(...){/* ... */}}

但它不直接改元数据,而是把请求包装成事件丢给CoordinatorEventProcessor的队列:

CoordinatorEventProcessor --> RebalanceManager : 触发重平衡

为什么多此一举?因为单线程事件循环是最简单可靠的并发模型。所有的元数据变更、Tablet 分配、Rebalance 决策,都在这一个线程里串行执行,天然避免了锁竞争和状态不一致。这也是 Kafka Controller 在 KRaft 模式下坚持的设计。

3.5 TabletServer:更简单的四个组件

classTabletServerextendsServerBase{privateLogManagerlogManager;privateKvManagerkvManager;privateReplicaManagerreplicaManager;privateRpcServerrpcServer;privatevoidloadTablet(TabletAssignmentassignment){/* ... */}privatevoidregisterToCoordinator(){/* ... */}}

TabletServer 没有状态机,它启动后只做两件事:

  1. registerToCoordinator()— 向 Coordinator 注册自己
  2. loadTablet(assignment)— 按 Coordinator 下发的分配方案加载 Tablet

它是一个纯粹的"执行者":Coordinator 让它加载哪个 Tablet 它就加载,让它把哪个 Tablet 迁走它就迁走。这种"无脑执行"的设计让扩缩容变得非常安全——第 4 篇讲 Rebalance 时会看到这一点。


四、动手验证

看完图,建议立刻起一个集群验证。用官方 Docker 镜像:

dockerrun-d--namefluss-coordinator-server\-p9123:9123-p9124:9124\fluss/fluss:0.9.1-incubating coordinatorServerdockerrun-d--namefluss-tablet-server\-p9125:9125\fluss/fluss:0.9.1-incubating tabletServer

进入 SQL 客户端,验证 Coordinator 的两阶段启动留下了什么痕迹:

-- 默认数据库 "fluss" 是 initCoordinatorLeader() 里创建的SHOWDATABASES;-- 输出:flussSHOWTABLES;-- 空-- 建一张 Log 表(只有 LogStore)CREATETABLEclick_events(event_idBIGINT,user_idBIGINT,event_type STRING,event_timeTIMESTAMP(3))WITH('bucket.num'='8');-- 再建一张 PK 表(LogStore + KvStore)CREATETABLEuser_profile(user_idBIGINT,name STRING,city STRING,PRIMARYKEY(user_id)NOTENFORCED)WITH('bucket.num'='16','table.merge-engine'='deduplicate');

然后去 TabletServer 的数据目录看一眼,你会看到两类表在磁盘上的差异:

/data/fluss-server/data/ ├── log/ │ └── fluss/ │ ├── click_events/ ← Log 表:只有 log 目录 │ │ └── bucket-0/ │ │ ├── 00000000000000000000.log │ │ ├── 00000000000000000000.index │ │ └── 00000000000000000000.timeindex │ └── user_profile/ ← PK 表:log + kv 两个目录 │ └── bucket-0/ │ ├── 00000000000000000000.log ← WAL │ └── kv/ │ └── rocksdb/ │ ├── CURRENT │ ├── MANIFEST-000001 │ └── 000005.sst ← RocksDB 数据文件

看到这个目录结构,figure 1 里那句"PK 表 = LogStore + KvStore"就从抽象概念变成了磁盘上的真实文件。


五、生产实践要点

5.1 部署规格建议

节点CPU / 内存磁盘副本数说明
CoordinatorServer4C / 8G50G(日志)3不存业务数据,资源需求低
TabletServer16C / 64GSSD,按数据量≥3本地 SSD 是性能关键
ZooKeeper2C / 4G100G3 或 5独立部署,不要和业务混用

踩坑提醒:很多人图省事把 ZK 和 TabletServer 部署在同一台机器上。TabletServer 在高写入时磁盘 IO 打满,会导致 ZK 心跳超时,进而触发 Coordinator 频繁重选。这个故障现象是"集群周期性不可用",排查起来很痛苦。

5.2 关键配置

# coordinator-server.yaml coordinator.host: 0.0.0.0 coordinator.port: 9123 zookeeper.address: zk1:2181,zk2:2181,zk3:2181 zookeeper.root: /fluss # tablet-server.yaml tablet-server.host: 0.0.0.0 tablet-server.port: 9125 data.dir: /data/fluss-server/data tablet-server.num-network-threads: 8 tablet-server.num-worker-threads: 16

5.3 健康检查

# Coordinator 是否已经选出 Leaderechosrvr|nclocalhost2181|grepMode# 查看当前 Leader(ZK 上的临时节点)zkCli.sh get /fluss/coordinator/leader# Coordinator 日志中确认两阶段grep"initCoordinatorStandby"logs/coordinator-server.loggrep"initCoordinatorLeader"logs/coordinator-server.log

六、排障手册

现象可能原因排查方向
集群起来后所有 DDL 都超时没有 Leader检查 3 个 Coordinator 是否都在Standby,看 ZK 选举路径下有没有candidate_节点
出现两个 Leaderepoch fence 失效检查 ZK 的/fluss/coordinator/epoch节点值;确认 ZK 集群本身没有脑裂
Leader 频繁切换ZK 会话超时检查 GC 停顿(grep "Total time for which application threads were stopped")和网络 RTT
TabletServer 注册不上网络或时钟确认 9123 端口可达;检查机器时钟偏差(超过zookeeper.session-timeout的 1/3 会出问题)
建表报 “tablet allocation failed”TabletServer 不足bucket.num大于可用 TabletServer 数,或剩余磁盘低于水位线

七、小结

回顾这两张图的核心结论:

  1. 四层架构:计算引擎 → Fluss 集群(Coordinator + TabletServer)→ ZooKeeper → 对象存储/湖格式。
  2. Coordinator 不存数据,只管元数据和调度,这是存算分离的基础,也是它能做得很轻量的原因。
  3. Flink 双通道:DDL 走 Coordinator,数据读写直连 TabletServer,Coordinator 不在数据路径上。
  4. 两阶段启动 + ZK FenceinitCoordinatorStandby()起基础设施,initCoordinatorLeader()起 Leader 专属资源,配合递增 epoch 防止旧 Leader 诈尸。
  5. 单线程事件循环CoordinatorEventProcessor串行处理所有协调事件,用简单模型换可靠性。

下一篇,我们下沉到 TabletServer 内部,拆解LogStore 与 KvStore 双引擎,看看一次INSERT到底在磁盘上留下了什么。


需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询