OceanBase 的存储引擎采用 LSM-Tree 架构,所有写入首先进入内存表(Memtable),再按一定策略冻结并转储为 SSTable。ObMemtable 是单个 Tablet 的内存表实体,向上承接 SQL 执行层与事务层的 set/get/scan/lock/replay 请求,向下通过 MVCC 引擎、Query 引擎和 Memstore 分配器完成并发控制、索引维护与内存管理。本文选取其初始化、读写路径与冻结/刷盘流程做重点解读。
核心数据结构
- ObSingleMemstoreAllocator local_allocator_:事务级无锁内存分配器,Memtable 冻结后停止新写入,所有已分配内存随 Memtable 生命周期一起回收。
- ObQueryEngine query_engine_:Hash 表 + B+Tree 双索引,Hash 用于单行快速定位,B+Tree 保证范围扫描有序。
- ObMvccEngine mvcc_engine_:MVCC 并发控制引擎,维护
ObMvccRow/ObMvccTransNode版本链,实现读不阻塞写、写写冲突检测。 - ObMemtableKey / ObMemtableSetArg:rowkey 编码与 DML 参数包,封装了新行、旧行、列描述、update 列索引等信息。
- ObReportedDmlStat:定期向
ObOptStatMonitorManager上报增删改行数,用于优化器统计信息。
初始化:三大引擎的装配
ObMemtable::init 是典型的“错误码瀑布”风格:按依赖顺序依次检查参数、绑定 MemtableMgr/Freezer、初始化分配器/Query 引擎/MVCC 引擎、设置 table_key 与 LS 句柄,最后将状态置为 ACTIVE。任意步骤失败都会调用 destroy() 清理已分配资源。
int ObMemtable::init(const ObITable::TableKey &table_key,
ObLSHandle &ls_handle,
storage::ObFreezer *freezer,
storage::ObTabletMemtableMgr *memtable_mgr,
const int64_t schema_version,
const uint32_t freeze_clock)
{
int ret = OB_SUCCESS;
if (is_inited_) { ret = OB_INIT_TWICE; }
else if (!table_key.is_valid() || OB_ISNULL(freezer) ||
OB_ISNULL(memtable_mgr) || schema_version get_ls_id();
mode_ = table_key.get_tablet_id().is_sys_tablet()
? lib::Worker::CompatMode::MYSQL
: MTL(lib::Worker::CompatMode);
state_ = ObMemtableState::ACTIVE;
init_timestamp_ = ObTimeUtility::current_time();
set_freeze_state(TabletMemtableFreezeState::ACTIVE);
is_inited_ = true;
}
if (OB_SUCCESS != ret && IS_NOT_INIT) { destroy(); }
return ret;
}
写入入口 set:从 DML 到 MVCC
ObMemtable::set 是单条 insert/update/delete 的公共入口。它先通过 ObMvccWriteGuard 校验写入权限,再用 ObMemtableKeyGenerator 把 schema rowkey 列编码为 ObMemtableKey,随后进入内部 set_ 链路。CompatModeGuard 保证 MySQL/Oracle 语义隔离;写入成功后还会异步上报 DML 统计。
int ObMemtable::set(const storage::ObTableIterParam ¶m,
storage::ObTableAccessContext &context,
const ObMemtableSetArg &arg)
{
int ret = OB_SUCCESS;
ObMvccWriteGuard guard(ret);
const blocksstable::ObDatumRow *new_row = arg.new_row_;
const ObIArray *columns = arg.columns_;
if (IS_NOT_INIT) { ret = OB_NOT_INIT; }
else if (!param.is_valid() || !context.is_valid() || !arg.is_valid()) {
ret = OB_INVALID_ARGUMENT;
} else if (OB_FAIL(guard.write_auth(*context.store_ctx_))) {
} else {
ObMemtableKeyGenerator key_gen(param.get_schema_rowkey_count(), *columns);
if (OB_FAIL(key_gen.init())) {
} else if (OB_FAIL(key_gen.generate_memtable_key(*new_row))) {
} else {
lib::CompatModeGuard compat_guard(mode_);
ret = set_(param, context, arg, key_gen.get_memtable_key());
guard.set_memtable(this);
}
}
if (OB_SUCC(ret)) {
int tmp_ret = OB_SUCCESS;
if (OB_TMP_FAIL(try_report_dml_stat_(param.table_id_))) {
TRANS_LOG_RET(WARN, tmp_ret, "fail to report dml stat");
}
}
return ret;
}
MVCC 原子写:mvcc_write_
所有写操作最终汇聚到 mvcc_write_。它分三步完成一次原子写入:先在 hash 表中创建或复用 ObMemtableKey → ObMvccRow;再由 MVCC 引擎在版本链上挂入新节点并检测写写冲突、主键重复、事务集违反;最后把 key/value 插入 b+tree 以保证扫描有序。若中间失败则调用 mvcc_undo 回滚,确保接口原子性。
int ObMemtable::mvcc_write_(ObStoreCtx &ctx,
const ObMemtableKey &memtable_key,
const ObTxNodeArg &tx_node_arg,
const bool check_exist,
ObMvccWriteResult &res)
{
int ret = OB_SUCCESS;
ObMemtableKey stored_key;
ObMvccRow *value = NULL;
ObMemtableCtx *mem_ctx = ctx.mvcc_acc_ctx_.get_mem_ctx();
// 1. 创建或复用 hash 表中的 key / value
if (OB_FAIL(mvcc_engine_.create_kv(&memtable_key,
blocksstable::ObDmlFlag::DF_INSERT == tx_node_arg.data_->dml_flag_,
&stored_key, value))) {
TRANS_LOG(WARN, "create kv failed", K(ret));
} else {
res.mtk_.encode(stored_key);
res.value_ = value;
}
// 2. MVCC 写入:检测冲突并挂新版本
if (OB_SUCC(ret) && OB_FAIL(mvcc_engine_.mvcc_write(ctx, *value,
tx_node_arg, check_exist, res))) {
if (OB_TRY_LOCK_ROW_CONFLICT == ret) {
mem_ctx->on_wlock_retry(memtable_key, res.lock_state_.lock_trans_id_);
ret = post_row_write_conflict_(ctx.mvcc_acc_ctx_, memtable_key,
res.lock_state_,
value->get_last_compact_cnt(),
value->get_total_trans_node_cnt());
} else if (OB_TRANSACTION_SET_VIOLATION == ret) {
mem_ctx->on_tsc_retry(memtable_key, ctx.mvcc_acc_ctx_.snapshot_.version(),
value->get_max_trans_version(),
value->get_max_trans_id());
}
// 3. 确保 key/value 进入 b+tree
} else if (OB_SUCC(ret) && OB_FAIL(mvcc_engine_.ensure_kv(&stored_key, value))) {
TRANS_LOG(WARN, "prepare kv after lock fail", K(ret));
}
// 4. 失败回滚
if (OB_FAIL(ret) && res.has_insert()) {
(void)mvcc_engine_.mvcc_undo(value);
res.is_mvcc_undo_ = true;
}
return ret;
}
读取入口 get / scan
点查 get 将 rowkey 编码为 ObMemtableKey,调用 mvcc_engine_.get 获取版本链,再通过 ObReadRow::iterate_row 把所需列投影到 ObDatumRow。若遇到行锁冲突,由 ObRowConflictHandler 进入锁等待或向上层返回冲突码。
范围扫描 scan 则根据是否为多版本小合并场景,分别分配 ObMemtableMultiVersionScanIterator 或 ObMemtableScanIterator,由 query_engine_ 的 b+tree 做有序遍历,MVCC 引擎负责版本可见性判断。
冻结与刷盘:从内存到 SSTable
Memtable 写满或日志回收需要时会触发冻结。is_frozen_memtable 比较 LS 级 freeze_clock 与本 Memtable 的 freeze_clock 判断是否已被冻结。ready_for_flush_ 则进一步检查:写引用归零、未提交事务数归零、当前 right_boundary 已推进到 max_end_scn 之后、左边界已解析。满足条件后状态变为 READY_FOR_FLUSH。
bool ObMemtable::ready_for_flush_()
{
bool is_frozen = is_frozen_memtable();
int64_t write_ref_cnt = get_write_ref();
int64_t unsubmitted_cnt = get_unsubmitted_cnt();
bool bool_ret = is_frozen && 0 == write_ref_cnt && 0 == unsubmitted_cnt;
SCN current_right_boundary = ObScnRange::MIN_SCN;
if (bool_ret) {
if (OB_FAIL(resolve_snapshot_version_())) {
} else if (OB_FAIL(resolve_max_end_scn_())) {
} else if (OB_FAIL(get_ls_current_right_boundary_(current_right_boundary))) {
} else if (current_right_boundary >= get_max_end_scn()) {
resolve_right_boundary();
bool_ret = get_resolved_active_memtable_left_boundary();
}
if (bool_ret) {
set_freeze_state(TabletMemtableFreezeState::READY_FOR_FLUSH);
}
}
return bool_ret;
}
一旦可刷盘,flush 会构造一次 MINI_MERGE 后台任务,把 Memtable 数据落盘为 Mini SSTable,并填充 occupy_size、replay_interval、last_end_scn 等调度参数。最后 finish_freeze 固化右边界与转储点元数据,完成整个 freeze 生命周期。
int ObMemtable::flush(share::ObLSID ls_id)
{
int ret = OB_SUCCESS;
if (get_is_flushed()) { return OB_NO_NEED_UPDATE; }
ObTabletMergeDagParam param;
param.ls_id_ = ls_id;
param.tablet_id_ = key_.tablet_id_;
param.merge_type_ = MINI_MERGE;
param.merge_version_ = ObVersion::MIN_VERSION;
fill_compaction_param_(ObTimeUtility::current_time(), param);
if (OB_FAIL(compaction::ObScheduleDagFunc::schedule_tablet_merge_dag(param))) {
if (OB_EAGAIN != ret && OB_SIZE_OVERFLOW != ret) {
TRANS_LOG(WARN, "failed to schedule tablet merge dag", K(ret));
}
}
return ret;
}
小结
ObMemtable 是 OceanBase 写入路径的必经之地,也是 LSM-Tree 内存侧的核心抽象。它通过 Memstore 分配器 管理事务内存,通过 Query 引擎 维护 Hash + B+Tree 双索引,通过 MVCC 引擎 实现并发控制与版本链,再通过 freeze/flush 机制把内存数据有序地转储到磁盘。理解它的生命周期与读写链路,是理解后续 compaction、事务、日志回收等模块的重要基础。