OceanBase源码解读:事务服务总控入口 ob_trans_service.cpp

在 OceanBase 的存储引擎中,事务模块负责保证 ACID、多版本并发控制(MVCC)以及分布式两阶段提交(2PC)。ob_trans_service.cpp 位于 storage/tx 层,是事务引擎在 observer 进程内的租户级门面与总控。它通过 MTL(Multi-Tenant Layer)机制为每个租户独立实例,对上承接 SQL 层的事务请求,对下统管事务上下文管理器、事务描述符管理器、GTS 源、RPC、定时器以及重复表(dup table)等子系统。本文聚焦这个“总控入口”,看它如何把散落的子系统装配成一套可运行、可停机、可扩展的事务服务。

核心数据结构

ObTransService 是整个事务引擎的租户级单例,关键成员构成如下:

class ObTransService : public ObSimpleThreadPool {
  bool is_inited_;
  bool is_running_;
  ObITransRpc *rpc_;                 // 事务 RPC 门面
  ObIDupTableRpc *dup_table_rpc_;    // 重复表 RPC
  ObILocationAdapter *location_adapter_; // 位置服务适配器
  ObIGtiSource *gti_source_;         // 全局时间戳源
  ObTsMgr *ts_mgr_;                  // 时间戳管理器
  ObTransTimer timer_;               // 事务时间轮定时器
  ObTxCtxMgr tx_ctx_mgr_;            // 事务上下文管理器(按 LS 组织)
  ObTxDescMgr tx_desc_mgr_;          // 事务描述符管理器(按会话事务组织)
  ObDupTableLoopWorker dup_table_loop_worker_;
  ObTabletLSCache tablet_to_ls_cache_;
  TenantConfigCache tenant_config_cache_;
  ...
};

其中 ObTxDesc 描述一次用户事务的生命周期,跨多条 SQL 语句保持;ObPartTransCtx 描述单个日志流(LS)上的参与者事务状态机;SCN 是 System Change Number,决定事务可见性与提交版本。三者共同构成“事务请求 → 描述符 → 参与者上下文”的三层模型。

生命周期:从装配到停机

mtl_init:注入基础设施

MTL 框架在租户启动时回调 mtl_init,把全局上下文中的 RPC 传输、位置服务、模式服务等基础设施注入当前租户实例:

int ObTransService::mtl_init(ObTransService *&it) {
  const ObAddr &self = GCTX.self_addr();
  share::ObLocationService *location_service = GCTX.location_service_;
  share::schema::ObMultiVersionSchemaService *schema_service = GCTX.schema_service_;
  obrpc::ObBatchRpc *batch_rpc = GCTX.batch_rpc_;
  obrpc::ObSrvRpcProxy *rpc_proxy = GCTX.srv_rpc_proxy_;
  share::ObAliveServerTracer *server_tracer = GCTX.server_tracer_;
  ObSrvNetworkFrame *net_frame = GCTX.net_frame_;
  rpc::frame::ObReqTransport *req_transport = net_frame->get_req_transport();
  if (OB_FAIL(it->rpc_def_.init(it, req_transport, self, batch_rpc))) {
    TRANS_LOG(ERROR, "rpc init error", KR(ret));
  } else if (OB_FAIL(it->location_adapter_def_.init(schema_service, location_service))) {
    TRANS_LOG(ERROR, "location adapter init error", KR(ret));
  } else if (OB_FAIL(it->gti_source_def_.init(self, req_transport))) {
    TRANS_LOG(ERROR, "gti source init error", KR(ret));
  } else if (OB_FAIL(it->init(self, &it->rpc_def_, ..., rpc_proxy, schema_service, server_tracer))) {
    TRANS_LOG(ERROR, "trans-service init error", KR(ret), KPC(it));
  }
  return ret;
}

这里只做“依赖装配”,真正的状态机初始化交给 init。一旦 mtl_init 失败,该租户的事务服务就无法启动,observer 启动流程会感知错误。

init:按租户装配状态机

init 是事务服务的“装配车间”。它先根据租户内存容量计算事务消息任务队列大小,再顺序初始化定时器、重复表扫描定时器、线程池、事务描述符管理器、事务上下文管理器、重复表工作线程、回滚段消息管理器、tablet→LS 缓存、只读检查器和租户配置缓存:

