图解 Fluss三读写全链路 —— 一次写入、一次扫描、一次点查阅读本文你将了解一条记录从客户端到落盘要经过哪些步骤、Exactly-Once 是怎么保证的、列裁剪为什么能把网络传输降两个数量级、以及 Lookup Join 的三级缓存如何让维表关联做到亚毫秒级。配套图表seq-01-write-path、seq-02-read-path、seq-03-pk-lookup难度⭐⭐⭐ | 适合人群要做性能调优与 Exactly-Once 方案设计的开发一、场景一个实时大屏的三个请求某电商平台要做一个实时大屏每秒处理 50 万条订单事件。这个大屏背后有三个完全不同的数据请求请求 A订单写入INSERTINTOordersSELECT*FROMkafka_source;-- 50 万 QPS要求不能丢、不能重复。这是写入路径。请求 B大屏聚合SELECTSUM(amount),COUNT(*)FROMordersWHEREevent_timeNOW()-INTERVAL5MINUTE;orders表有 180 个字段但这个查询只用 2 个。这是扫描路径要求高吞吐、低延迟。请求 C订单详情关联用户SELECTo.order_id,u.user_name,u.cityFROMorders oLEFTJOINuser_profileFORSYSTEM_TIMEASOFo.proc_timeASuONo.user_idu.user_id;每条订单都要查一次用户表QPS 和订单写入一样高。这是点查路径要求极低延迟。三条路径三种完全不同的优化思路。这一篇我们逐条跟踪。二、图 1PK 表写入路径整条链路分三个阶段图里用分隔线标得很清楚。2.1 阶段一元数据定位Client - Coordinator : GetTableMetadata(tablePath) Coord -- Client : TableDescriptor Bucket 映射 Client - Client : 计算 bucketId hash(pk) % bucket.num Client - Client : 定位 bucket 的 Leader TabletServer注意这个阶段只在初始化时执行一次不是每条记录都问一次 Coordinator。客户端拿到 TableDescriptor 后会缓存// 客户端侧的逻辑简化示意// TableBucket 的真实路径org.apache.fluss.metadata.TableBucketpublicvoidwrite(RowDatarow){// 1. 从本地缓存取表元数据首次或失效时才请求 CoordinatorTableDescriptordescriptormetadataCache.get(tablePath);// 2. 计算 bucket本地计算零网络开销intbucketIdBucketingFunction.hash(descriptor.getBucketKeys(),descriptor.getBucketNum(),row);// 等价于MurmurHash3.hash(primaryKey) % bucket.num// 3. 从缓存的路由表找 Leader首次或 Leader 变更时才请求 CoordinatorServerNodeleaderroutingCache.get(newTableBucket(tableId,bucketId));// 4. 直连 TabletServer 写入client.send(newWriteRecordRequest(leader,bucketId,batch,acks));}这个设计是 Fluss 高吞吐的前提之一Coordinator 不在数据路径上。如果每条记录都要问一次 Coordinator我该写到哪50 万 QPS 会把 Coordinator 打爆。三个关键缓存缓存内容失效时机表元数据缓存Schema、bucket.num、分区信息表结构变更Coordinator 主动推送路由缓存bucketId → Leader TabletServer 地址Leader 变更请求返回NOT_LEADER时刷新连接池到各 TabletServer 的 TCP 长连接连接断开2.2 阶段二写入先 WAL 后 RocksDB图中两个group块标出了两次写入// TabletServer 侧处理 WriteRecordRequest简化publicCompletableFutureWriteRecordResponsehandleWrite(WriteRecordRequestrequest){intbucketIdrequest.getBucketId();// ---- 写入 LogTablet (WAL) ----LogTabletlogTabletlogManager.getOrCreateLogTablet(tablePath,bucketId);LogAppendInfoappendInfologTablet.append(request.getBatch());// 内部检查是否需要 Roll Segment段写满 1GB 或超过时间阈值// 返回新记录的起始 offset// ---- 写入 KvTablet (RocksDB) ----if(tableDescriptor.hasPrimaryKey()){KvTabletkvTabletkvManager.getOrCreateKvTablet(tablePath,bucketId);for(Recordrecord:request.getBatch()){kvTablet.put(record.getKey(),record.getValue());// 内部先写 RocksDB 自身 WAL再写 MemTable// 如果配置了 merge-engine这里会先读旧值再合并}}// ---- 等待副本确认 ----returnwaitForAcks(appendInfo.getBaseOffset()appendInfo.getCount()-1,request.getAcks());}图右侧的 note 总结了顺序保证写入顺序保证: 1. 先写 WAL (LogTablet) 保证持久性 2. 再写 RocksDB (KvTablet) 提供查询能力 3. acksall 时等待所有 ISR 确认后才返回为什么这个顺序不能反上一篇已经答过LogTablet 是参与副本复制的共享 WALRocksDB 的本地 WAL 不是。反过来的话Follower 就没有办法追平数据。2.3 阶段三副本同步与 High WatermarkLeader - Follower : Replicate(batch, offset) Follower - Follower : 追加到本地 LogTablet Follower -- Leader : ACK(offset) Leader - Leader : 等待所有 ISR 确认 Leader - Leader : advanceHighWatermark(offset) Leader -- Client : 写入成功ISRIn-Sync Replicas是能够跟上 Leader 的副本集合。High WatermarkHW的定义是HW min(所有 ISR 副本的 LEO)其中 LEOLog End Offset是每个副本的日志末端 offset。举例说明 HW 的作用LEO HW Leader: [0 1 2 3 4 5 6] → 7 4 Follower1:[0 1 2 3 4 5] → 6 - Follower2:[0 1 2 3 4] → 5 - HW min(7, 6, 5) 5HW 的意义offset HW 的记录已经至少被所有 ISR 副本确认消费者一定能读到它即使 Leader 立刻宕机也不会丢。这就是acks参数的语义acks行为持久性延迟acks0写入 Leader 内存即返回最低可能丢最低acks1写入 Leader 的 LogTablet 即返回Leader 宕机可能丢低acksall等待所有 ISR 确认最高较高生产环境推荐acksallmin.insync.replicas2# 客户端配置 client.request.acks: all # 服务端配置ISR 最少要有 2 个副本否则拒绝写入 tablet-server.min.insync.replicas: 2这个组合很重要如果只配acksall而不配min.insync.replicas当 ISR 只剩 Leader 一个时acksall会退化成acks1——你以为数据有三副本实际上只有一份。2.4 Exactly-Once 是怎么来的图里最后一行写着写入成功 (Exactly-Once)。这需要两层机制叠加第一层Flink Sink 的两阶段提交classFlussSinkCommitterimplementsSinkCommitter{// Checkpoint 触发时调用预提交数据已写入但对外不可见publicListCommitRequestprepareCommit(){/* ... */}// 所有并行 writer 都预提交成功后调用真正提交数据对外可见publicvoidcommit(ListCommitRequestcommits){/* ... */}}第二层服务端的幂等写入classWriterStateManager{// 每个 Writer (producer) 有唯一的 writerId 和递增的 sequence number// 服务端记录每个 writerId 的最后 sequence// - seq lastSeq 1 → 正常接收// - seq lastSeq → 重复写入直接返回成功幂等// - seq lastSeq 1 → 数据丢失返回错误}两层叠加的效果场景Flink Checkpoint 失败后从上次成功的 checkpoint 恢复重放一批数据 传统 Kafka Sink这批数据会被重复写入 → 下游看到重复 Fluss Sink 两阶段提交保证预提交的数据没有 commit 重放时 sequence number 相同服务端识别为重复写入并丢弃 → 下游只看到一次配置方式-- Flink SQLSETexecution.checkpointing.modeEXACTLY_ONCE;SETexecution.checkpointing.interval30s;CREATETABLEorders_sink(...)WITH(connectorfluss,bootstrap.serverscoordinator:9123,sink.delivery-guaranteeexactly-once-- 默认值);2.5 请求 A 的答案回到开头的订单写入 50 万 QPS优化点效果元数据本地缓存Coordinator 零压力不在数据路径上bucket 路由本地计算无额外网络往返顺序写 WAL磁盘顺序 IO单盘可到 500MB/s批量追加攒批减少 RPC 次数acksallmin.insync.replicas2不丢数据调优参数# 客户端批量参数 client.writer.batch.size: 1MB # 攒批大小 client.writer.linger.ms: 10 # 最多等 10ms 攒批 client.writer.buffer.memory: 64MB # 客户端缓冲区三、图 2流式读取路径3.1 关键FetchRequest 带着投影列和过滤条件注意图里第一条消息Client - Server : FetchRequest(tabletId, offset, projectedColumns, predicates)这是 Fluss 和 Kafka 最本质的区别之一。Kafka 的 FetchRequest 只有topic partition offset服务端只能把完整的字节流吐给你而 Fluss 的请求里带了projectedColumns我只要这几列predicates我的 WHERE 条件服务端因此可以做两件优化图中用两个group块标出。3.2 优化一列裁剪Column ProjectionLog - Arrow : 读取 Arrow RecordBatch (全列) Arrow - Arrow : 通过 TransferPair 零拷贝投影只保留 projectedColumns Arrow -- Log : 投影后的 RecordBatch (仅需要的列)为什么能零拷贝因为底层数据是 Arrow 列式格式。Arrow 的内存布局是按列连续存储的行式Kafka / 传统 row1: [id|name|age|city|amount|...] row2: [id|name|age|city|amount|...] → 要读第 2 列必须把每一行的完整字节都读出来 列式Arrow col_id: [1, 2, 3, ...] ← 连续内存块 A col_name: [a,b,c, ...] ← 连续内存块 B col_age: [20, 25, 30, ...] ← 连续内存块 C → 要读第 2 列直接引用内存块 B不需要碰 A 和 C投影的本质就是把不需要的列的FieldVector引用丢掉只保留需要的。Arrow 的TransferPair机制让这个操作是 O(1) 的引用转移不产生任何数据拷贝classColumnProjector{publicVectorSchemaRootproject(VectorSchemaRootfull,int[]projectedColumns){// 对每个需要的列// TransferPair pair vector.getTransferPair(allocator);// pair.transfer(); ← 零拷贝只是转移 buffer 所有权// result.add(pair.getTo());// 不需要的列直接 close() 释放}}3.3 优化二谓词下推Predicate PushdownLog - Log : 应用 WHERE 过滤条件对 Arrow Vector 向量化过滤这一步是向量化执行不是一行一行判断而是对一个 Arrow Vector 批量判断生成一个 selection vector位图。// 伪代码向量化过滤publicSelectionVectorapplyPredicate(FieldVectoramountVector,Predicatepred){// 传统for each row: if (row.amount 100) keep// 向量化一次性处理 1024 个值生成 bitmask// long[] bitmask new long[1024 / 64];// for (int i 0; i 1024; i) {// if (amountVector.get(i) 100) bitmask[i/64] | (1L (i%64));// }}向量化过滤的好处CPU cache 友好连续内存访问可被 JIT 向量化现代 JVM 能把这种循环编译成 SIMD 指令分支预测友好没有 per-row 的 if 分支跳转3.4 效果量化图中 note 给出了一个非常直观的数字列裁剪 谓词下推在服务端完成: 200 列的表只读 2 列 → 网络传输减少 99% WHERE 过滤在 TabletServer 执行 → 减少无效数据传输我们算一下那个180 字段只用 2 个的大屏查询传统方案KafkaFluss网络传输50万 QPS × 2KB/行 1 GB/s50万 QPS × 20B 10 MB/s反序列化 CPU180 字段全解析只解析 2 列Flink TM 内存完整对象GC 压力大只有 2 列的对象过滤位置Flink 侧数据已经传过来了TabletServer 侧网络传输降低 99%这个量级的差距足以让一个跑不动的作业变成跑得很轻松。3.5 客户端侧的收尾Client - Client : Arrow → Flink RowData 转换 Client - Client : 更新 Checkpoint offset第二步是 Flink 容错的关键offset 随 Checkpoint 一起持久化。Flink 的 Checkpoint 完成后offset 才被提交故障恢复时从 Checkpoint 里的 offset 重新消费。3.6 请求 B 的答案SELECTSUM(amount),COUNT(*)FROMordersWHEREevent_timeNOW()-INTERVAL5MINUTE;Flink 优化器会把SUM(amount)用到的列推导出来生成projectedColumns [amount, event_time]把event_time ...下推为 predicate。TabletServer 只返回这两列的过滤后数据。验证方法——看 Flink 的执行计划EXPLAINSELECTSUM(amount),COUNT(*)FROMordersWHEREevent_timeNOW()-INTERVAL5MINUTE;如果看到 Source 上有projectedFields或类似的提示说明列裁剪生效了。四、图 3PK Lookup 路径4.1 三级优化链图下方的 note 总结了 Lookup 的性能优化链性能优化链: 1. LRU 缓存命中 → 零网络开销 2. Bloom Filter → 快速排除不存在的 key 3. RocksDB LSM → 亚毫秒级点查询这三级的成本递减、命中率递减构成了一个典型的缓存层次结构第 1 级LRU 缓存Flink TaskManager 本地堆内存 成本~100 ns 命中率取决于维表大小和缓存容量热点数据可到 80% 第 2 级Bloom FilterTabletServer 内存 成本~1 μs只判断可能存在还是一定不存在 作用对于不存在的 key直接返回 null避免 RocksDB 查询 第 3 级RocksDB LSM 查询 成本~100 μs ~ 1 ms 路径MemTable → Immutable MemTable → L0 → L1 → ... → L64.2 第一级LRU 缓存classFlussLookupFunctionextendsLookupFunction{privatefinalCacheRowData,RowDatalookupCache;// Guava/Caffeine LRUprivatefinalFlussConnectionconnection;publicvoideval(Object...joinKeys){RowDatakeytoRowData(joinKeys);RowDatacachedlookupCache.getIfPresent(key);if(cached!null){collect(cached);// 命中零网络开销return;}// 未命中走下面的远程查询// ...lookupCache.put(key,value);// 回填缓存}}配置CREATETABLEuser_profile(user_idBIGINT,user_name STRING,city STRING,PRIMARYKEY(user_id)NOTENFORCED)WITH(bucket.num16,lookup.cache.max-rows100000,-- 最多缓存 10 万行lookup.cache.ttl1h-- 缓存 1 小时过期);调优要点参数影响建议lookup.cache.max-rows命中率 vs 内存维表小 100万行可以全缓存大表按热点比例设 10%-20%lookup.cache.ttl数据新鲜度维表更新频繁 → 设短分钟级静态维表 → 设长或不设缓存一致性的坑如果维表在 Fluss 里被更新了Flink 侧的 LRU 缓存不会自动失效要等 TTL 过期。所以维表更新频率高且要求强一致的场景TTL 不能设太长。4.3 第二级Bloom FilterServer - Bloom : mayContain(key) Bloom -- Server : 不存在 (快速返回) → 直接返回 null Bloom -- Server : 可能存在 → 继续查 RocksDBBloom Filter 的特点是可能误报存在但绝不会误报不存在。返回不存在 → 100% 确定这个 key 没有可以安全返回 null返回可能存在 → 大概率存在需要真的查一次 RocksDB上一篇讲过BLOOM_BITS_PER_KEY 10.0这个配置下假阳性率约 1%。也就是说查 10000 个不存在的 key - 约 9900 个被 Bloom Filter 拦截 → 零磁盘 IO - 约 100 个假阳性 → 白跑一次 RocksDB 查询但返回 null这个优化对于大量 key 命不中的场景效果拔群。比如风控场景大部分用户不在黑名单里Bloom Filter 能拦掉 99% 的无谓查询。4.4 第三级RocksDB LSM 查询Kv - Kv : 查 MemTable → Immutable → L0~L6 SSTRocksDB 的读路径图中标注的顺序1. Active MemTable内存跳表最新写入 ~ 100 ns 2. Immutable MemTable正在 flush 的 ~ 100 ns 3. Block Cache读缓存 ~ 1 μs 4. L0 SST 文件可能有重叠要查多个 ~ 10-100 μs 5. L1 ~ L6 SST 文件每层有序二分查找 ~ 10-100 μs/层读放大问题LSM 树的层级越多查询可能要读的文件越多。这就是为什么BLOCK_CACHE_SIZE和 Bloom Filter 都很重要——它们能把大部分查询挡在第 4 步之前。4.5 请求 C 的答案SELECTo.order_id,u.user_name,u.cityFROMorders oLEFTJOINuser_profileFORSYSTEM_TIMEASOFo.proc_timeASuONo.user_idu.user_id;图下方的 note 也标注了这条 SQL对应 SQL: LEFT JOIN dim FOR SYSTEM_TIME AS OF o.time LRU 缓存降低服务端查询压力性能对比方案单次关联延迟50万 QPS 下的表现Redis 维表0.5 - 1 ms网络 RTT需要 Redis 集群且仍受网络限制HBase 维表2 - 10 ms扛不住需要大集群Flink 双流 Join无网络但状态巨大状态 TB 级Checkpoint 慢Fluss PK Lookup缓存命中 ~0.1μs未命中 ~0.5ms热点命中率高时整体 P99 1ms核心优势Fluss 把维表数据存在 TabletServer 本地配合 LRU Bloom RocksDB 三级优化把远程查询变成了大部分时候的本地内存查询。五、动手验证5.1 观察列裁剪的效果-- 建一张宽表CREATETABLEwide_table(idBIGINT,c01 STRING,c02 STRING,...c50 STRING,-- 50 个 STRING 字段amountDECIMAL(18,2))WITH(bucket.num8);-- 只查 2 列EXPLAINSELECTid,amountFROMwide_tableWHEREamount100;在 TabletServer 的 metrics 里观察curlhttp://tablet-server:9125/metrics|grep-ibytes_out对比全表扫描和只查 2 列的网络字节数差距应该在两个数量级。5.2 验证 Lookup 缓存命中率curlhttp://tablet-server:9125/metrics|grep-ilookup关注指标fluss_lookup_cache_hit_count fluss_lookup_cache_miss_count fluss_lookup_cache_hit_ratio ← 目标 0.8 fluss_lookup_remote_latency_p99 ← 目标 1ms5.3 验证 Exactly-Once-- 开启 CheckpointSETexecution.checkpointing.interval10s;SETexecution.checkpointing.modeEXACTLY_ONCE;-- 写入并观察INSERTINTOorders_sinkSELECT*FROMsource;手动杀掉 Flink 作业从 Checkpoint 恢复检查目标表的数据条数是否有重复SELECTCOUNT(*),COUNT(DISTINCTorder_id)FROMorders_sink;-- 两者相等 → 无重复六、生产实践要点6.1 写入调优清单# 客户端 client.request.acks: all client.writer.batch.size: 1MB client.writer.linger.ms: 10 client.writer.compression.type: lz4 # 压缩牺牲 CPU 换带宽 # 服务端 tablet-server.min.insync.replicas: 2 log.segment.size: 1GB log.flush.interval.messages: 10000 # 攒多少条刷盘 log.flush.interval.ms: 10006.2 读取调优清单# 增大 fetch 批次减少 RPC 往返 client.scanner.fetch.max-bytes: 64MB client.scanner.fetch.max-wait-ms: 500 # 服务端 Arrow 内存池 arrow.allocator.memory.limit: 2GB6.3 Lookup 调优清单-- 维表设计的黄金法则-- 1. bucket.num 要足够大避免热点 bucket-- 2. 开启缓存并设置合理 TTL-- 3. 维表字段不要太多投影只取需要的列CREATETABLEdim_user(user_idBIGINT,user_name STRING,city STRING,PRIMARYKEY(user_id)NOTENFORCED)WITH(bucket.num64,-- 足够分散lookup.cache.max-rows500000,-- 缓存 50 万行lookup.cache.ttl30min-- 30 分钟过期);6.4 常见反模式反模式问题正确做法acks0追求低延迟数据可能丢用acksall 批量写入来降延迟而不是降低持久性Lookup 缓存 TTL 设成 1 天维表更新后读到脏数据按维表更新频率设 TTL或改用双流 Join用 PK 表存纯日志数据存储成本翻 2.5 倍纯追加场景用 Log 表bucket.num设成 1完全无法并行单 bucket 成为瓶颈至少设为 TabletServer 数的 2-4 倍七、排障手册现象可能原因排查方向写入延迟突然升高ISR 收缩触发acksall等待检查 Follower 是否有 GC 或网络问题curl .../metrics | grep isr写入报NOT_LEADER_FOR_BUCKETLeader 刚切换客户端会自动刷新路由重试如果持续报错检查 Coordinator 是否在频繁 Rebalance读取吞吐上不去列裁剪没生效检查 Flink 执行计划确认谓词可下推不支持 UDF 过滤下推Lookup P99 高缓存命中率低检查lookup_cache_hit_ratio调大max-rows或缩小维表数据出现重复Checkpoint 没开或两阶段提交没生效确认execution.checkpointing.mode EXACTLY_ONCE和sink.delivery-guarantee exactly-once数据丢失min.insync.replicas未配置补上配置检查 ISR 最小副本数八、小结三条路径三种优化思路路径核心优化关键机制写入元数据本地缓存 顺序写 WALbucket 路由本地计算、ISR 复制、两阶段提交 幂等写入扫描服务端列裁剪 谓词下推Arrow 列式格式、TransferPair零拷贝投影、向量化过滤点查三级缓存层次LRU 本地缓存 → Bloom Filter 快速排除 → RocksDB LSM三句话记住写入Coordinator 不在数据路径上先 WAL 后 RocksDBacksallmin.insync.replicas2是持久性的底线。扫描请求里带投影列和谓词服务端做列裁剪和过滤宽表场景网络传输能降 99%。点查LRU → Bloom → RocksDB 三级优化把远程查询变成大部分时候的本地内存查询。下一篇我们关注集群本身启动选举、副本状态机、以及扩缩容时的数据迁移。