OceanBase源码解读:存储引擎内存表 ObMemtable

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 &param,
                    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 则根据是否为多版本小合并场景,分别分配 ObMemtableMultiVersionScanIteratorObMemtableScanIterator,由 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、事务、日志回收等模块的重要基础。

发表回复

您的邮箱地址不会被公开。 必填项已用 * 标注