int ObTransService::init(...) {
  const int64_t tenant_memory_limit = lib::get_tenant_memory_limit(tenant_id);
  int64_t msg_task_cnt = MSG_TASK_CNT_PER_GB * (tenant_memory_limit / (1024 * 1024 * 1024));
  msg_task_cnt = std::max(msg_task_cnt, MSG_TASK_CNT_PER_GB);
  msg_task_cnt = std::min(msg_task_cnt, MAX_MSG_TASK_CNT);
  if (OB_FAIL(timer_.init("TransTimeWheel"))) { ... }
  else if (OB_FAIL(ObSimpleThreadPool::init(2, msg_task_cnt, "TransService", tenant_id))) { ... }
  else if (OB_FAIL(tx_desc_mgr_.init(...))) { ... }
  else if (OB_FAIL(tx_ctx_mgr_.init(tenant_id, ts_mgr, this))) { ... }
  else if (OB_FAIL(dup_table_loop_worker_.init())) { ... }
  ...
  else { is_inited_ = true; }
  return ret;
}

所有子系统初始化成功后才会把 is_inited_ 置为 true,否则调用方不会进入 start。这种“瀑布式错误处理”避免了半初始化状态继续运行。

start / stop / wait / destroy:标准 MTL 四步

start 按依赖顺序启动已初始化的子系统:定时器 → 重复表扫描任务 → 事务 RPC → GTS 源 → 事务上下文管理器 → 事务描述符管理器 → 租户配置刷新。任意子系统启动失败都会回传错误,保证不会进入“部分运行”状态。

stop 则与 start 顺序相反:先停事务上下文管理器与描述符管理器(停止接收新事务),再停定时器与 GTS 源,最后停 RPC。特别注意 tx_ctx_mgr_ 必须在 timer_ 之前停止,否则定时器回调可能访问已释放的上下文,引发悬空指针。

获取快照:get_gts_

强一致读需要获取一个“足够新”的快照版本号。get_gts_ 向全局时间戳服务申请当前最新时间戳,并与本地 tx_version_mgr_.max_commit_ts_ 取较大值:

int ObTransService::get_gts_(SCN &snapshot_version, MonotonicTs &receive_gts_ts,
                             const int64_t trans_expired_time,
                             const int64_t stmt_expire_time,
                             const uint64_t tenant_id) {
  const MonotonicTs request_ts = MonotonicTs::current_time();
  const int64_t WAIT_GTS_US = 500;
  SCN gts;
  MonotonicTs tmp_receive_gts_ts;
  do {
    if (ObClockGenerator::getClock() >= trans_expired_time) {
      ret = OB_TRANS_TIMEOUT;
    } else if (ObClockGenerator::getClock() >= stmt_expire_time) {
      ret = OB_TRANS_STMT_TIMEOUT;
    } else if (OB_FAIL(ts_mgr_->get_gts(tenant_id, stc_ahead, NULL, gts, tmp_receive_gts_ts))) {
      if (OB_EAGAIN != ret) { TRANS_LOG(WARN, "get gts failed", KR(ret)); }
      else { ob_usleep(WAIT_GTS_US); }
    } else if (!gts.is_valid()) { ret = OB_ERR_UNEXPECTED; }
    else {
      const SCN max_commit_ts = tx_version_mgr_.get_max_commit_ts(true);
      snapshot_version = SCN::max(max_commit_ts, gts);
      receive_gts_ts = tmp_receive_gts_ts;
    }
  } while (OB_EAGAIN == ret);
  return ret;
}

这里体现了 OceanBase 对“单调读”与“外部一致性”的追求:即使 GTS 服务返回的时间戳略落后于本地已提交事务的最大版本,也要取二者较大值,确保读到最新已提交数据。

一阶段提交:end_1pc_trans

对于单分区事务,OceanBase 走一阶段提交(1PC)路径。end_1pc_trans 是 SQL 层 COMMIT / ROLLBACK 在存储引擎侧的落点之一:

int ObTransService::end_1pc_trans(ObTxDesc &trans_desc,
                                  ObITxCallback *endTransCb,
                                  const bool is_rollback,
                                  const int64_t expire_ts) {
  if (is_rollback) {
    interrupt(trans_desc, ObTxAbortCause::EXPLICIT_ROLLBACK);
    if (OB_FAIL(rollback_tx(trans_desc))) { ... }
  } else if (OB_FAIL(submit_commit_tx(trans_desc, expire_ts, *endTransCb))) { ... }
  TRANS_LOG(INFO, "end 1pc trans", KR(ret), K(trans_desc.tid()));
  return ret;
}

