在 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 继承自 ObSimpleThreadPool,handle 是线程池消费入口。它处理的事务异步任务包括:
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 状态机。