十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Fluss 核心功能点源码阅读计划

Apache Fluss 核心功能点源码阅读计划 Apache Fluss 核心功能点源码阅读计划按核心功能模块拆解的深度源码阅读路线每个功能点标注了精确的文件路径和阅读顺序阅读总览本计划将 Fluss 源码按12 个核心功能点拆解每个功能点包含 功能概述 源码路径清单精确到文件 推荐阅读顺序由浅入深 核心类和关键方法标注预计总阅读量约 200 核心源文件预计 4-6 周完成。第 1 周基础架构功能点 1-3 → 了解骨架 第 2 周存储引擎功能点 4-6 → 深入核心 第 3 周分布式机制功能点 7-8 → 理解协调 第 4-6 周高级特性功能点 9-12 → 全景掌握功能点 1项目骨架与启动流程概述理解整个项目的模块组织、构建系统、启动流程和配置加载机制。源码路径项目根目录 ├── pom.xml ← Maven 多模块构建入口 ├── fluss-dist/ ← 发行版打包 ├── fluss-server/ │ └── src/main/java/org/apache/fluss/server/ │ ├── coordinator/ │ │ └── CoordinatorServer.java ★ 入口类Coordinator 启动 │ └── tablet/ │ └── TabletServer.java ★ 入口类TabletServer 启动 │ ├── fluss-common/ │ └── src/main/java/org/apache/fluss/ │ ├── config/ │ │ └── ConfigOption.java ★ 配置项定义框架 │ └── cluster/ │ ├── ServerNode.java ★ 节点表示 │ └── ServerType.java ★ Coordinator vs TabletServer 枚举 │ └── fluss-rpc/ └── src/main/java/org/apache/fluss/rpc/ ├── RpcServer.java ★ RPC 服务启动 └── RpcGateway.java ★ RPC 网关接口阅读顺序序号文件关注重点1pom.xml模块依赖关系、Java 版本、关键依赖版本2ServerType.javaCoordinator 和 TabletServer 的角色划分3CoordinatorServer.java启动流程ZK 连接 → 元数据初始化 → RPC 启动 → 事件循环4TabletServer.java启动流程注册 → 加载 Tablet → 开始服务5RpcServer.javaProtobuf RPC 如何注册和处理请求6ConfigOption.javaFluss 的配置系统设计模式核心问题CoordinatorServer 和 TabletServer 启动时分别做了哪些初始化它们如何发现彼此通过 ZooKeeper 注册和心跳配置项是如何定义和加载的功能点 2元数据管理与 Table 生命周期概述理解 Fluss 如何管理 Database、Table、Partition、Schema 等元数据以及创建表、修改表、删除表的完整生命周期。源码路径fluss-server/src/main/java/org/apache/fluss/server/ ├── coordinator/ │ ├── MetadataManager.java ★ 元数据管理核心与 ZK 交互 │ ├── TableManager.java ★ 表生命周期管理创建/删除/修改 │ ├── TableLifecycleThrottler.java ★ 表操作限流 │ └── SchemaUpdate.java ★ Schema 变更ADD COLUMN │ ├── metadata/ │ ├── TableMetadata.java ★ 表元数据结构 │ ├── PartitionMetadata.java ★ 分区元数据 │ └── TabletMetadata.java ★ Tablet 元数据 │ └── zk/ ├── ZooKeeperClient.java ★ ZK 客户端封装 ├── ZkMetadataStorage.java ★ ZK 上的元数据存储格式 └── ZkTableOperations.java ★ 表的 ZK 操作 fluss-common/src/main/java/org/apache/fluss/ ├── metadata/ │ ├── TableDescriptor.java ★ 表定义类型、Schema、Bucket 数 │ ├── TablePath.java ★ Database.Table 路径 │ └── Schema.java ★ 列定义、数据类型 │ └── config/ └── ConfigOptions.java ★ 所有可配置参数定义阅读顺序序号文件关注重点1TableDescriptor.java表类型的枚举、Schema 结构、Bucket 配置2Schema.java列的数据类型系统3TablePath.javadatabase.table 的层级表示4TableMetadata.java表在服务端的完整元数据结构5MetadataManager.javacreateTable/deleteTable 的完整流程6TableManager.java表的生命周期状态机7SchemaUpdate.javaADD COLUMN 的实现细节8ZkTableOperations.java元数据在 ZooKeeper 上的序列化格式核心问题Fluss 如何区分 Log Table 和 Primary Key TableCREATE TABLE从 SQL 到落盘经历了哪些步骤Schema Evolution 的 ADD COLUMN 对已有数据有何影响元数据在 ZooKeeper 上是如何组织的功能点 3Bucket、Partition 与数据分布概述理解 Fluss 如何将一张表的数据切分为 Partition 和 Bucket以及 Bucket 到 TabletServer 的映射、Rebalance 机制。源码路径fluss-server/src/main/java/org/apache/fluss/server/ ├── coordinator/ │ ├── AutoPartitionManager.java ★ 自动分区管理 │ └── rebalance/ │ ├── RebalanceCoordinator.java ★ 重平衡协调器 │ ├── RebalancePlan.java ★ 重平衡计划 │ └── TabletBalancer.java ★ Tablet 分配算法 │ └── tablet/ ├── TabletManager.java ★ Tablet 运行时管理 └── TabletAssignment.java ★ Tablet 到 Server 的映射 fluss-common/src/main/java/org/apache/fluss/ ├── bucketing/ │ ├── BucketingFunction.java ★ 分桶函数Hash/MurmurHash │ └── BucketKeyExtractor.java ★ 从 Row 提取 Bucket Key │ └── cluster/ └── BucketLocation.java ★ Bucket 物理位置Server 路径阅读顺序序号文件关注重点1BucketingFunction.javagetBucket(key, numBuckets)哈希算法2BucketKeyExtractor.java如何从行数据中提取分桶键3AutoPartitionManager.java分区自动创建的触发条件和流程4TabletBalancer.javaTablet 分配到 TabletServer 的算法轮询容量感知5RebalanceCoordinator.java扩缩容时如何重新分配 Tablet6RebalancePlan.java迁移计划的数据结构哪些 Tablet 从哪移到哪核心问题bucket.num参数如何影响数据分布Rebalance 过程中数据一致性如何保证PK 表的分区列为什么必须是主键子集功能点 4LogStore —— 日志存储引擎概述这是 Fluss 写入路径的核心。理解 LogTablet → LogSegment → .log/.index 文件的完整存储机制。源码路径fluss-server/src/main/java/org/apache/fluss/server/log/ ├── LogTablet.java ★ 核心门面append/read/delete ├── LogSegment.java ★ 段文件.log .index 管理 ├── LogSegments.java ★ 段集合管理TreeMap 组织 ├── LocalLog.java ★ 本地日志目录管理 ├── LogManager.java ★ 所有 LogTablet 的生命周期管理 ├── LogLoader.java ★ 启动时加载已有日志段 ├── OffsetIndex.java ★ 偏移量稀疏索引二分查找 ├── TimeIndex.java ★ 时间戳索引 ├── AbstractIndex.java ★ 索引基类mmap 实现 ├── LazyIndex.java ★ 延迟加载索引 ├── FetchParams.java ★ 读取参数offset、maxBytes ├── FetchIsolation.java ★ 读取隔离级别 ├── LogAppendInfo.java ★ 追加操作的返回信息 ├── LogReadInfo.java ★ 读取操作的返回信息 ├── RollParams.java ★ 日志段滚动Rolling的条件 ├── SnapshotFile.java ★ 日志快照文件 │ ├── WriterStateManager.java ★ 写入器幂等性Epoch Sequence ├── WriterStateEntry.java ★ 单个写入器的状态 │ ├── FilterContext.java ★ 读取时的过滤上下文 ├── PredicateSchemaResolver.java ★ 谓词下推的 Schema 解析 │ ├── checkpoint/ │ ├── CheckpointFile.java ★ 检查点文件格式 │ └── OffsetCheckpoint.java ★ Offset 检查点记录已刷盘的 offset │ └── remote/ ├── RemoteLogManager.java ★ 远程日志管理 └── RemoteLogSegment.java ★ 远程日志段阅读顺序序号文件关注重点Phase 1: 索引系统1AbstractIndex.javammap 实现、二分查找2OffsetIndex.javalookup(offset) → position映射3TimeIndex.javalookup(timestamp) → offset映射Phase 2: 段管理4LogSegment.javaappend(batch),read(offset, maxBytes), 何时 Roll5LogSegments.javaConcurrentSkipListMapbaseOffset, LogSegment6LocalLog.java目录结构、日志文件命名规则Phase 3: Tablet 层7LogTablet.java★ 核心append/read/highWatermark/ISR8LogManager.java创建/删除/清理 LogTablet9LogLoader.java启动恢复扫描目录、加载段、重建索引Phase 4: 高级特性10WriterStateManager.java幂等写入Epoch 和 Sequence Number11FilterContext.java服务端过滤的实现12PredicateSchemaResolver.java谓词下推与 Arrow Schema 结合核心问题append()的完整链路是什么内存 → OS Page Cache → 磁盘Segment Roll 的触发条件是什么大小时间稀疏索引如何做到 O(log N) 查找启动恢复时如何从 .log 和 .index 文件重建状态功能点 5KvStore —— RocksDB 键值存储引擎概述理解 PK 表的可查询能力来源——KvStore 如何基于 RocksDB 实现高效的 Upsert/Get/Delete 操作。源码路径fluss-server/src/main/java/org/apache/fluss/server/kv/ ├── KvFlushScheduler.java ★ KV 刷盘调度 │ ├── rocksdb/ │ ├── RocksDBKv.java ★ RocksDB 封装核心 │ ├── RocksDBConfig.java ★ RocksDB 参数配置 │ ├── RocksDBWriteBatch.java ★ 批量写入 │ └── RocksDBSnapshot.java ★ 快照管理 │ ├── wal/ │ ├── KvWalManager.java ★ KV 的 WAL 管理 │ └── WalRecovery.java ★ WAL 恢复逻辑 │ ├── rowmerger/ │ ├── RowMerger.java ★ 行合并接口 │ ├── DeduplicateRowMerger.java ★ 去重合 │ ├── PartialUpdateRowMerger.java ★ 部分更新合并 │ └── AggregationRowMerger.java ★ 聚合合并 │ ├── partialupdate/ │ └── PartialUpdateHandler.java ★ 部分更新处理 │ ├── prewrite/ │ └── PrewriteManager.java ★ 预写管理 │ ├── scan/ │ └── KvScan.java ★ 范围扫描 │ └── autoinc/ └── AutoIncrementManager.java ★ 自增列管理阅读顺序序号文件关注重点1RocksDBConfig.javaBlock Cache、Write Buffer、Compaction 策略2RocksDBKv.javaput/get/delete 的封装、Bloom Filter 配置3KvWalManager.java写入流程先写 WAL 再写 RocksDB4WalRecovery.java故障恢复从 WAL 重放到 RocksDB5RowMerger.java合并引擎的抽象接口6DeduplicateRowMerger.java默认去重保留最新值7PartialUpdateRowMerger.java部分更新NULL 字段不覆盖已有值8AggregationRowMerger.java聚合sum/max/min/count/last_value9KvFlushScheduler.java何时将 MemTable 刷到 SST 文件10RocksDBSnapshot.java快照的创建和恢复11KvScan.java前缀扫描、范围扫描核心问题RocksDB 的 LSM 结构如何支持高吞吐写入和亚毫秒读取WAL 和 RocksDB 的写入顺序是什么WAL 先还是 RocksDB 先Partial Update 两个 NULL 列如何正确合并快照是如何实现一致性备份的功能点 6Arrow 列式格式与数据读写概述理解 Fluss 如何利用 Apache Arrow 实现列式数据存储以及列裁剪和零拷贝读取的底层机制。源码路径fluss-common/src/main/java/org/apache/fluss/ ├── record/ │ ├── ArrowLogReader.java ★ Arrow 格式日志读取器 │ ├── ArrowLogWriter.java ★ Arrow 格式日志写入器 │ ├── ArrowRecordBatch.java ★ Arrow 批次封装 │ ├── LogRecordBatch.java ★ 日志批次抽象 │ └── ColumnProjector.java ★ 列投影裁剪 │ ├── row/ │ ├── RowData.java ★ 行数据接口 │ ├── ArrowRowData.java ★ Arrow 行数据实现 │ └── InternalRow.java ★ 内部行表示 │ ├── types/ │ ├── DataType.java ★ 数据类型系统 │ ├── DataTypes.java ★ 类型工厂 │ └── RowType.java ★ 行类型Schema │ ├── memory/ │ └── ArrowMemoryAllocator.java ★ Arrow 内存分配器 │ └── shaded/arrow/ └── (shaded Arrow 依赖) fluss-server/src/main/java/org/apache/fluss/server/log/ └── FilterContext.java ★ 列裁剪 谓词下推的过滤上下文阅读顺序序号文件关注重点1DataType.java/DataTypes.javaFluss 支持哪些数据类型与 Arrow 类型如何映射2LogRecordBatch.java批次的数据结构schema rows offset3ArrowRecordBatch.javaArrow VectorSchemaRoot 的封装4ArrowLogWriter.java★ 列式写入逐列填充 Arrow Vector5ArrowLogReader.java★ 列式读取从 Arrow IPC 格式反序列化6ColumnProjector.java列裁剪只提取需要的 Vector零拷贝引用7ArrowMemoryAllocator.javaDirect Memory 管理、chunk 分配核心问题Arrow IPC 格式在.log文件中是如何布局的列裁剪时如何做到「零拷贝」引用而非复制数据200 列的表只读 2 列I/O 如何降低 99%Arrow 格式如何与 Flink 的 RowData 互相转换功能点 7副本机制与 ISR概述理解 Fluss 如何基于 ISR 模型实现数据副本的高可用——Leader 选举、日志同步、故障切换。源码路径fluss-server/src/main/java/org/apache/fluss/server/ ├── replica/ │ ├── ReplicaManager.java ★ 副本管理器核心 │ ├── Replica.java ★ 单个副本的表示 │ ├── ReplicaFetcher.java ★ Follower 从 Leader 拉取数据 │ ├── LeaderReplica.java ★ Leader 副本的读写处理 │ └── FollowerReplica.java ★ Follower 副本的同步逻辑 │ ├── coordinator/ │ ├── CoordinatorLeaderElection.java ★ Coordinator 的 Leader 选举ZK │ └── statemachine/ │ ├── ReplicaStateMachine.java ★ 副本状态机 │ └── PartitionStateMachine.java ★ 分区状态机 │ └── kv/ └── KvSnapshotManager.java ★ KV 快照管理用于 Follower 恢复阅读顺序序号文件关注重点1Replica.java副本的抽象id、角色Leader/Follower、状态2ReplicaStateMachine.java状态转换NewReplica → Online → Offline3LeaderReplica.javaLeader 的 append 和 read 处理4FollowerReplica.javaFollower 的同步逻辑fetch → append → ack5ReplicaFetcher.javaFollower 如何从 Leader 拉取数据类似 Kafka Fetcher6ReplicaManager.java★ 管理所有副本becomeLeader/becomeFollower7CoordinatorLeaderElection.javaCoordinator 故障时如何选举新 Leader8KvSnapshotManager.javaKvStore 如何为 Follower 提供快照核心问题ISRIn-Sync Replica的判定标准是什么Leader 故障时如何从 ISR 中选择新 LeaderFollower 落后太多时如何通过快照快速追平日志复制是同步的还是异步的一致性如何保证功能点 8Tiering 分层存储与 Lakehouse 集成概述理解 Fluss 如何将热数据Arrow自动 Compaction 为冷数据Parquet并写入 Iceberg/Paimon/Lance。源码路径fluss-server/src/main/java/org/apache/fluss/server/ ├── coordinator/ │ ├── LakeTableTieringManager.java ★ 分层存储管理器 │ ├── LakeCatalogDynamicLoader.java ★ Lake Catalog 加载器 │ └── RemoteStorageCleaner.java ★ 远程存储清理孤儿文件 │ └── tiering/ ├── TieringService.java ★ Tiering 服务入口 ├── TieringTask.java ★ 单个 Tiering 任务 └── compactor/ └── ArrowParquetCompactor.java ★ Arrow → Parquet 转换 fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/ ├── tiering/ │ ├── FlinkTieringService.java ★ Flink 驱动的 Tiering │ └── TieringCommitter.java ★ Tiering 提交器 │ └── lake/ ├── LakeTableHelper.java ★ Lake 表工具 └── LakeUnionRead.java ★ Union Read 实现 fluss-common/src/main/java/org/apache/fluss/lake/ ├── LakeFormat.java ★ Lake 格式接口 ├── iceberg/ │ ├── IcebergLakeFormat.java ★ Iceberg 格式实现 │ └── IcebergCommitter.java ★ Iceberg 提交器 ├── paimon/ │ ├── PaimonLakeFormat.java ★ Paimon 格式实现 │ └── PaimonCommitter.java ★ Paimon 提交器 └── lance/ ├── LanceLakeFormat.java ★ Lance 格式实现 └── LanceCommitter.java ★ Lance 提交器阅读顺序序号文件关注重点1LakeFormat.java湖格式的抽象接口2IcebergLakeFormat.javaIceberg 集成细节3PaimonLakeFormat.javaPaimon 集成细节4TieringService.java★ Tiering 的触发条件和执行流程5ArrowParquetCompactor.javaArrow → Parquet 的列式转换6LakeTableTieringManager.java哪些表需要 Tiering何时执行7LakeUnionRead.java★ 联合查询如何合并 Hot Cold 数据8RemoteStorageCleaner.java清理已 Tiering 的本地 Segment核心问题Tiering 是增量还是全量的Arrow → Parquet 转换时如何保证 Schema 兼容Union Read 在查询计划中如何表示和优化Iceberg 和 Paimon 作为冷存储有何实现差异功能点 9Flink Connector —— Catalog、Source、Sink概述理解 Fluss 如何注册为 Flink Catalog、如何实现流式和批量读写、以及 Exactly-Once 语义。源码路径fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/ ├── FlinkConnectorOptions.java ★ 所有 Connector 配置参数 │ ├── catalog/ │ ├── FlussCatalog.java ★ Flink Catalog 实现核心 │ ├── FlussCatalogFactory.java ★ Catalog 工厂 │ └── FlussDatabase.java ★ Database 抽象 │ ├── source/ │ ├── FlussSource.java ★ Flink Source 实现 │ ├── FlussSourceEnumerator.java ★ Split 发现与分配 │ ├── FlussSourceReader.java ★ 数据读取器 │ ├── FlussSourceSplit.java ★ Split 定义 │ └── lookup/ │ └── FlussLookupFunction.java ★ Lookup Join 函数 │ ├── sink/ │ ├── FlussSink.java ★ Flink Sink 实现 │ ├── FlussSinkWriter.java ★ Sink 写入器 │ ├── FlussSinkCommitter.java ★ Two-Phase Commit 提交器 │ └── writer/ │ ├── UpsertWriter.java ★ Upsert 写入 │ ├── AppendWriter.java ★ 追加写入 │ └── DeleteWriter.java ★ 删除写入 │ ├── adapter/ │ └── FlinkAdapter.java ★ Flink ↔ Fluss 适配器 │ ├── row/ │ └── RowDataSerializationSchema.java ★ RowData → Arrow 序列化 │ ├── lake/ │ └── LakeUnionRead.java ★ Lakehouse 联合读取 │ ├── metrics/ │ └── FlinkMetricReporter.java ★ Flink 指标报告 │ └── utils/ └── FlinkConnectorUtils.java ★ 工具函数阅读顺序序号文件关注重点Phase 1: Catalog1FlussCatalogFactory.javaCREATE CATALOG ... WITH (typefluss)的入口2FlussCatalog.java★ getTable/createTable/tableExists 的实现Phase 2: Source3FlussSource.javaSource 接口实现创建 Enumerator 和 Reader4FlussSourceEnumerator.javaSplit 发现从 Coordinator 获取 Tablet 列表5FlussSourceReader.java数据读取从 TabletServer 拉取 Arrow Batch6FlussSourceSplit.javaSplit TablePath Partition BucketIdPhase 3: Sink7FlussSink.javaSink 接口实现8FlussSinkWriter.java写入数据到 TabletServer9FlussSinkCommitter.javaTwo-Phase Commit 实现 Exactly-Once10UpsertWriter.javaUpsert 写入PK 表11AppendWriter.java追加写入Log 表Phase 4: Lookup12FlussLookupFunction.java★FOR SYSTEM_TIME AS OF的实现Phase 5: Lake 集成13LakeUnionRead.java如何构建 Hot Cold 联合查询计划核心问题Fluss Catalog 的getTable()如何将 Fluss 表元数据转换为 Flink DynamicTableSourceSource 的 Split 分配策略是什么轮询本地优先Sink 如何实现 Two-Phase Commit 保证 Exactly-OnceLookup Join 的缓存是如何实现的LRU Cache TTL功能点 10Delta Join 与状态外部化概述理解 Fluss 如何通过 Delta Join 将 Flink 的 Join 状态外部化到 TabletServer。源码路径fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/ ├── source/ │ └── delta/ │ ├── DeltaJoinOperator.java ★ Delta Join 算子 │ ├── DeltaJoinState.java ★ Join 状态管理 │ └── DeltaJoinConfig.java ★ Join 配置 │ └── sink/ └── delta/ └── DeltaJoinWriter.java ★ Join 结果回写 fluss-server/src/main/java/org/apache/fluss/server/ └── tablet/ └── delta/ ├── DeltaJoinManager.java ★ 服务端 Delta Join 管理 └── JoinStateStore.java ★ Join 状态的 RocksDB 存储阅读顺序序号文件关注重点1DeltaJoinConfig.javaJoin 配置左表/右表、Join Key、Buffer 大小2DeltaJoinOperator.java★ Flink 端无状态的 Join 处理逻辑3DeltaJoinState.java状态的外部化接口pointLookup4JoinStateStore.java★ 服务端KvStore 中存储 Join 状态5DeltaJoinManager.java服务端管理多个 Join 的状态核心问题Delta Join 如何做到 Flink 算子完全无状态Join 状态在 Fluss 端如何存储KvStore 中的 key join_key, value right_row故障恢复时新的 Flink 算子如何立即获取最新 Join 状态Delta Join vs 传统 Flink Managed State 的性能差异来源是什么功能点 11客户端与网络通信概述理解 Fluss 客户端如何与 Coordinator/TabletServer 通信以及 Protobuf RPC 协议的设计。源码路径fluss-client/src/main/java/org/apache/fluss/client/ ├── FlussClient.java ★ 客户端主入口 ├── FlussConnection.java ★ 连接管理 ├── ConnectionManager.java ★ 连接池 │ ├── read/ │ ├── LogScanner.java ★ 日志扫描器 │ ├── KvLookupClient.java ★ KV 查询客户端 │ └── ArrowBatchReader.java ★ Arrow 批次读取 │ ├── write/ │ ├── LogWriter.java ★ 日志写入器 │ ├── BatchBuilder.java ★ 批次构建器 │ └── WriterId.java ★ 写入器 ID幂等性 │ └── admin/ ├── AdminClient.java ★ 管理客户端DDL └── FlussAdmin.java ★ 管理操作接口 fluss-rpc/src/main/java/org/apache/fluss/rpc/ ├── RpcClient.java ★ RPC 客户端 ├── RpcServer.java ★ RPC 服务端 ├── RpcGateway.java ★ RPC 网关接口 ├── protobuf/ │ ├── ProtobufSerializer.java ★ Protobuf 序列化 │ └── RpcRequest.java ★ RPC 请求封装 └── netty/ ├── NettyClient.java ★ Netty 客户端 └── NettyServer.java ★ Netty 服务端阅读顺序序号文件关注重点1FlussClient.java客户端入口connect/close/read/write/admin2FlussConnection.java到 Coordinator 和 TabletServer 的连接3LogScanner.java★ 日志扫描指定 offset 读取批次数据4KvLookupClient.java★ KV 点查询PK Lookup 请求5LogWriter.java★ 日志写入batch 写入实现6BatchBuilder.java行数据 → Arrow RecordBatch7NettyClient.java/NettyServer.java底层网络框架8RpcGateway.java定义了哪些 RPC 方法核心问题客户端如何发现 Coordinator 和 TabletServer 的地址写入请求的负载均衡策略是什么PK Lookup 请求如何路由到正确的 TabletServer功能点 12监控、安全与运维工具概述理解 Fluss 的 Metrics 系统、认证授权机制、CLI 工具和运维脚本。源码路径fluss-metrics/src/main/java/org/apache/fluss/metrics/ ├── MetricRegistry.java ★ 指标注册中心 ├── MetricReporter.java ★ 指标报告器接口 ├── prometheus/ │ └── PrometheusReporter.java ★ Prometheus 格式输出 └── groups/ ├── ServerMetricGroup.java ★ 服务端指标组 └── TabletMetricGroup.java ★ Tablet 指标组 fluss-server/src/main/java/org/apache/fluss/server/ ├── metrics/ │ ├── ServerMetrics.java ★ 服务端指标定义 │ └── TabletMetrics.java ★ Tablet 级别指标 │ ├── authorizer/ │ └── FlussAuthorizer.java ★ 授权接口 │ ├── cli/ │ ├── FlussClusterCommand.java ★ 集群管理命令 │ └── FlussTableCommand.java ★ 表管理命令 │ └── tools/ └── ClusterTool.java ★ 运维工具集 fluss-server/src/main/java/org/apache/fluss/server/security/ ├── SASL/PLAIN 认证相关 ★ 认证机制阅读顺序序号文件关注重点1MetricRegistry.java指标如何注册和获取2ServerMetrics.javaCoordinator 暴露哪些指标3TabletMetrics.javaTabletServer 暴露哪些指标写入速率、延迟 P994PrometheusReporter.java如何输出 Prometheus 格式5FlussAuthorizer.java权限校验接口6FlussClusterCommand.java集群运维命令实现建议阅读节奏Week 1 ─ 功能点 1 2 骨架与元数据 ─ 建立Fluss 是做什么的全局认知 关键文件CoordinatorServer, TableManager, MetadataManager Week 2 ─ 功能点 3 4 数据分布 日志存储 ─ 理解数据如何切分、如何落盘 关键文件LogTablet, LogSegment, BucketingFunction, TabletBalancer Week 3 ─ 功能点 5 6 KV 存储 Arrow 列式 ─ 理解 PK 表的工作原理和性能基础 关键文件RocksDBKv, ArrowLogWriter, ColumnProjector, RowMerger Week 4 ─ 功能点 7 8 副本机制 分层存储 ─ 理解高可用和湖仓统一 关键文件ReplicaManager, TieringService, LakeUnionRead Week 5 ─ 功能点 9 Flink Connector ─ 理解 Fluss 与 Flink 的深度集成 关键文件FlussCatalog, FlussSource, FlussSink, FlussLookupFunction Week 6 ─ 功能点 10-12 Delta Join 客户端 运维 ─ 理解状态外部化、网络通信、运维工具 关键文件DeltaJoinOperator, FlussClient, MetricRegistry附录快速索引关键词核心文件PK Table 写入链路LogWriter→LogTablet.append()→KvTablet.put()Log Table 写入链路LogWriter→LogTablet.append()PK LookupKvLookupClient→RocksDBKv.get()SQL DDLFlussCatalog.getTable()→MetadataManager.createTable()Lookup JoinFlussLookupFunction→KvLookupClient→RocksDBKvDelta JoinDeltaJoinOperator→JoinStateStoreTieringTieringService→ArrowParquetCompactor→IcebergCommitter列裁剪ColumnProjector→ Arrow Vector 零拷贝RebalanceRebalanceCoordinator→TabletBalancer故障恢复LogLoaderWalRecovery本计划基于 https://github.com/apache/fluss 仓库 main 分支2026年8月
返回列表