图解 Fluss(二):存储引擎双引擎 —— LogStore 与 KvStore
阅读本文你将了解: LogSegment 的稀疏索引为什么和 Kafka 长得一样、PK 表"先写 WAL 再写 RocksDB"的顺序为什么不能反、RocksDB 那几个关键参数的含义、以及三种 Merge Engine 该怎么选。
配套图表:class-02-storage-engine、class-04-merge-engine
难度:⭐⭐⭐ | 适合人群:要做表设计与性能调优的开发
一、场景:一次"加个字段"引发的事故
先看一个真实场景。
某风控团队有一张用户行为表,200 多个字段,日均 20 亿条。他们用 Flink 做实时特征计算,实际只用到其中 5 个字段:
SELECT user_id, event_type, amount, device_id, event_time
FROM user_behavior
WHERE event_type = 'PAY';
但底层是 Kafka,Flink 的 Source 必须把整行 200 个字段全部拉下来,再在内存里做投影和过滤。结果:
- 网络带宽:高峰期 12 Gbps,其中 97% 是无效字段
- Flink TM 内存:经常 OOM,因为反序列化后的完整对象太大
- 端到端延迟:P99 达到 3 秒
后来他们加了第二个需求:风控规则要查用户最近 30 天的累计交易额。原来的做法是把结果写到 Redis,但维表关联又引入了新的网络往返。
这两个需求,本质上是在问两件事:
Fluss 的答案是:同一张表,同时给你 LogStore 和 KvStore 两套引擎。
这一篇我们就拆开看这两套引擎。
二、图 1:存储引擎类图

