云计算百科
云计算领域专业知识百科平台

图解 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 个字段:

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,但维表关联又引入了新的网络往返。

这两个需求,本质上是在问两件事:

  • 能不能只读需要的列?(列式存储 + 列裁剪)
  • 能不能按主键直接查?(KV 存储 + 点查)
  • 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) ─┘

    注意事项:

  • 如果业务上确实需要"把某列更新成 NULL",Partial Update 做不到——你需要显式写入一个哨兵值(如空字符串、-1)。
  • 多个源同时写同一列时,最后一次写入生效(顺序由 LogTablet 的 offset 决定)。
  • 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)

    经验值:

    场景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: 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 区间

    七、小结

  • 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 状态。
  • 下一篇,我们把这两套引擎串起来,跟踪一次写入、一次流式读取、一次点查的完整调用链。


    赞(0)
    未经允许不得转载:网硕互联帮助中心 » 图解 Fluss(二):存储引擎双引擎 —— LogStore 与 KvStore
    分享到: 更多 (0)

    评论 抢沙发

    评论前必须登录!