OceanBase 是分布式 Shared-Nothing 关系型数据库。每个节点对等运行单个 observer 进程
(入口 src/observer/main.cpp → ObServer::start()),节点内自带 SQL、存储、事务、分布式日志四大引擎,
基于 Multi-Paxos(在 PaLF 中实现)保证副本一致性与高可用(RPO=0,RTO<8s)。
关键特性(摘自 README.md / docs/docs/en/architecture.md):
透明水平扩展
HTAP 单引擎
向量检索
MySQL 兼容
RPO=0 / RTO<8s
多租户隔离
src/logservice/palf),
显著降低 Paxos 心跳/选主开销。Tablet 通过 LS 之间的「迁移」实现负载均衡。仓库根目录的关键模块(节选 src/):
| 目录 | 角色 | 说明 |
|---|---|---|
src/observer | OBServer 进程框架 | 包含进程入口 main.cpp、网络层 net/、MySQL 协议 mysql/、
多租户运行容器 omt/、RPC 派发 ob_srv_deliver/xlator。 |
src/sql | SQL 引擎 | parser → resolver → rewrite → optimizer → code_generator → engine → executor,外加 PX 并行
engine/px 与数据访问层 das,计划缓存 plan_cache。 |
src/storage | 存储引擎(基于 LSM-Tree) | MemTable memtable/、SSTable / 微块 blocksstable/、合并 compaction/、
访问路径 access/、LS 管理 ls/、事务 tx/、tablet 元数据 meta_mem/、
列式 column_store/、检查点 checkpoint/。 |
src/logservice | 分布式日志服务 | PaLF(Multi-Paxos 实现)palf/、回放 replayservice/、应用回调 applyservice/、
归档 archiveservice/、CDC cdcservice/、选举 palf/election/。 |
src/rootserver | RootService(集群管控) | DDL、负载均衡 balance/、容灾 ob_disaster_recovery_*、
冻结/合并调度 freeze/、备份 backup/。RootService 本身也是一个 LS 的 Leader。 |
src/share | 共享基础 | Schema、location cache、partition table、object pool 等。 |
src/pl | PL/SQL 引擎 | Oracle/MySQL 兼容的存储过程。 |
src/objit | 表达式 JIT | LLVM 表达式编译。 |
src/plugin | 插件机制 | FTS、向量、JSON 插件挂载点。 |
deps/oblib | 基础库 | 容器、内存、网络、RPC、协议、加密。 |
一个 OBServer 进程内的整体分层(参考 src/observer/ob_server.h、
ob_srv_network_frame、ob_srv_deliver、ob_srv_xlator、omt/):
启动流程(ObServer::start() 简化):
main()
├─ ObServer::init() // 解析 cfg、初始化 schema/location/log_block_mgr 等
└─ ObServer::start()
├─ start ObSrvNetworkFrame // 监听 2881/2882
├─ start OMT(多租户线程池)
├─ start ObLSService / ObTenantFreezer / ObCheckpointService
├─ start LogService (PaLF) + ApplyService + ReplayService
├─ register heartbeat to RootService
└─ start RootService(若本机为 RS Leader 副本)
SQL 引擎入口 src/sql/ob_sql.h ObSql::stmt_query()。整条 pipeline:
语种路径 src/sql/plan_cache,ObPlanCache 用 SQL 文本(参数化后)+ schema_version + literval 模式为 key 缓存 ObPhysicalPlan。
ob_plan_set.h 区分 FAST / PL"_TYPE / PS 模式;literal 类型维度去重。SQL Plan Management(src/sql/spm)维护 evolved baseline。当优化器选择新计划且代价优于现有 baseline 阈值,会被加入 evolve 队列;DBA 或自动任务执行 verify 后确认升级为新 baseline。Outline(sql/ob_outline.h)和 UDR''udr` 提供人工 pin 计划能力。
ObLogicalOperator(sql/optimizer/ob_log_plan.h)树形 IR。ob_join_order_enum_idp.cpp(IDP 动态规划)+ ob_join_order_enum_permutation.cpp。多表时折叠使用 System-R / IDP-2 with DPVol。ob_access_path_estimation.cpp,对每个候选索引基于 ObOptStatManager(直方图、NDV、Density)计算 cost。索引估算区分 ss_rowcount/emory/seek/selectivity 并考虑 covering、skip scan、回表 bitcost。ob_dynamic_sampling.cpp,统计信息不足时即席 OpenGauss 风采样(默认 level=2)。ob_unnest)。ObStaticEngineCG 遍历 CBO 计划树 → 构造 ObOpSpec(算子规格)+ 表达式 spec,全部序列化可缓存可 RPC 复用。ObExpr,exec 路径 ObExprCG 生成为函数指针 vec 求值;复杂表达式可被 src/objit(LLVM JIT)编译为原生机器码。ObBatchRows(默认 2048 行,由 ObExecContext::batch_size_ 控制)。HashJoin / HashAgg 走。ob_exec_hash_struct_vec.cpp。ObSQLMemMgrProcessor 跟踪内存,超 workarea_percentage 走落盘(temp store)。算子工厂在 sql/engine/ob_operator_factory.{h,cpp}。生成两种模式:
ObOperator::get_next_row() 兼容路径(仅 PL 临时表或特定虚拟表内部走)。storage/access/ob_sstable_index_filter 把可在 IndexBlock 上求值的谓词涂色,减少数据微块 IO。storage/access/ob_pushdown_aggregate_vec 在 micro block 维度聚合,上层只需 reduce。ob_index_skip_scanner + ob_skip_index_sortedness:对前缀列范围整块跳过,性能类似 MySQL Index Skip Scan。sql/engine/aggregate/ob_adaptive_bypass_ctrl:HashAgg 小基数直接 bypass 到 partial/单分区模式,避免哪些需要二次聚合的开销。ob_optimizer_trace_impl 把 plan 选择全过程埋点导出表 opt_trace 用于调优。ObPlanCache::flush_cache by `refine_task`)。ObExpr 都有 eval 与 eval_batch 两种入口。如果某个表达式只支持 row 级(如 PL SPI 调用),CG 会保留变体并在 batch 角标 fallback 为逐行调用。代价是劣化但正确;不会全盘退化为火山式。parallel_max_servers + pool CPU)、数据量估算、小表 broadcasting 决定;Exchange 分布类型选择按 join 类型和表分布 hash/pkey/broadcast,并与 hash subpartition 的 redistribute cost 模型比较。Admission 在执行入口(PX Coord)拦截不合规 DOP。ob_enable_jit + objit not stripped。OLTP 一般关闭,AP 大 query scene 才 enabled。LLVM IR compile 需要~几十 ms,因此 sample 入库 first/warm-up 时不受限于。算子工厂在 sql/engine/ob_operator_factory.{h,cpp}。代码生成两种模式:
ObStaticEngineCG 生成基于 ObBatchRows/ObBitVector 的批量化算子,
表达式可被 JIT 编译(src/objit 基于 LLVM)。ObOperator::get_next_row() 兼容路径。主要算子族:
| 家族 | 代码位置 | 代表算子 |
|---|---|---|
| 聚合 | engine/aggregate | HashGroupBy / MergeGroupBy / HashDistinct(含向量化 vec 版本) |
| Join | engine/join | HashJoin / NestedLoop / MergeJoin / NLJ_Vec |
| 排序/集合 | engine/sort, set | 外排序 + 落盘、UNION/INTERSECT/EXCEPT |
| 窗口/CTE | engine/window_function,recursive_cte | Window、递归 CTE |
| 子查询 | engine/subquery | SubPlan Filter、Subplan Scan |
| 表/扫描 | engine/table | Table Scan、Direct Receive、Function Table、向量索引扫描 |
| DML | engine/dml | Insert/Update/Delete/Merge/Replace |
| PDML | engine/pdml | 并行 DML |
| 并行 | engine/px | PX Coord、DFO、Granule、Exchange |
DAS 是 SQL 与存储之间的数据访问中间层(src/sql/das),把每个 Scan/DML 拆分成对各 Tablet 的 ObDasTask,
并按数据位置远程/本地下推。核心类:ObDataAccessService、ObDASRef、ObDASScanOp、
ObDASInsert/Update/Delete/LockOp、定位 ObDASLocationRouter、重试 ObDASRetryCtrl。
PX 实现 MPP 风格多机并行。源码在 src/sql/engine/px。关键概念:
ob_dfo.h)。ob_granule_pump)。sql/dtl(Data Transfer Layer)传输。ob_dfo_scheduler)、Admission(ob_px_admission)。OceanBase 的存储引擎是准 LSM-Tree 架构,对应代码 src/storage:
storage/memtable),同时通过 PaLF 写多副本 REDO。storage/blocksstable)。位于 src/storage/access。读需要把同一行在 MemTable 与多层 SSTable 的版本「融合」:
blocksstable/ob_block_manager 管理。blocksstable/encoding 提供 RAW/Dict/RLE/Const/Diff/Prefix/HexString 等列式编码器,
并有 NEON/SIMD 优化解码(encoding/neon、ob_raw_decoder_simd.cpp、ob_dict_decoder_simd.cpp)。blocksstable/cs_encoding + storage/column_store 实现真正按列组织的 SSTable,
支持 HTAP 中的 AP 查询;CONVERT_CO_MAJOR_MERGE 把行存 major 转列存。| 缓存 | 路径 | 用途 |
|---|---|---|
| BlockCache | blocksstable/ob_block_cache_* | 缓存解码后的微块 |
| RowCache | blocksstable/ob_row_cache_* | 点查 row 级缓存 |
| FuseRowCache | blocksstable/ob_fuse_row_cache | 融合多 SSTable 的结果行 |
| BloomFilter | blocksstable/ob_bloom_filter_* | SSTable 上 rowkey 是否存在 |
数据源码:src/storage/memtable。每行结构:
ObMemtable
├─ ObQueryEngine // 用 KeyBtree + HashIndex 索引
│ ├─ ob_mvcc/ob_keybtree.{h,cpp} // B+树,叶节点为 ObMvccRow*
│ └─ ob_mt_hash.h // JIT hash 索引适配点查
├─ ObMvccEngine // 行级 latch + write handler
│ ├─ ObMvccRow // 同 rowkey 多版本链表 ObMvccTransNode
│ ├─ ObMvccTransNode // 一个版本:trans_version + data
│ └─ row_latch.h // 行写锁
├─ ObRedoLogGenerator // 与 clog 衔接
├─ ObRowCompactor // 行级 flush
└─ ObMemtableCtx // TX 局部上下文:回调链 + undo 链
ObMvccRow 上持有 row_latch(ob_row_latch.h)写锁后追加新 TransNode。若发现已有 大于自身 trans_seq 的未提交版本,则触发 write-write 冲突,回滚自身(ob_row_conflict_handler.{h,cpp})。ObMvccEngine::get 中遍历版本链,只接受 trans_version <= snapshot_scn 或当前事务自己写的版本。读时刻保存 read_snapshot_scn(来自 ObTxReadCtx,可能取自 GTS / Standby Timestamp / Pure Local)。对每个版本 n:
if n.trans_version <= read_snapshot_scn AND
n.commit_version 已经通过 PA + TxData落盘稳定:
可见 n
elif n.write_seq < = own_tx_seq (是自己写的): // self-modification
在 curr_stmt 内可见
else:
skip n
读会按 rowkey 融合:MemTable(活跃 / frozen) + 各层 SSTable(Mini / Minor / Major) → ObMultipleMerge 通过 LoserTree 多路归并,结合 bloom filter、index filter、push down filter 提前裁剪。
| 隔离级别 | 实现要点 | 出入点 |
|---|---|---|
| READ-COMMITTED | 每条 SQL 重新申请 snapshot_scn(来自 GTS) | ObTransService::get_read_snapshot |
| REPEATABLE-READ | 事务开始拿一次 snapshot_scn,复用整个事务 | ObSqlCtx::isolation_ |
| SERIALIZABLE | snapshot SCN + 写时检查 write-write 冲突(等于 SI + 首冲突回滚) | ob_row_conflict_handler |
| 弱一致读 / Standby Read | 用 standby_time_service 或某个 PAlog位点回放后的 SCN | ob_standby_timestamp_service、ob_tx_sby_read_define |
OB 不再像一维文件系统(一行一文件),而是在指定目录(squarerouter/ob_file_system_router.h)下按 tenant/ls/tablet 分层;macro block 由 ob_block_manager 以一组大小为 2MB 的 PGE 文件统一划块回收,写入通过 libeasy / io_uring 提交 IO。
memtable/mvcc/ob_keybtree),并可选 hash 索引(ob_mt_hash.h)作为点查加速。选 B+树的理由:1)多版本追加写入需要顺序扫描性能稳定;2)缓存的 chunk 集中便于复用;3)回放 / dump 时天然顺序读 rowkey。RocksDB 老牌 SkipList 是为了无锁插入,OB 用 latch + B+ 树 在长 scan 上占优。ObRedoLogGenerator::submit);redo log 多 Buff写后回调成功后再"提交到" MemTable 内存结构。即使 MemTable 还在内存中刷盘前进程崩溃,重启从 clog replay 即可恢复。Memtable 本身 flush 出的小 sstable 文件 + clog 双重持久化,clog的多数派可追保住 RPO=0。ob_replay_handler);(2)只读 LS 上的恢复操作 + reconfirm会启动恢复。并行回放(多 LS 线程)、多 worker 加速;8s 对故障的短极限目标是指 leader 损失场景下 lease 过期 + 准连锁 reconfirm 后新 leader 接写。并非进程全量重启时间。合并类型在 src/storage/compaction/ob_compaction_util.h 定义:
enum MergeType {
INVALID_MERGE_TYPE,
MINOR_MERGE, // 多个 Mini -> 更大 Mini
HISTORY_MINOR_MERGE,
META_MAJOR_MERGE,
MINI_MERGE, // MemTable flush -> Mini SSTable
MAJOR_MERGE, // 全量基线合并(每日)
MEDIUM_MERGE, // 中型增量合并
...
CONVERT_CO_MAJOR_MERGE, // 行存 Major -> 列存 CG SSTables
INC_MAJOR_MERGE, // 增量 Major
};
compaction/ob_compaction_dag_ranker、ob_schedule_dag_func。ob_partition_merger.*、ob_partition_merge_fuser.*(多 SSTable 行融合 fuser)。i_compaction_filter 与 compaction/filter:丢弃已删除/过期版本,回收空间。ob_compaction_diagnose、ob_compaction_suggestion。column_store。snapshot_version + medium_list 实现,确保所有副本 Major 后产生相同基线版本(见 ob_medium_list_checker、ob_extra_medium_info、ob_medium_compaction_mgr)。| 类型 | 阶段 | 输入 | 输出 | 触发条件 | 是否停写 | 关键代码 |
|---|---|---|---|---|---|---|
| MINI_MERGE | 实时 | 1 个冻结 MemTable | 1 个 Mini SSTable(增量基线 0) | MemTable 冻结后立即触发;阈值受 ob_tenant_freezer 控制(memstore 内存 / 行数 / 时间) |
否(后台) | ob_tablet_memtable_mgr::schedule_mini_merge、ob_mini_merge |
| MINOR_MERGE | 近实时 | N 个 Mini SSTable + 历史 minor | 1 个更大的 Minor SSTable | Mini 数达到 minor_compact_trigger(默认 2~3) | 否 | ob_schedule_tablet_func::schedule_minor_merge、ob_partition_merge_iter |
| HISTORY_MINOR_MERGE | 后台 | 历史 Minor + Minis | 归并历史 | 清理 stale minor、控制 minor 层数 | 否 | HISTORY_MINOR_MERGE 分支 |
| MEDIUM_MERGE | 中型增量 | Major + Minor/Mini(range 内) | 新 Medium SSTable | RootService 在 medium_list 为每个 tablet 颁发 MediumCompactionInfo;watermark 之上的 minor/mini 都是被合并对象 | 否 | ob_medium_compaction_mgr、ob_medium_loop、ob_medium_compaction_func |
| MAJOR_MERGE | 全量基线 | 当前 Major + 其上所有增量 | 1 个新 Major SSTable(全表 flat) | 全局冻结(ObTenantFreezer::do_major_freeze)后,RootService 下发到所有 tablet;所有副本按相同 snapshot_version 做 | 否(不停写) | ob_root_minor_freeze、ob_tenant_freezer、ob_basic_schedule_tablet_func::schedule_major_merge、ob_partition_merge_fuser |
| META_MAJOR_MERGE | 全量 | meta tablet(如 LOB/aux) | 新 meta SSTable | 随 Major 联动 | 否 | META_MAJOR_MERGE 分支 |
| CONVERT_CO_MAJOR_MERGE | 行列转换 | 行存 Major | 列存 CG(Column Group)SSTables | 用户切列存策略 / 列存化升级 | 读路径切换 | column_store/、cs_encoding/、ob_dag_macro_block_writer |
| INC_MAJOR_MERGE | 增量基线 | Major + 增量 | 在新 Major 中保留部分增量 | 大表渐进合并 | 否 | ob_progressive_merge_helper |
| MDS_MINI / MDS_MINOR | 多源数据 | multi_data_source 制数据 | — | MDS 模块自带 mini/minor | 否 | ob_mds_filter_info、multi_data_source/ |
| BATCH_FREEZE_TABLETS | 批量冻结 | 多小 tablet memstore | 同时冻结 | 节省 clog 流量 | 否 | ob_batch_freeze_tablets_dag |
ObTenantFreezer(storage/tx_storage/ob_tenant_freezer.h)控制租户级写入门槛:
freeze_trigger_percentage(默认 20%)后触发该 tablet 的小冻结。memstore_limit_percentage(默认 90%)开始让写操作回滚并走 Throttle(log_throttle、ob_memtable_write_throttle)。ALTER SYSTEM MAJOR FREEZE / 介质变更 passed (compaction_suggestion)。所有合并统一进入 DAG 调度器(share/scheduler/ + compaction/ob_compaction_dag_ranker):每个 tablet 合并抽象为 dag(含多 task)。ObDagRanker 根据如下维度打分:
suspend_merging 控制)ob_compaction_diagnose 检测 "schedule too slow"/"merge waiting" 等,ob_compaction_suggestion 给出自动放大资源/放宽时机的提议ob_partition_merge_iter.h):打开各参与 sstable,按 rowkey 顺序吐行。ob_partition_merge_fuser.*):把同行多个版本合并成「最终有效版」,处理 uncommitted row(来自 MemTable)和 multi-version 行。ob_partition_merger.*):按行输出 macro block + micro block,通过 IndexBlockBuilder 构建索引。i_compaction_filter.h 的 ObICompactionFilter。compaction/filter/ + DDL filter staleness filter + TTL filter + MDS filter。ObProgressiveMergeHelper(ob_progressive_merge_helper.h):当 Major 数据量巨大、一刀切会打垮 IO 时,按 macro block 「轮转合并」,由 progressive_merge_round 控制剩余未合并 range,逐 round 推进。INC_MAJOR_MERGE 类型即应用此机制。
dag_status_manager 报告并指数退避重试(ob_schedule_status_cache)。TabletPointer switch)。minor_compact_trigger(2~3)就 minor;Minor 数不会无限堆叠。MediumCompactionInfo(snapshot_version + schema_version + medium_type 等)。所有副本 Major/Medium 不可超过 medium_watermark;完成后写相同的 extra_medium_info。因此多副本 Major 完成后产出的新 SSTable 的 rowkey 范围、版本谱、schema_version 三者一致;合并输出的字节流上方相等(同 macro block hash 又可验)。ob_dag_macro_block_writer)。Scheduler 为了吞吐会在不同 tablet 间并发。代码位置:src/storage/tx 与 src/storage/tx_storage。核心类:
ObTransService(ob_trans_service.h)/ V4 ob_trans_service_v4:向 SQL 层提供 begin/commit/rollback。ObTxDesc:事务描述符(客户端视角事务)。ObPartTransCtx(ob_trans_part_ctx):某 LS 上该事务的参与者上下文。ObTxCtxMgrV4:每个 LS 一个 ctx 管理器。ObITSManager / GTS / Standby Timestamp:分布式读时间戳 / 全局单调时间戳服务(ob_ts_mgr、ob_gts_*、ob_timestamp_service)。ObMvccRow(memtable/mvcc/ob_mvcc_row),下行链表为本 row 的多个版本,每个版本带 trans_version(事务提交 SCN)。snapshot / SCN,决定可见版本;MemTable 索引用 KeyBTree(memtable/mvcc/ob_keybtree)。ob_row_conflict_handler + ob_row_latch(行级 latch)+ write-write 冲突回滚。ob_tx_elr_handler。2PC 角色与状态见 src/storage/tx/ob_committer_define.h:
enum class Ob2PCRole { UNKNOWN=0, ROOT, INTERNAL, LEAF }; // 树形分布式事务
enum class ObTxState : uint8_t {
UNKNOWN=0, INIT=10, REDO_COMPLETE=20, PREPARE=30,
PRE_COMMIT=40, COMMIT=50, ABORT=60, CLEAR=70, MAX=100
};
enum class ObTwoPhaseCommitLogType { OB_LOG_TX_INIT, OB_LOG_TX_COMMIT_INFO,
OB_LOG_TX_PREPARE, OB_LOG_TX_PRE_COMMIT, OB_LOG_TX_COMMIT, OB_LOG_TX_ABORT, OB_LOG_TX_CLEAR };
2PC 协调器核心:ObTwoPhaseCommitter(storage/tx/ob_two_phase_committer.h),与上下游 commit 链
ob_two_phase_upstream_committer / ob_two_phase_downstream_committer 形成多级 commit tree,
IDA 的「分布式事务」形态通过 ob_tx_2pc_ctx_impl、ob_tx_2pc_msg_handler 消息驱动。
ob_redo_log_generator 流式提交,织入 PaLF(ob_tx_redo_submitter.h)。ob_multi_data_source.*):把 DDL/DUMP/Lob/外挂消息等跨模块信息以 commit 阶段离线写入 commit_info log(OB_LOG_TX_COMMIT_INFO),保证事务与外部副作用一致。ob_tx_replay_executor 把提交信息回放到 follower 上,支持 StandbyReadable。跨多个 LS 的分布式协调不是直接的 N-方 2PC,而是按依赖关系组织成 commit tree:
Ob2PCRole 分 ROOT / INTERNAL / LEAF,由 root 服务(通常事主参与 LS)作为 coordinator;
每个分发到的 LS 都是一个 downstream committer;上游者应答前驱 agnostic。ObTwoPhaseCommitter 同时
向上游发送 PRE_COMMIT / COMMIT / ABORT,在 ob_two_phase_upstream_committer/ob_two_phase_downstream_committer
上完成树形级联。这样可以樹形隐藏 single-coordinator bottleneck,时延近似 O(log L),避免某分片持有 participant 名称表。
源码:ob_tx_free_route.*。允许同一逻辑事务在物理 多个 OBServer 节点之间路由(即便事务尚未提交,
客户端 SQL 也可被 ObProxy 路由到任意 server 执行 next part),而事务状态不再严格"绑在某一 server 上"。状态机
TxFreeRouteState:SAVEPOINT / COMMIT_INFO 等通过 commit_info log + 内存 session state snapshot 在 server 间传递。
这是 OLTP 负载均衡与秒级切关键。
ob_tx_sby_read_ctx_helper / ob_standby_timestamp_service 维护可读 SCN。ob_weak_read 路径)。storage/tx/deadlock_adapter)等待样在行锁图谱上跑 DFS;检测到 cycle 则 abort 一个 victim。ob_xa_*.{h,cpp} 支持 XA 协议(含 dblink)。ob_gti_source.h、ob_trans_id_service:租户级全局唯一 ID service,基于 clog 推进。ob_gts_source 提供 monotone counter,作为快照 SCN 来源;多副本备份。enum ObTxState(10~70)。PRE_COMMIT 是 OB 加的优化阶段,让 leader 端的多数派持久化"提 prepare"后告知下游,减少 COMMIT 阶段 RTT。retransmit_upstream_msg_);root 若已发布 COMMIT 则下发 CLAIM;若严格未发布、且 timeout,会按 PREPARE emission 时间决定 ABORT/COMMIT:通常 commit_info 已包含 redo/结果 → 不会丢;保护由 retain_ctx_mgr(ob_tx_retain_ctx_mgr)保留结构直到 commit/abort 完成。ObTxDataTable::check_tx_status 验证目标 trans 是否已 commit 与是否 <= snapshot_version;同时脏读未被可见 access 限制(每次读到 ObMvccRow 时检查被事务是否在 commit_version 中"确定状态",使用 TxData 临时把 committed/unfinished 分类)。ObPartTransCtx 持有事务态、回叙并旃;当客户端路由到新节点、从 commit_info log 中恢复了主要 participant 中的 snapshot。强一致性由 clog + 2PC + commit info-log 保证,不在 server 内存里维持"transactor host"。ObRowLatch 在被多事务并发更新同 row 时 serialize,影响仅限冲突同 row。deadlock_adapter 通过 "wait-for" graph(加锁链表)定期发现有向环;同时强超时(ob_tx_timeout)自毫 fault-injection 与高 trauma 又 印证。PaLF = Paxos Log Facility(src/logservice/palf),是 OceanBase 自研的、面向 LS 的、基于
Multi-Paxos + Lease-based 选举的日志复制引擎。每个 LS 自己有一组 PaLF 实例。它不是 Fast Paxos
(Fast Paxos 需要 2F+1 全部参与 prepare 阶段且不需 Leader),OceanBase 使用 Leader-based Multi-Paxos:
只有 Leader 接读写,其余副本只持久化和回放。Leader 通过独立的选举模块(palf/election) lease 选出。
| 模块 | 路径 | 职责 |
|---|---|---|
| PalfEnv / PalfHandle | palf_env_impl, palf_handle_impl | 管理一组 LS 的 PaLF 实例及句柄 |
| LogStateMgr | log_state_mgr.{h,cpp} | 维护本副本的 (role, state) 状态机并驱动切换 |
| LogConfigMgr | log_config_mgr.{h,cpp} | 成员变更(加/减副本、降级、仲裁、切主) |
| LogSlidingWindow | log_sliding_window.{h,cpp} | 提案编号、leader 上的 ack 滑窗、follower 接收点 |
| LogEngine | log_engine.{h,cpp} | 阻塞 IO、checksum、落盘、回收、网络收发 |
| LogModeMgr | log_mode_mgr | 访问模式切换(RAW_WRITE / APPEND) |
| LogReconfirm | log_reconfirm.{h,cpp} | 新主上线时把多数派未一致的日志拉全、补齐 |
| Election | palf/election/ | 独立选举模块(非 Paxos 缺一不可的部分) |
| ApplyService / ReplayService | logservice/applyservice, replayservice | 把已提交日志应用到存储引擎 |
角色 ObRole:LEADER / FOLLOWER;副本状态 ObReplicaState(log_define.h):
enum ObReplicaState {
INVALID_STATE = 0, INIT = 1, ACTIVE = 2, RECONFIRM = 3, PENDING = 4,
};
状态转移节选自 log_state_mgr.cpp:
LogReconfirm::State(log_reconfirm.h)。这是 OceanBase Multi-Paxos 重要扩展——确保新 Leader 在服务写之前,
把多数派未确认的日志全部补齐,这是不同 Paxos 组从「无主 / 切主 / 宕机」直奔稳定阶段的恢复路径:
log_config_mgr 支持。源码:log_config_mgr.cpp 的 LogConfigChangeType 包括 ADD/REMOVE/CHANGE_NUM/DEGRADE/UPGRADE/ADD_ARB/REMOVE_ARB/CHANGE_LEADER/SETR_LOGONLY/...。一次变更分四步:
LogReconfirm(log_reconfirm.{h,cpp})做新主补全多副本不一致日志。关键变量:
max_lsn_:本机当前最大已 flush LSN。majority_max_log_server_:在多数派中具有最大 lsn 的 server。fetching_lsn_:正在从多数派拉取的目标 LSN。RECONFIRM_PREPARE_RETRY_INTERVAL_US=2s。进程:选举胜选 → 本机 flush 完毕 → 向多数派发 prepare → 找出 lsn 最大者 → 拉 (max_lsn_local, majority_max_lsn] 区间日志 → 写 START_WORKING log(多 raft 唯一标志)→ 切 LEADER_ACTIVE。 START_WORKING commit 后才允许对该 LS append 新日志——这是 Multi-Paxos "断主重主不准漏日志" 的关键实现。
ob_replay_handler),follower 一直 hot,不需要冷启动。实例:3F normal 切主一般在 1~3s 完成 START_WORKING;整 RTO 阈值 8s 留 buffer 给网络抖动 + brick 恢复。
OB 把多条 redo log entry 合成 LogGroupEntry(log_group_entry)批量落盘,减少 IO 次数。Sliding Window(log_sliding_window)维护 ack 位图:leader 维护每 follower 的 acked lsn;当"多数派 ack lon >= entry 的 lsn" 则该 entry committed 走 ApplyService。
Multi-Paxos 本身只是协议,需要「谁来做 Leader」这一前提。OceanBase 把选举与共识拆开:选举模块
(src/logservice/palf/election)用 Lease 选举每 LS 独立维护一个稳定 Leader,
Paxos 只在 Lease 失效或切主时被重新启动(见 LogStateMgr::switch_state)。
election/interface/election.h):
DevoteToBeLeader(无主选举登基)、ChangeLeaderToBeLeader(切主登基)、LeaseExpiredToRevoke(租约过期退位)、
ChangeLeaderToRevoke(切主退位)、StopToRevoke(停服退位)。关键代码:
palf/election/algorithm/election_proposer.(发起 Prepare/Accept)。election_acceptor(应答 Prepare/Accept、记录 promised_proposal_id)。palf/election/message/election_message:Prepare/Accept 请求与应答,带优先级(election_priority)。temporarily_downgrade_protocol_priority:短暂降低本副本优先级(如刚恢复但日志较少)。Leader 持有严格租约,租约期内不开新一轮选举;租约即将到期时由 Leader 自我续约(lease renew),
减少选举震荡。租约过期后由 Follower 触发新一轮。RCHandle、RequestChecker 处理过期与消息时序。
ElectionPriority(palf/election/interface/election_priority.h)的比较维度(从高到低):
log_lsn / proposal_id(日志更全者优先)因此在切主时原则上由日志最全者当选,避免后续 leader 还要回拉 follower 日志的代价。temporarily_downgrade_protocol_priority 提供外部介入:刚恢复的 server 短时下调自己优先级,避免它当选后立刻被更全者"切回来"。
索引定义见 src/share/schema/ob_schema_struct.h(ObIndexType)。OceanBase 的二级索引并不直接坐落
B+Tree 数据结构上,而是「索引表」——每张索引表本身就是一个 Tablet,存储布局与数据表一致(LSM-Tree)。
用户表与索引表之间通过 rowkey 维持对应关系。
| 大类 | 典型 index_type | 存储形态 | 用途 |
|---|---|---|---|
| 本地索引(local) | IS_LOCAL/UNIQUE_LOCAL | 主表 Tablet 内的辅助 SSTable | 同一 Tablet 范围的快速定位 |
| 全局索引(global) | IS_GLOBAL/UNIQUE_GLOBAL | 独立 Tablet,跨分区分布 | 跨分区查询、不绑定主表分区 |
| 向量索引 IVF | INDEX_TYPE_VEC_IVFFLAT_CENTROID_LOCAL / CID_VECTOR_LOCAL / ROWKEY_CID_LOCAL | 多张辅助 Tablet:质心表 + 倒排 | IVF_FLAT / IVF_SQ8 向量 ANN |
| 向量索引 HNSW | vec_hnsw_*(见 schema_struct.h:429) | 列存 HNSW 图数据 | 近似向量检索(高召回) |
| 全文检索 FTS | is_local/global/global_local_fts_index*、fts_doc_word_aux | doc-word 倒排辅助表 | 全文 MATCH AGAINST |
| 多值索引 | is_multivalue_index* | 多值列倒排表 | JSON 数组成员查询 |
| Domain 索引 | domain_index 相关 | 插件可扩展 | 插件化自定义索引 |
每张 SSTable 由 Index Block 树(多级 + 微块)组织,源码在 src/storage/blocksstable/index_block:
索引扫描核心组件(src/storage/access):
ob_index_block_tree_traverser:B+ 树遍历。ob_index_tree_prefetcher:预取优化。ob_index_skip_scanner + ob_skip_index_sortedness:Skip Scan(跳过不匹配段)。ob_sstable_index_filter:下推到 SSTable 的索引层过滤(列裁剪 + 谓词)。位于 src/sql/das/ob_das_vec_define、ob_vector_index_lookup_op、ob_das_search_index_utils、
src/sql/optimizer/ob_access_path_estimation。IVF 系列:
Centroid Local(质心)+ CID-Vector Local(每个质心所属向量)+ Rowkey-CID Local(反向映射),3 张辅助 Tablet 联合检索;
HNSW:单独的图结构辅助表。SQL 层支持 VECTOR 数据类型、距离函数(L2/IP/Cosine)与混合检索
(sql/code_generator/ob_hybrid_search_cg_service_*)。
位于 src/storage/fts、sql/code_generator/ob_hybrid_search_cg_service_fulltext。结构:分词后的
doc-word 倒排辅助表(is_fts_doc_word_aux),查 MATCH AGAINST 时复用 parser 分词结果做倒排合并。支持 local/global/global-local 三种形态。
用于 JSON / ARRAY 列:把数组每个元素展开到一张倒排辅助表,使 MEMBER OF、JSON_CONTAINS 等可以走索引。
ObAccessPathEstimation 计算时判断 select/where 列是否全部在索引前缀 + 主键中;是则可避免回表 lookup cost。ob_index_skip_scanner、ob_skip_index_sortedness):当前缀列基数小但取值离散时,对每个前缀值都做 1 次 range scan,等效消除前缀谓词。代价模型上要算多次 micro block 开销,CBO 决定是否启用。ob_sstable_index_filter):能在索引块上解出的谓词 → 直接过滤;其余谓词 push 到 row-level(push down filter)。ObOptStatManager 的直方图,索引列同样采直方图 + density;global/local 索引采用 global stats.ob_clustered_index_block_writer):与数据 macro block 同构存储,将某索引列集和数据内容物理交错,减少回表。ob_multiple_merge 中按列需求拼接:AP 查询列打开 CG,要小行 rowkey lookup 仍取 minor 行存。Major 切换保证 schema 一致。ObDASScanOp 抽象 + ObDASExtraData 适配不同索引类型:FTS 用 hybrid_search_cg、向量用 ob_vector_index_lookup_op,多值用倒排 task。SQL 算子侧对外仍是 TableScan / SubPlan Scan,专门索引插.min scan_descriptors(ob_das_def_reg)注册以避免上层重新发明算子。RootService(src/rootserver)是一个集群级单点服务(实际上也是根 LS 的 Leader 副本承担),负责:
ob_bootstrap)、server 上下线、心跳(ob_heartbeat_service)、ob_all_server_checker。ob_ddl_service、ob_ddl_operator、ob_index_builder、ob_lob_meta_builder、catalog(ob_catalog_ddl_*)、AI model(ob_ai_model_ddl_*)。balance/、ob_root_balancer、ob_partition_balance、ob_balance_ls_primary_zone:Unit 与 LS 在 zone 间迁移。ob_disaster_recovery_*、教师副本切换、补副本(ob_lost_replica_checker、ob_empty_server_checker)。freeze/ob_root_minor_freeze、compaction_ttl/:调度 Major / Medium。backup/。mview/ 与 ob_mlog_builder。ddl_task/、direct_load/。| 主题 | 关键文件 / 目录 |
|---|---|
| 进程入口 | src/observer/main.cpp、ob_server.{h,cpp}、ob_srv_network_frame、ob_srv_deliver/xlator |
| SQL 入口 | src/sql/ob_sql.{h,cpp}、ob_result_set |
| Parser | src/sql/parser/(ob_sql_parser.l/y、raw_token.h) |
| Resolver | src/sql/resolver/(ob_resolver.h、dml/ddl/cmd/tcl/dcl 子目录) |
| Rewrite | src/sql/rewrite/ob_transform_rule.*、udr/ |
| Optimizer | src/sql/optimizer/ob_optimizer.h、ob_log_plan.h、ob_join_order_enum_idp |
| Code Generator | src/sql/code_generator/ob_static_engine_cg.{h,cpp}、ob_expr_generator_impl |
| 执行算子 | src/sql/engine/*(ob_operator.h、ob_operator_factory、basic/join/sort/aggregate/px/dml) |
| DAS | src/sql/das/ob_data_access_service.h、ob_das_scan_op.h、ob_das_location_router.h |
| PX | src/sql/engine/px/ob_dfo.h、ob_dfo_scheduler、ob_granule_pump、ob_px_coord |
| 存储访问 | src/storage/access/ob_multiple_merge、ob_index_block_tree_traverser |
| MemTable / MVCC | src/storage/memtable/ob_memtable.h、mvcc/ob_keybtree.h、mvcc/ob_mvcc_engine.h |
| SSTable | src/storage/blocksstable/ob_block_manager、ob_micro_block_*、index_block/ob_index_block_builder |
| 列式编码 | src/storage/blocksstable/encoding/*、cs_encoding/* |
| Compaction | src/storage/compaction/ob_compaction_util.h、ob_partition_merger、ob_partition_merge_fuser、ob_medium_compaction_mgr |
| Tablet / LS | src/storage/ls/、src/storage/tx_storage/ob_ls_service/ob_ls_map、src/storage/meta_mem/ |
| 事务 | src/storage/tx/ob_trans_service.h、ob_trans_part_ctx.h、ob_two_phase_committer.h、ob_committer_define.h |
| PaLF | src/logservice/palf/:palf_env_impl、palf_handle_impl、log_state_mgr、log_config_mgr、log_sliding_window、log_reconfirm、log_engine |
| Apply / Replay | src/logservice/applyservice/ob_log_apply_service、replayservice/ob_log_replay_service、ob_replay_handler |
| 选举 | src/logservice/palf/election/algorithm/:election_impl、election_proposer、election_acceptor;interface/election.h |
| 归档 / CDC | src/logservice/archiveservice/、cdcservice/、libobcdc/ |
| RootService | src/rootserver/ob_root_service.h、ob_ddl_service.h、ob_root_balancer、ob_bootstrap |
| Schema | src/share/schema/(ob_schema_struct.h:ObIndexType 定义) |
| FTS | src/storage/fts/、sql/code_generator/ob_hybrid_search_cg_service_fulltext.cpp |
| 向量索引 | src/sql/das/ob_das_vec_define、ob_vector_index_lookup_op、sql/code_generator/ob_hybrid_search_cg_service_vec.cpp |
ObSql::stmt_query() 跟到 ObPhysicalPlan,再进入
ObTableScanOp → ObDASScanOp → ObMultipleScanMerge → MemTable/SSTable 算子,
同时从 ObTransService::start_tx 跟到 ObPartTransCtx::submit_log → PalfHandleImpl::append,
即可贯通 SQL / Storage / Tx / Paxos 四层。