图解 Fluss(二):存储引擎双引擎 —— LogStore 与 KvStore

📅 发布时间:2026/8/31 4:46:34
图解 Fluss(二):存储引擎双引擎 —— LogStore 与 KvStore 图解 Fluss二存储引擎双引擎 —— LogStore 与 KvStore阅读本文你将了解LogSegment 的稀疏索引为什么和 Kafka 长得一样、PK 表先写 WAL 再写 RocksDB的顺序为什么不能反、RocksDB 那几个关键参数的含义、以及三种 Merge Engine 该怎么选。配套图表class-02-storage-engine、class-04-merge-engine难度⭐⭐⭐ | 适合人群要做表设计与性能调优的开发一、场景一次加个字段引发的事故先看一个真实场景。某风控团队有一张用户行为表200 多个字段日均 20 亿条。他们用 Flink 做实时特征计算实际只用到其中 5 个字段SELECTuser_id,event_type,amount,device_id,event_timeFROMuser_behaviorWHEREevent_typePAY;但底层是 KafkaFlink 的 Source 必须把整行 200 个字段全部拉下来再在内存里做投影和过滤。结果网络带宽高峰期 12 Gbps其中 97% 是无效字段Flink TM 内存经常 OOM因为反序列化后的完整对象太大端到端延迟P99 达到 3 秒后来他们加了第二个需求风控规则要查用户最近 30 天的累计交易额。原来的做法是把结果写到 Redis但维表关联又引入了新的网络往返。这两个需求本质上是在问两件事能不能只读需要的列列式存储 列裁剪能不能按主键直接查KV 存储 点查Fluss 的答案是同一张表同时给你 LogStore 和 KvStore 两套引擎。这一篇我们就拆开看这两套引擎。二、图 1存储引擎类图这张图分了两个包蓝色是 LogStore日志存储米色是 KvStore键值存储。2.1 LogStore和 Kafka 同构的设计先看类关系链LogManager *-- LogTablet *-- LogSegment *-- OffsetIndex *-- TimeIndex LogTablet .. LogLoader (启动恢复)逐个说明类职责类比LogManager管理一个 TabletServer 上所有 LogTablet 的生命周期Kafka 的LogManagerLogTablet单个 Tablet 的日志持有ConcurrentSkipListMapLong, LogSegmentKafka 的Log分区LogSegment单个日志段 .log.index.timeindexKafka 的LogSegmentOffsetIndex稀疏的 offset → 物理位置索引Kafka 的OffsetIndexTimeIndex稀疏的 timestamp → offset 索引Kafka 的TimeIndexLogLoader启动时扫描目录、恢复段、校验完整性Kafka 的LogLoader熟悉 Kafka 的同学看到这张表应该会心一笑——Fluss 的 LogStore 在设计上几乎和 Kafka 的分区存储是一一对应的。这不是巧合因为追加写 稀疏索引 顺序读已经被证明是日志存储的最优解没必要重新发明。LogSegment 的三个文件classLogSegment{privatefinallongbaseOffset;// 段的起始 offsetprivatefinalFilelogFile;// 00000000000000000000.logprivatefinalOffsetIndexoffsetIndex;// ...indexprivatefinalTimeIndextimeIndex;// ...timeindexpublicLogAppendInfoappend(LogRecordsBatchbatch){/* ... */}publicLogReadInforead(longoffset,intmaxSize){/* ... */}publiclongsize(){/* ... */}}磁盘上的样子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) ...查 offset150 的流程abstractclassAbstractIndex{protectedMappedByteBuffermmap;// 内存映射零拷贝protectedlongbaseOffset;publicabstractOffsetPositionlookup(longkey);// 内部二分查找找到 150 的最大索引项 (128, 8192)// 然后从文件 8192 字节处开始顺序扫描直到找到 150}权衡全量索引查找是 O(1) 但索引文件会和数据一样大稀疏索引查找是 O(log n) 少量顺序扫描但索引只有数据量的千分之一。对于日志这种绝大多数情况都是顺序读的访问模式稀疏索引明显更优。这就是图中AbstractIndex用mmap字段的原因——索引文件用MappedByteBuffer内存映射读索引不产生系统调用。2.2 KvStoreRocksDB 的封装KvManager *-- KvTablet *-- RocksDBKv .. RocksDBConfig KvTablet .. LogTablet (作为 WAL) KvTablet .. WalRecovery (故障恢复)关键在KvTablet这个类classKvTablet{privatefinalRocksDBKvrocksDB;// 实际存储privatefinalLogTabletwalLogTablet;// 指向同一 Tablet 的 LogTabletpublicvoidput(byte[]key,byte[]value){/* ... */}publicbyte[]get(byte[]key){/* ... */}publicvoiddelete(byte[]key){/* ... */}publicvoidrecoverFromLog(LogTabletlogTablet){/* ... */}}注意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 中之后的记录这个回放就是图中的WalRecoveryclassWalRecovery{publicvoidrecover(LogTabletlogTablet,RocksDBKvkvStore,longfromOffset){// 从 fromOffset 扫描到 LEO// 对每条记录重新执行 KvTablet.put()// RocksDB 的写入是幂等的同 key 覆盖所以重复回放安全}}这个设计和 MySQL 的 redo log / HBase 的 WAL 是同一个思想用一份顺序写日志来保证持久性用一份随机写结构来保证查询性能。RocksDBConfig 的几个关键参数classRocksDBConfig{publicstaticfinallongBLOCK_CACHE_SIZE256*1024*1024;// 256MBpublicstaticfinallongWRITE_BUFFER_SIZE64*1024*1024;// 64MBpublicstaticfinaldoubleBLOOM_BITS_PER_KEY10.0;publicstaticfinalintNUM_LEVELS7;}参数值作用调优建议BLOCK_CACHE_SIZE256MB读缓存缓存解压后的 SST block点查多、内存足 → 调到 1-2GBWRITE_BUFFER_SIZE64MBMemTable 大小写满后 flush 成 L0写入量大 → 调大减少 flush 次数BLOOM_BITS_PER_KEY10.0布隆过滤器精度10 bits 的假阳性率约 1%足够NUM_LEVELS7LSM 层数一般不动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 表CREATETABLEuser_behavior(user_idBIGINT,event_type STRING,amountDECIMAL(18,2),device_id STRING,event_timeTIMESTAMP(3),-- ... 还有 195 个字段PRIMARYKEY(user_id,event_time)NOTENFORCED)WITH(bucket.num64,table.merge-enginededuplicate);需求 1只读 5 列靠 LogStore 的 Arrow 列式格式 列裁剪解决第 3 篇详述。需求 2查累计交易额靠 KvStore 的点查解决-- 直接点查不用再往 Redis 写一份SELECTbalanceFROMuser_profileWHEREuser_id10086;三、图 2Merge Engine 类图有了 KvStore就带来一个新问题同一个 key 被写入多次时应该保留哪个值这就是 Merge Engine 要回答的。图中RowMerger接口有三个实现interfaceRowMerger{RowDatamerge(RowDataoldRow,RowDatanewRow);}3.1 Deduplicate默认新值覆盖旧值CREATETABLEuser_profile(user_idBIGINT,name STRING,ageINT,city STRING,PRIMARYKEY(user_id)NOTENFORCED)WITH(bucket.num16,table.merge-enginededuplicate-- 默认值可省略);行为classDeduplicateRowMergerimplementsRowMerger{publicRowDatamerge(RowDataoldRow,RowDatanewRow){returnnewRow;// 就这么简单}}操作结果INSERT (1, Alice, 28, Beijing)(1, Alice, 28, Beijing)INSERT (1, Alice, 29, Shanghai)(1, Alice, 29, Shanghai)全行覆盖适用场景维表、CDC 同步的目标表源端发来的是完整的行镜像。3.2 PartialUpdateNULL 不覆盖这是 Fluss 最有实用价值的特性之一。CREATETABLEuser_tags(user_idBIGINT,tag1 STRING,tag2 STRING,tag3 STRING,PRIMARYKEY(user_id)NOTENFORCED)WITH(table.merge-enginepartial-update);-- 三个不同的业务系统各自写入自己负责的列INSERTINTOuser_tagsVALUES(1,vip,NULL,NULL);-- 会员系统INSERTINTOuser_tagsVALUES(1,NULL,tech,NULL);-- 内容系统INSERTINTOuser_tagsVALUES(1,NULL,NULL,gamer);-- 游戏系统-- 最终结果(1, vip, tech, gamer)classPartialUpdateRowMergerimplementsRowMerger{publicRowDatamerge(RowDataoldRow,RowDatanewRow){// 逐列newRow 该列为 NULL 则保留 oldRow 的值否则用新值}}这个特性解决了一个经典的工程难题宽表的多源拼接。传统做法要用 Flink 做多流 Join-- 传统方案三流 Join状态巨大SELECTa.user_id,a.tag1,b.tag2,c.tag3FROMmember_stream aLEFTJOINcontent_stream bONa.user_idb.user_idLEFTJOINgame_stream cONa.user_idc.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) ─┘注意事项如果业务上确实需要把某列更新成 NULLPartial Update 做不到——你需要显式写入一个哨兵值如空字符串、-1。多个源同时写同一列时最后一次写入生效顺序由 LogTablet 的 offset 决定。3.3 Aggregation把聚合下推到存储层CREATETABLEuser_stats(user_idBIGINT,total_spentDECIMAL(18,2),order_countINT,last_loginTIMESTAMP(3),PRIMARYKEY(user_id)NOTENFORCED)WITH(table.merge-engineaggregation,fields.total_spent.aggregate-functionsum,fields.order_count.aggregate-functionsum,fields.last_login.aggregate-functionlast_value);写入效果INSERTINTOuser_statsVALUES(1,100.00,1,TIMESTAMP2026-08-01 10:00:00);INSERTINTOuser_statsVALUES(1,50.00,1,TIMESTAMP2026-08-01 11:30:00);INSERTINTOuser_statsVALUES(1,200.00,1,TIMESTAMP2026-08-01 12:00:00);-- 表中最终只有一行(1, 350.00, 3, 2026-08-01 12:00:00)图中AggregationRowMerger持有MapString, AggregateFunctionclassAggregationRowMergerimplementsRowMerger{privatefinalMapString,AggregateFunctionfunctions;// 每列可配置publicRowDatamerge(RowDataoldRow,RowDatanewRow){// 遍历每一列// 配置了聚合函数 → 执行 functions.get(col).aggregate(oldVal, newVal)// 未配置 → 新值覆盖退化成 Deduplicate}}五种AggregateFunction实现函数语义典型用途SumAggregate累加累计金额、累计次数MaxAggregate取最大最高单价、最近更新时间MinAggregate取最小最低价CountAggregate计数行为次数LastValueAggregate取最新值最后登录时间、最后状态这个特性的价值把聚合状态从 Flink 下推到存储层。传统做法-- Flink 里做聚合状态存 RocksDB StateBackendSELECTuser_id,SUM(amount),COUNT(*)FROMordersGROUPBYuser_id;-- 问题状态随 user_id 数量线性增长亿级用户 TB 级状态-- Checkpoint 慢恢复慢用 Aggregation Merge Engine-- Flink 只需要做简单的转发INSERTINTOuser_statsSELECTuser_id,amount,1,event_timeFROMorders;-- 聚合在 Fluss 服务端完成Flink 作业无状态这样一来Flink 作业变成了无状态的纯转发Checkpoint 秒级完成扩缩容秒级生效。这就是 Fluss 宣传的状态外部化的一部分。3.4 Merge Engine 选型决策树需要按主键查询/更新吗 │ ├─ 否 → Log 表无需 Merge Engine只有 LogStore │ └─ 是 → PK 表选哪个 Merge Engine │ ├─ 每次写入都是完整行镜像如 CDC │ └─ deduplicate默认 │ ├─ 多个数据源各自写不同的列 │ └─ partial-update │ └─ 需要累加/取最大最小/计数 └─ aggregation3.5 三种引擎的代价对比写放大读性能存储开销适用写模式deduplicate1x最优1x全行覆盖partial-update1x需先读旧值优1x多源列更新aggregation1x需先读旧值优1x只增不减的累加注意partial-update和aggregation在合并时需要先读出旧值所以它们的写入路径比deduplicate多一次 RocksDB 读。如果旧值在 block cache 里这次读是微秒级如果不在会有一次磁盘 IO。调优方向就是把BLOCK_CACHE_SIZE调大让热 key 常驻内存。四、动手验证4.1 观察 LogSegment 的滚动# 创建一个段大小很小的表生产不要这么干这里是为了观察CREATE TABLE seg_test(idBIGINT, payload STRING)WITH(bucket.num1,log.segment.size1MB);-- 持续写入观察文件滚动watch-n1ls -la /data/fluss-server/data/log/fluss/seg_test/bucket-0/你会看到00000000000000000000.log写满 1MB 后出现00000000000000012345.log文件名就是新段的起始 offset。4.2 验证 Partial Update 的合并行为-- Flink SQL 客户端CREATETABLEuser_tags(user_idBIGINT,tag1 STRING,tag2 STRING,tag3 STRING,PRIMARYKEY(user_id)NOTENFORCED)WITH(table.merge-enginepartial-update);INSERTINTOuser_tagsVALUES(1,vip,NULL,NULL);INSERTINTOuser_tagsVALUES(1,NULL,tech,NULL);INSERTINTOuser_tagsVALUES(1,NULL,NULL,gamer);SELECT*FROMuser_tagsWHEREuser_id1;-- 预期1, vip, tech, gamer4.3 验证 Aggregation 的累加CREATETABLEuser_stats(user_idBIGINT,total_spentDECIMAL(18,2),order_countINT,PRIMARYKEY(user_id)NOTENFORCED)WITH(table.merge-engineaggregation,fields.total_spent.aggregate-functionsum,fields.order_count.aggregate-functionsum);INSERTINTOuser_statsVALUES(1,100.00,1);INSERTINTOuser_statsVALUES(1,50.00,1);INSERTINTOuser_statsVALUES(1,200.00,1);SELECT*FROMuser_statsWHEREuser_id1;-- 预期1, 350.00, 3五、生产实践要点5.1 Bucket 数量怎么定bucket.num是 PK 表最重要的参数因为它决定了并发度和数据分布的均匀性且后期调整代价大。-- 估算公式bucket.num ≈max(预期峰值写入 QPS/单 bucket 吞吐,TabletServer 数 ×2~4)经验值场景bucket.num 建议小维表百万行以内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: 365.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 涨到 50msblock cache 命中率下降监控 RocksDBblock-cache-hit-rate检查是否有大批量扫描打乱了缓存写入吞吐上不去L0 文件堆积触发 write stall看日志有没有 “Stopping writes”调大 write buffer 或加快 compactionPartial Update 结果不符合预期两个源写了同一列检查各写入任务的列映射确认没有把需要保留的列写成 NULLAggregation 数值偏小部分写入失败被静默丢弃检查 Flink 作业是否有反压导致的超时重试磁盘空间持续增长不释放LogSegment 未过期检查log.retention配置确认 PK 表的 KvTablet 有在做 snapshot重启后恢复特别慢WAL 回放量大增大 KvTablet 的 snapshot 频率减小需要回放的 offset 区间七、小结LogStore 与 Kafka 同构LogSegment.log.index.timeindex稀疏索引 mmap 二分查找这是日志存储的最优解。KvStore 用 LogTablet 做共享 WAL写入顺序必须是先 WAL 后 RocksDB原因是 RocksDB 本地 WAL 不参与副本复制。Bloom Filter 是点查性能的关键它把查不存在的 key的代价从几十次磁盘 IO 降到零。三种 Merge Engine 解决三类问题deduplicate全行覆盖、partial-update多源列拼接、aggregation聚合下推。partial-update的工程价值最大它让多流 Join 拼宽表变成了多源独立写入 存储层合并彻底消除了 Join 状态。下一篇我们把这两套引擎串起来跟踪一次写入、一次流式读取、一次点查的完整调用链。