OceanBase源码解读:物化视图刷新 ob_mview_refresh.cpp

物化视图(Materialized View)是 OceanBase 4.3 的旗舰特性之一:把一条查询 SQL 的结果实体化存储成一张真实表,查询时直接读表,免去每次实时聚合计算。但实体化带来一个经典问题——基表数据一直在变,物化视图里的”快照”怎么跟上?答案就是本篇的主角 ob_mview_refresh.cpp(约 1000 行):物化视图的刷新执行器 ObMViewRefresher,负责把 MV 的数据从上次刷新点安全地推进到当前时点。它提供两条刷新路径:FAST 增量刷新(消费基表 MLog 变更日志,只算差量)和 COMPLETE 全量刷新(把定义 SQL 重新算一遍),并提供自适应选路逻辑在两者间智能切换。

一、核心数据结构:SCN 是一切的锚点

理解刷新逻辑前先认识三个关键角色。整个文件的主线,其实是围绕”SCN 区间”展开的版本管理。

  • ObMViewInfo:持久化在内表里的 MV 元数据,记录 last_refresh_scn(上次刷到哪个版本)、data_sync_scn(嵌套一致刷新的对齐版本)、刷新类型与耗时等。它是刷新的”账本”。
  • ObScnRange:[start_scn, end_scn] 版本区间。FAST 刷新的取数窗口就是它——只消费落在这个区间内的 MLog 变更。
  • MLog:基表上的物化视图日志,一张自动维护的内表。基表每次 DML 都会往里追加变更记录,是增量刷新的数据源。每条 MLog 有清理水位 last_purge_scn。

其中 SCN(System Commit Number)是 OceanBase 的提交版本号,MV 数据本质上就是”某个 SCN 快照下的查询结果”,所以刷新的全部一致性保障都建立在 SCN 对齐之上。

二、主流程:refresh() 的四步走

入口函数 refresh() 的骨架非常清晰:

int ObMViewRefresher::refresh()
{
  // ① 加锁:对 MV 元数据记录 try_lock,防并发刷新
  if (OB_FAIL(lock_mview_for_refresh())) {
  } else if (OB_FAIL(prepare_for_refresh())) {
    // ② 决策:收集元数据、算 SCN 区间、定 FAST 还是 COMPLETE
  } else {
    if (ObMVRefreshType::FAST == refresh_type) {
      // ③a 增量路径:先向 MDS 注册 mview 操作记录(供嵌套 MV 感知)
      arg.mview_op_type_ = MVIEW_OP_TYPE::FAST_REFRESH;
      arg.read_snapshot_ = refresh_ctx_->mview_refresh_scn_range_.end_scn_...;
      ObMViewMdsOpHelper::register_mview_mds(...);
      fast_refresh();
    } else if (ObMVRefreshType::COMPLETE == refresh_type) {
      complete_refresh();  // ③b 全量路径:转发 RootService 重算
    }
  }
  // ④ 刷新后收集统计信息,供下次自适应选路使用
  if (OB_SUCC(ret) && nullptr != refresh_stats_collection_) {
    refresh_stats_collection_->collect_after_refresh(*refresh_ctx_);
  }
  // OB_ERR_TASK_SKIPPED 是"已到目标版本、无需再刷"的正常跳过,吞掉转成功
  if (OB_ERR_TASK_SKIPPED == ret) {
    ret = OB_SUCCESS;
  }
}

加锁细节值得一提:lock_mview_for_refresh() 用 try_lock 模式拿不到锁就 sleep 100ms 循环重试,每轮都调用 check_status() 响应 kill 请求——并发刷新同一 MV 的任务会排队,但会话不会被卡死在锁上。

三、决策大脑:prepare_for_refresh() 怎么选路

prepare_for_refresh() 是全文件最复杂的函数,做了五件事:取当前 SCN、校验 MV schema、拉取持久化的 mview_info、计算 SCN 区间,最后决定刷新类型。选路逻辑如下:

// 用户指定 FAST / FORCE / FORCE_AUTO,结合可行性判定
if (ObMVRefreshMethod::COMPLETE == refresh_method ||
    (!can_fast_refresh && ObMVRefreshMethod::FORCE == refresh_method)) {
  refresh_type = ObMVRefreshType::COMPLETE;
} else if (!can_fast_refresh && ObMVRefreshMethod::FAST == refresh_method) {
  // 明确要增量但增量不了,直接报错
  ret = OB_ERR_MVIEW_CAN_NOT_FAST_REFRESH;
} else if (OB_FAIL(check_fast_refreshable_(...))) {
  // FORCE 模式下增量失败可降级全量:依赖变更 / MLog 过期都不算错误
  if (OB_ERR_MVIEW_MISSING_DEPENDENCE == ret) {
    refresh_type = ObMVRefreshType::COMPLETE;
    ret = OB_SUCCESS;
  }
} else {
  refresh_type = ObMVRefreshType::FAST;
}