这张图分了两个包:蓝色是 LogStore(日志存储),米色是 KvStore(键值存储)。
2.1 LogStore:和 Kafka 同构的设计
先看类关系链:
LogManager *– LogTablet *– LogSegment *– OffsetIndex
*– TimeIndex
LogTablet ..> LogLoader (启动恢复)
逐个说明:
| LogManager | 管理一个 TabletServer 上所有 LogTablet 的生命周期 | Kafka 的 LogManager |
| LogTablet | 单个 Tablet 的日志,持有 ConcurrentSkipListMap<Long, LogSegment> | Kafka 的 Log(分区) |
| LogSegment | 单个日志段 = .log + .index + .timeindex | Kafka 的 LogSegment |
| OffsetIndex | 稀疏的 offset → 物理位置索引 | Kafka 的 OffsetIndex |
| TimeIndex | 稀疏的 timestamp → offset 索引 | Kafka 的 TimeIndex |
| LogLoader | 启动时扫描目录、恢复段、校验完整性 | Kafka 的 LogLoader |
熟悉 Kafka 的同学看到这张表应该会心一笑——Fluss 的 LogStore 在设计上几乎和 Kafka 的分区存储是一一对应的。这不是巧合,因为"追加写 + 稀疏索引 + 顺序读"已经被证明是日志存储的最优解,没必要重新发明。
LogSegment 的三个文件
class LogSegment {
private final long baseOffset; // 段的起始 offset
private final File logFile; // 00000000000000000000.log
private final OffsetIndex offsetIndex; // …index
private final TimeIndex timeIndex; // …timeindex
public LogAppendInfo append(LogRecordsBatch batch) { /* … */ }
public LogReadInfo read(long offset, int maxSize) { /* … */ }
public long size() { /* … */ }
}
磁盘上的样子:
bucket-0/
├── 00000000000000000000.log 1073741824 (1GB,滚动阈值)
├── 00000000000000000000.index 1048576 (1MB,固定大小)
├── 00000000000000000000.timeindex 524288
├── 00000000000000368944.log 843215678 (当前活跃段)
├── 00000000000000368944.index
└── 00000000000000368944.timeindex
文件名就是该段的起始 offset,这个设计让"给定一个 offset,找到它在哪个段"变成一次对 ConcurrentSkipListMap 的 floorEntry() 查询,O(log n)。
稀疏索引:为什么不建全量索引?
OffsetIndex 不是每条消息都记一条索引,而是每隔 N 字节记一条。假设 N = 4096:
索引项:(相对 offset, 物理位置)
(0, 0)
(57, 4096) ← 第 57 条消息在文件的第 4096 字节
(128, 8192)
(201, 12288)
…
查 offset=150 的流程:
abstract class AbstractIndex {
protected MappedByteBuffer mmap; // 内存映射,零拷贝
protected long baseOffset;
public abstract OffsetPosition lookup(long key);
// 内部:二分查找找到 <= 150 的最大索引项 (128, 8192)
// 然后从文件 8192 字节处开始顺序扫描,直到找到 150
}
权衡:全量索引查找是 O(1) 但索引文件会和数据一样大;稀疏索引查找是 O(log n) + 少量顺序扫描,但索引只有数据量的千分之一。对于日志这种"绝大多数情况都是顺序读"的访问模式,稀疏索引明显更优。
这就是图中 AbstractIndex 用 mmap 字段的原因——索引文件用 MappedByteBuffer 内存映射,读索引不产生系统调用。
2.2 KvStore:RocksDB 的封装
KvManager *– KvTablet *– RocksDBKv ..> RocksDBConfig
KvTablet ..> LogTablet (作为 WAL)
KvTablet ..> WalRecovery (故障恢复)
关键在 KvTablet 这个类:
class KvTablet {
private final RocksDBKv rocksDB; // 实际存储
private final LogTablet walLogTablet; // 指向同一 Tablet 的 LogTablet
public void put(byte[] key, byte[] value) { /* … */ }
public byte[] get(byte[] key) { /* … */ }
public void delete(byte[] key) { /* … */ }
public void recoverFromLog(LogTablet logTablet) { /* … */ }
}
注意 walLogTablet 这个字段——这是理解 PK 表的关键。
图中右下角那条最重要的关系:
LogTablet <.. KvTablet : "写入顺序:先 WAL 后 RocksDB"
为什么必须"先 WAL 后 RocksDB"?
因为 RocksDB 自身的 WAL 和 Fluss 的副本机制对不上。
假设我们只用 RocksDB(不写 Fluss 的 LogTablet),会发生什么:
1. Leader 写入 RocksDB 成功
2. Leader 宕机,还没来得及同步给 Follower
3. Follower 晋升为新 Leader
4. 新 Leader 的 RocksDB 里没有这条数据
5. Follower 重启后,如何追平?—— 没办法,因为 RocksDB 的 WAL 是本地的,
不参与副本复制,也没有统一的 offset 概念
所以 Fluss 的解法是:把 LogTablet 当作 KvTablet 的共享 WAL。
写入顺序:
① LogTablet.append(batch) → 产生全局递增 offset,参与 ISR 复制
② RocksDBKv.put(k, v) → 更新本地 LSM 树,提供点查能力
恢复顺序(节点重启):
① 读取 RocksDB 的持久化状态
② 从 RocksDB 记录的 lastWrittenOffset 开始,回放 LogTablet 中之后的记录
这个"回放"就是图中的 WalRecovery:
class WalRecovery {
public void recover(LogTablet logTablet, RocksDBKv kvStore, long fromOffset) {
// 从 fromOffset 扫描到 LEO
// 对每条记录重新执行 KvTablet.put()
// RocksDB 的写入是幂等的(同 key 覆盖),所以重复回放安全
}
}
这个设计和 MySQL 的 redo log / HBase 的 WAL 是同一个思想:用一份顺序写日志来保证持久性,用一份随机写结构来保证查询性能。
RocksDBConfig 的几个关键参数
class RocksDBConfig {
public static final long BLOCK_CACHE_SIZE = 256 * 1024 * 1024; // 256MB
public static final long WRITE_BUFFER_SIZE = 64 * 1024 * 1024; // 64MB
public static final double BLOOM_BITS_PER_KEY = 10.0;
public static final int NUM_LEVELS = 7;
}
| BLOCK_CACHE_SIZE | 256MB | 读缓存,缓存解压后的 SST block | 点查多、内存足 → 调到 1-2GB |
| WRITE_BUFFER_SIZE | 64MB | MemTable 大小,写满后 flush 成 L0 | 写入量大 → 调大,减少 flush 次数 |
| BLOOM_BITS_PER_KEY | 10.0 | 布隆过滤器精度 | 10 bits 的假阳性率约 1%,足够 |
| NUM_LEVELS | 7 | LSM 层数 | 一般不动 |
Bloom Filter 在这里的作用非常关键,它直接决定了"查一个不存在的 key"的代价:
没有 Bloom Filter:查 key → MemTable 未命中 → 遍历 L0 所有文件 → L1 → … → L6
最坏情况下要读 7 层的多个 SST 文件,几十次磁盘 IO
有 Bloom Filter: 查 key → MemTable 未命中 → Bloom 判断"一定不存在" → 直接返回 null
零次磁盘 IO
对于维表关联这种"大量 key 命不中"的场景,Bloom Filter 能把 P99 延迟压到亚毫秒级。下一节讲 Lookup 时序图时会再看到它。
2.3 回到场景:这套设计解决了什么?
回到开头的风控场景,如果改用 Fluss PK 表:
CREATE TABLE user_behavior (
user_id BIGINT,
event_type STRING,
amount DECIMAL(18, 2),
device_id STRING,
event_time TIMESTAMP(3),
— … 还有 195 个字段
PRIMARY KEY (user_id, event_time) NOT ENFORCED
) WITH (
'bucket.num' = '64',
'table.merge-engine' = 'deduplicate'
);
需求 1(只读 5 列):靠 LogStore 的 Arrow 列式格式 + 列裁剪解决(第 3 篇详述)。
需求 2(查累计交易额):靠 KvStore 的点查解决:
— 直接点查,不用再往 Redis 写一份
SELECT balance FROM user_profile WHERE user_id = 10086;
三、图 2:Merge Engine 类图