回滚路径先标记事务为中断,再执行 rollback_tx;提交路径调用 submit_commit_tx 进入本地日志提交流程。这个函数虽然短,却是用户事务“落盘”的关键分叉口。

异步任务总控:handle

ObTransService 继承自 ObSimpleThreadPoolhandle 是线程池消费入口。它处理的事务异步任务包括:

  • END_TRANS_CB_TASK:事务提交回调,若未触发则重新排队,实现“软定时”效果。
  • ADVANCE_LS_CKPT_TASK:推进日志流 checkpoint 时间戳。
  • STANDBY_CLEANUP_TASK:备库事务清理。
  • DUP_TABLE_TX_REDO_SYNC_RETRY_TASK:重复表 redo 同步重试。
void ObTransService::handle(void *task) {
  ATOMIC_FAA(&output_queue_count_, 1);
  ObTransTask *trans_task = static_cast(task);
  if (!trans_task->ready_to_handle()) {
    push(trans_task); // 尚未 ready,重新排队
  } else if (ObTransRetryTaskType::END_TRANS_CB_TASK == trans_task->get_task_type()) {
    ObTxCommitCallbackTask *cb_task = static_cast(task);
    bool has_cb = false;
    if (cb_task->get_need_wait_us() > 0) { ob_usleep(cb_task->get_need_wait_us()); }
    if (OB_FAIL(cb_task->callback(has_cb))) { ... }
    if (has_cb) { ObTxCommitCallbackTaskFactory::release(cb_task); }
    else { push(cb_task); }
  } else if (ObTransRetryTaskType::ADVANCE_LS_CKPT_TASK == ...) {
    ...
  } else if (REACH_TIME_INTERVAL(10 * 1000 * 1000)) {
    // 每 10s 打印队列与 TxDesc 统计,便于排查堆积与句柄泄漏
  }
}

Multi-Data Source 注册

MDS(Multi-Data Source)用于承载 DDL、隐藏行、location 变更等“非用户表数据”的事务化写入。register_mds_into_tx 把一段二进制数据注册到指定 LS 的事务上下文中,随事务一起提交到日志:

int ObTransService::register_mds_into_tx(ObTxDesc &tx_desc, const ObLSID &ls_id,
                                         const ObTxDataSourceType &type,
                                         const char *buf, const int64_t buf_len, ...) {
  create_implicit_savepoint(tx_desc, tx_param, savepoint); // 失败可回滚
  if (!seq_no.is_valid()) { seq_no = tx_desc.inc_and_get_tx_seq(0); }
  arg.init(tx_desc.tenant_id_, tx_desc, ls_id, type, str, seq_no, request_id, register_flag);
  do {
    if (OB_NOT_MASTER == ret) { ob_usleep(RETRY_INTERVAL); retry_cnt++; }
    if (ls_leader_addr == self_) {
      register_mds_into_ctx_(tx_desc, ls_id, type, buf, buf_len, seq_no, register_flag);
    } else if (self_ == tx_desc.addr_) {
      rpc_proxy_->to(ls_leader_addr).by(tx_desc.tenant_id_)
                 .timeout(remain_timeout_us).register_tx_data(arg, result);
    }
  } while (OB_NOT_MASTER == ret && self_ == tx_desc.addr_);
  return ret;
}

该函数先创建隐式 savepoint,若注册失败则回滚,保证原子性;若本地不是 leader,则通过 RPC 转发到 leader,遇到 OB_NOT_MASTER 时刷新位置缓存并重试。最终统一收集 tx_result 并合并到事务描述符中。

小结

ob_trans_service.cpp 是理解 OceanBase 事务引擎的最佳入口之一。它不承担具体的 MVCC 锁、2PC 状态机细节,而是负责把 RPC、GTS、定时器、事务上下文、事务描述符等子系统装配成一套租户级服务,并提供统一的生命周期管理、快照获取、一阶段提交和 MDS 注册接口。最值得带走的三点设计:一是“瀑布式初始化 + 全部成功才置 is_inited_”,避免半初始化状态;二是 stop 与 start 的严格逆序,防止定时器回调访问已释放上下文;三是 MDS 的隐式 savepoint + 失败重试机制,保证非用户数据也能事务化地安全写入。要继续深挖,可以顺着 tx_ctx_mgr_ 进入 ob_tx_ctx_mgr.cpp,或顺着 submit_commit_tx 进入 ObPartTransCtx 的 2PC 状态机。

发表回复

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