can_fast_refresh 由 MVProvider 探测,背后是 check_fast_refreshable_() 的三连检:

  1. 依赖数量一致——MV 定义没被重建过;
  2. 依赖对象 ID 逐个比对——基表删了重建 ID 会变,增量窗口随之断裂;
  3. MLog 清理水位检查——MLog 的 last_purge_scn 不能比上次刷新点更新。MLog 只保留最近的数据,若清理水位已越过上次刷新 SCN,中间的变更就丢了,只能全量补。

SCN 区间计算在 calc_scn_range() 中完成,两个区间的语义不同:MV 自己的推进窗口是 [last_refresh_scn, current_scn],而读基表 MLog 的窗口左边界在嵌套一致刷新场景下要用 data_sync_scn(通常更新),避免重复消费变更。

四、两条路径的执行细节

FAST 增量路径fast_refresh())全程在一个事务里,失败整体回滚:

// ① 版本对账:回读 mview_info,确认 last_refresh_scn 未被并发刷新挪动
// ② 清掉 MLog 相关的动态采样统计缓存(即将被消费,旧采样失效)
// ③ 抬高 DML 并行度(Oracle/MySQL 语法不同)
sql.assign_fmt("SET _force_parallel_dml_dop = %lu", parallelism);
// ④ 核心循环:逐条执行 MVProvider 预生成的内部 SQL
for (int64_t i = 0; OB_SUCC(ret) && i < refresh_sqls.count(); ++i) {
  const ObString &fast_refresh_sql = refresh_sqls.at(i);
  trans.write(tenant_id, fast_refresh_sql.ptr(), affected_rows);
  // 每条 SQL 记录耗时统计
}
// ⑤ 复查 MLog 未被中途 purge(清理水位不能越过本次刷新点)
// ⑥ 推进账本:last_refresh_scn = end_scn,记录刷新类型/耗时/trace_id

值得注意的设计是:增量刷新的 SQL 是 prepare 阶段由 MVProvider 预生成的一组内部 INSERT/UPDATE/DELETE 语句(数据源就是各基表 MLog 在 SCN 区间内的变更),执行器只负责逐条跑,把”怎么合并差量”和”什么时候跑”解耦。

COMPLETE 全量路径complete_refresh())很有意思:它并不在本节点重算,而是组装 RPC 参数(连 NLS 时区、日期格式等会话环境都带上)发给 RootService,复用 DDL 框架异步执行,然后轮询 task_id 等待完成。设计动机很直接——整表重写是重资源操作,交给 RS 统一调度 DAG 并行任务,失败重试和进度观测都能复用 DDL 的既有基础设施。

五、自适应选路:差量大了就别增量了

4.3.5 还引入了一个聪明的启发式 check_adaptive_refresh_method()

const int64_t MLOG_ROWS_THRESHOLD = 500;
// 遍历刷新前的变更统计
int64_t num_mlog_rows = change_stats.num_rows_ins_ + change_stats.num_rows_del_
                     + change_stats.num_rows_upd_;
int64_t num_base_table_rows = change_stats.num_rows_ - change_stats.num_rows_ins_
                           + change_stats.num_rows_del_;
if (num_mlog_rows > MLOG_ROWS_THRESHOLD
    && static_cast(num_mlog_rows)/num_base_table_rows
       > complete_refresh_ratio_threshold) {
  should_adaptive_complete = true;  // 改判为 COMPLETE
}

当某基表的变更行数超过 500 行、且变更量占基表行数的比例超过租户配置的阈值(_mv_adaptive_complete_refresh_threshold)时,自动改走全量。这是典型的代价感知思维:FAST 的优势全在差量小,差量占比大到一定程度,逐行应用变更反而比整表重算更慢。

小结

ob_mview_refresh.cpp 展示了 OceanBase 实现”可更新的查询快照”的完整方法论:用 SCN 给所有数据打版本,用 MLog 记录差量流水,用行锁 + 单事务保证并发安全与原子性,用依赖检测 + 清理水位校验守住增量正确性,最后用代价感知的自适应启发式在增量和全量之间动态取舍。整个模块只有约千行,却同时涉及内表元数据、内部 SQL 执行、RPC 转发 DDL 框架、统计收集等多层机制,是理解 4.3 存储层特性实现的绝佳样本。下一篇我们继续解析 4.3 新特性系列的其他文件。

发表回复

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