有了 KvStore,就带来一个新问题:同一个 key 被写入多次时,应该保留哪个值?
这就是 Merge Engine 要回答的。图中 RowMerger 接口有三个实现:
interface RowMerger {
RowData merge(RowData oldRow, RowData newRow);
}
3.1 Deduplicate(默认):新值覆盖旧值
CREATE TABLE user_profile (
user_id BIGINT,
name STRING,
age INT,
city STRING,
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'bucket.num' = '16',
'table.merge-engine' = 'deduplicate' — 默认值,可省略
);
行为:
class DeduplicateRowMerger implements RowMerger {
public RowData merge(RowData oldRow, RowData newRow) {
return newRow; // 就这么简单
}
}
| INSERT (1, 'Alice', 28, 'Beijing') | (1, 'Alice', 28, 'Beijing') |
| INSERT (1, 'Alice', 29, 'Shanghai') | (1, 'Alice', 29, 'Shanghai') 全行覆盖 |
适用场景:维表、CDC 同步的目标表(源端发来的是完整的行镜像)。
3.2 PartialUpdate:NULL 不覆盖
这是 Fluss 最有实用价值的特性之一。
CREATE TABLE user_tags (
user_id BIGINT,
tag1 STRING,
tag2 STRING,
tag3 STRING,
PRIMARY KEY (user_id) NOT ENFORCED
) WITH ('table.merge-engine' = 'partial-update');
— 三个不同的业务系统各自写入自己负责的列
INSERT INTO user_tags VALUES (1, 'vip', NULL, NULL); — 会员系统
INSERT INTO user_tags VALUES (1, NULL, 'tech', NULL); — 内容系统
INSERT INTO user_tags VALUES (1, NULL, NULL, 'gamer'); — 游戏系统
— 最终结果:(1, 'vip', 'tech', 'gamer')
class PartialUpdateRowMerger implements RowMerger {
public RowData merge(RowData oldRow, RowData newRow) {
// 逐列:newRow 该列为 NULL 则保留 oldRow 的值,否则用新值
}
}
这个特性解决了一个经典的工程难题:宽表的多源拼接。
传统做法要用 Flink 做多流 Join:
— 传统方案:三流 Join,状态巨大
SELECT a.user_id, a.tag1, b.tag2, c.tag3
FROM member_stream a
LEFT JOIN content_stream b ON a.user_id = b.user_id
LEFT JOIN game_stream c ON a.user_id = c.user_id;
这个 Join 的状态量是三个流的状态之和,而且任何一条流迟到都会触发回撤。
用 Partial Update 之后,三个系统各自独立写入同一张表的不同列,完全不需要 Join:
会员系统 → INSERT (user_id, tag1, NULL, NULL) ─┐
内容系统 → INSERT (user_id, NULL, tag2, NULL) ─┼→ Fluss 自动合并
游戏系统 → INSERT (user_id, NULL, NULL, tag3) ─┘
注意事项:
3.3 Aggregation:把聚合下推到存储层
CREATE TABLE user_stats (
user_id BIGINT,
total_spent DECIMAL(18, 2),
order_count INT,
last_login TIMESTAMP(3),
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'table.merge-engine' = 'aggregation',
'fields.total_spent.aggregate-function' = 'sum',
'fields.order_count.aggregate-function' = 'sum',
'fields.last_login.aggregate-function' = 'last_value'
);
写入效果:
INSERT INTO user_stats VALUES (1, 100.00, 1, TIMESTAMP '2026-08-01 10:00:00');
INSERT INTO user_stats VALUES (1, 50.00, 1, TIMESTAMP '2026-08-01 11:30:00');
INSERT INTO user_stats VALUES (1, 200.00, 1, TIMESTAMP '2026-08-01 12:00:00');
— 表中最终只有一行:(1, 350.00, 3, 2026-08-01 12:00:00)
图中 AggregationRowMerger 持有 Map<String, AggregateFunction>:
class AggregationRowMerger implements RowMerger {
private final Map<String, AggregateFunction> functions; // 每列可配置
public RowData merge(RowData oldRow, RowData newRow) {
// 遍历每一列:
// 配置了聚合函数 → 执行 functions.get(col).aggregate(oldVal, newVal)
// 未配置 → 新值覆盖(退化成 Deduplicate)
}
}
五种 AggregateFunction 实现:
| SumAggregate | 累加 | 累计金额、累计次数 |
| MaxAggregate | 取最大 | 最高单价、最近更新时间 |
| MinAggregate | 取最小 | 最低价 |
| CountAggregate | 计数 | 行为次数 |
| LastValueAggregate | 取最新值 | 最后登录时间、最后状态 |
这个特性的价值:把聚合状态从 Flink 下推到存储层。
传统做法:
— Flink 里做聚合,状态存 RocksDB StateBackend
SELECT user_id, SUM(amount), COUNT(*)
FROM orders GROUP BY user_id;
— 问题:状态随 user_id 数量线性增长,亿级用户 = TB 级状态
— Checkpoint 慢,恢复慢
用 Aggregation Merge Engine:
— Flink 只需要做简单的转发
INSERT INTO user_stats SELECT user_id, amount, 1, event_time FROM orders;
— 聚合在 Fluss 服务端完成,Flink 作业无状态
这样一来,Flink 作业变成了无状态的纯转发,Checkpoint 秒级完成,扩缩容秒级生效。这就是 Fluss 宣传的"状态外部化"的一部分。
3.4 Merge Engine 选型决策树
需要按主键查询/更新吗?
│
├─ 否 → Log 表(无需 Merge Engine,只有 LogStore)
│
└─ 是 → PK 表,选哪个 Merge Engine?
│
├─ 每次写入都是完整行镜像(如 CDC)
│ └─ deduplicate(默认)
│
├─ 多个数据源各自写不同的列
│ └─ partial-update
│
└─ 需要累加/取最大最小/计数
└─ aggregation
3.5 三种引擎的代价对比
| deduplicate | 1x | 最优 | 1x | 全行覆盖 |
| partial-update | 1x(需先读旧值) | 优 | 1x | 多源列更新 |
| aggregation | 1x(需先读旧值) | 优 | 1x | 只增不减的累加 |
注意:partial-update 和 aggregation 在合并时需要先读出旧值,所以它们的写入路径比 deduplicate 多一次 RocksDB 读。如果旧值在 block cache 里,这次读是微秒级;如果不在,会有一次磁盘 IO。调优方向就是把 BLOCK_CACHE_SIZE 调大,让热 key 常驻内存。
四、动手验证
4.1 观察 LogSegment 的滚动
# 创建一个段大小很小的表(生产不要这么干,这里是为了观察)
CREATE TABLE seg_test (
id BIGINT, payload STRING
) WITH ('bucket.num' = '1', 'log.segment.size' = '1MB');
— 持续写入,观察文件滚动
watch -n 1 'ls -la /data/fluss-server/data/log/fluss/seg_test/bucket-0/'
你会看到 00000000000000000000.log 写满 1MB 后,出现 00000000000000012345.log,文件名就是新段的起始 offset。
4.2 验证 Partial Update 的合并行为
— Flink SQL 客户端
CREATE TABLE user_tags (
user_id BIGINT, tag1 STRING, tag2 STRING, tag3 STRING,
PRIMARY KEY (user_id) NOT ENFORCED
) WITH ('table.merge-engine' = 'partial-update');
INSERT INTO user_tags VALUES (1, 'vip', NULL, NULL);
INSERT INTO user_tags VALUES (1, NULL, 'tech', NULL);
INSERT INTO user_tags VALUES (1, NULL, NULL, 'gamer');
SELECT * FROM user_tags WHERE user_id = 1;
— 预期:1, vip, tech, gamer
4.3 验证 Aggregation 的累加
CREATE TABLE user_stats (
user_id BIGINT, total_spent DECIMAL(18,2), order_count INT,
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'table.merge-engine' = 'aggregation',
'fields.total_spent.aggregate-function' = 'sum',
'fields.order_count.aggregate-function' = 'sum'
);
INSERT INTO user_stats VALUES (1, 100.00, 1);
INSERT INTO user_stats VALUES (1, 50.00, 1);
INSERT INTO user_stats VALUES (1, 200.00, 1);
SELECT * FROM user_stats WHERE user_id = 1;
— 预期:1, 350.00, 3
五、生产实践要点
5.1 Bucket 数量怎么定
bucket.num 是 PK 表最重要的参数,因为它决定了并发度和数据分布的均匀性,且后期调整代价大。
— 估算公式
bucket.num ≈ max(预期峰值写入 QPS / 单 bucket 吞吐, TabletServer 数 × 2~4)
经验值:
| 小维表(百万行以内) | 4 – 16 |
| 中型表(千万行) | 32 – 64 |
| 大表(亿级以上) | 128 – 512 |
注意:bucket.num 只能在建表时指定。要修改只能重建表迁移数据。所以宁可一开始就给大一点。
5.2 RocksDB 调优
默认参数对中小规模够用,但如果点查 QPS 很高(> 10万/s),需要调:
# tablet-server.yaml
# 增大 block cache(但要给 Flink TM 留内存)
rocksdb.block.cache.size: 1GB
# 增大 MemTable,减少 flush(写多读少场景)
rocksdb.write.buffer.size: 128MB
# 限制 L0 文件数,避免读放大
rocksdb.level0.slowdown-writes-trigger: 20
rocksdb.level0.stop-writes-trigger: 36
5.3 磁盘选型
推荐:NVMe SSD
避免:网络盘(云盘)、HDD
原因:RocksDB 的 compaction 是随机 IO 密集型,网络盘的 IOPS 和延迟都不达标。如果只能用云盘,选择 ESSD PL1 以上规格。
5.4 容量规划
PK 表单副本占用 ≈ LogStore 数据 + KvStore 数据
≈ 原始数据量 × (1 + 1.5) # RocksDB 有空间放大
≈ 原始数据量 × 2.5
再乘以副本数(默认 3):总占用 ≈ 原始数据量 × 7.5
这也是为什么 Log 表比 PK 表"便宜"很多——Log 表只有 1x 数据 × 3 副本 = 3x。
六、排障手册
| 点查 P99 突然从 1ms 涨到 50ms | block cache 命中率下降 | 监控 RocksDB block-cache-hit-rate;检查是否有大批量扫描打乱了缓存 |
| 写入吞吐上不去 | L0 文件堆积触发 write stall | 看日志有没有 “Stopping writes”;调大 write buffer 或加快 compaction |
| Partial Update 结果不符合预期 | 两个源写了同一列 | 检查各写入任务的列映射;确认没有把需要保留的列写成 NULL |
| Aggregation 数值偏小 | 部分写入失败被静默丢弃 | 检查 Flink 作业是否有反压导致的超时重试 |
| 磁盘空间持续增长不释放 | LogSegment 未过期 | 检查 log.retention 配置;确认 PK 表的 KvTablet 有在做 snapshot |
| 重启后恢复特别慢 | WAL 回放量大 | 增大 KvTablet 的 snapshot 频率,减小需要回放的 offset 区间 |
七、小结
下一篇,我们把这两套引擎串起来,跟踪一次写入、一次流式读取、一次点查的完整调用链。
网硕互联帮助中心

评论前必须登录!
注册