PalfHandleImpl 是 OceanBase PALF(Paxos Asynchronous Log Flush)日志服务的核心句柄,一个实例对应一个日志流(palf_id)。它向上承接 SQL/事务层的日志提交请求,向下协调状态机、成员配置、选举、日志引擎与滑动窗口,完成日志的写入、复制、落盘与提交。可以把 PalfHandleImpl 理解为”单日志流控制器”:所有与该日志流相关的 Paxos 协议行为,最终都汇聚到这个类中。
核心数据结构
PalfHandleImpl 内部聚合了多个子模块,每个子模块负责日志流的一个切面:
| 成员 | 职责 |
|---|---|
| state_mgr_ | 角色与任期管理,决定当前能否接收/提交/复制日志 |
| config_mgr_ | 成员列表、配置版本、加减副本等成员变更 |
| mode_mgr_ | 访问模式(APPEND/RAW_WRITE)及模式版本 |
| log_engine_ | 磁盘 IO、网络 RPC、日志块池的底层引擎 |
| sw_(LogSlidingWindow) | 日志排序、滑动窗口提交、多数派 ACK 统计 |
| election_ | 基于仲裁的 Leader 选举 |
| log_cache_ | 读路径缓存,加速 Follower 追日志 |
| lock_ | 读写锁,保护句柄级状态 |
关键流程一:Leader 提交日志
业务层调用 submit_log 提交单条日志。该接口先进行参数与磁盘空间校验,然后在读锁保护下检查 state_mgr_ 是否允许 append,最后把实际排序与 LSN/SCN 分配交给滑动窗口。
int PalfHandleImpl::submit_log(
const PalfAppendOptions &opts,
const char *buf,
const int64_t buf_len,
const SCN &ref_scn,
LSN &lsn,
SCN &scn)
{
int ret = OB_SUCCESS;
if (IS_NOT_INIT) {
ret = OB_NOT_INIT;
} else if (NULL == buf || buf_len MAX_LOG_BODY_SIZE
|| !ref_scn.is_valid()) {
ret = OB_INVALID_ARGUMENT;
} else {
RLockGuard guard(lock_);
if (false == palf_env_impl_->check_disk_space_enough()) {
ret = OB_LOG_OUTOF_DISK_SPACE;
} else if (!state_mgr_.can_append(opts.proposal_id, opts.need_check_proposal_id)) {
ret = OB_NOT_MASTER;
} else if (OB_FAIL(sw_.submit_log(buf, buf_len, ref_scn, lsn, scn))) {
if (OB_EAGAIN != ret) {
PALF_LOG(WARN, "submit_log failed", KPC(this), KP(buf), K(buf_len));
}
} else {
PALF_LOG(TRACE, "submit_log success", K(ret), KPC(this), K(buf_len), K(lsn), K(scn));
}
}
return ret;
}
这里体现了 PALF 的分层思想:PalfHandleImpl 作为控制器只做”能不能写”的粗粒度判断,真正的并发序列化由 sw_ 完成。这样做既简化了控制器,又允许滑动窗口针对日志顺序做精细优化。
关键流程二:日志落盘
滑动窗口决定某条日志需要落盘后,会调用 inner_append_log,由 LogEngine 执行真正的 pwrite。该函数属于”内核层”写盘,不检查角色,只负责把 WriteBuf 写到指定 LSN。
int PalfHandleImpl::inner_append_log(const LSN &lsn,
const LogWriteBuf &write_buf,
const SCN &scn)
{
int ret = OB_SUCCESS;
const int64_t begin_ts = ObTimeUtility::current_time();
if (IS_NOT_INIT) {
ret = OB_NOT_INIT;
} else if (false == lsn.is_valid() || false == write_buf.is_valid()) {
ret = OB_INVALID_ARGUMENT;
} else if (OB_FAIL(log_engine_.append_log(lsn, write_buf, scn))) {
PALF_LOG(ERROR, "LogEngine pwrite failed", K(ret), KPC(this), K(lsn), K(scn));
} else {
const int64_t curr_size = write_buf.get_total_size();
const int64_t accum_size = ATOMIC_AAF(&accum_write_log_size_, curr_size);
const int64_t now = ObTimeUtility::current_time();
const int64_t time_cost = now - begin_ts;
append_cost_stat_.stat(time_cost);
if (time_cost >= 5 * 1000) {
PALF_LOG_RET(WARN, OB_ERR_TOO_MUCH_TIME, "write log cost too much time", ...);
}
}
return ret;
}
落盘完成后,IO 线程通过 inner_after_flush_log 回调通知滑动窗口。此时日志已持久化,可以推进滑动窗口、释放资源并触发后续提交。
int PalfHandleImpl::inner_after_flush_log(const FlushLogCbCtx &flush_log_cb_ctx)
{
int ret = OB_SUCCESS;
const int64_t begin_ts = ObTimeUtility::current_time();
RLockGuard guard(lock_);
if (IS_NOT_INIT) {
ret = OB_NOT_INIT;
} else if (OB_FAIL(sw_.after_flush_log(flush_log_cb_ctx))) {
PALF_LOG(WARN, "sw_.after_flush_log failed", K(ret), K(flush_log_cb_ctx));
} else {
const int64_t time_cost = ObTimeUtility::current_time() - begin_ts;
flush_cb_cost_stat_.stat(time_cost);
PALF_LOG(TRACE, "after_flush_log success", K(ret));
}
return ret;
}
关键流程三:Follower 接收日志
Leader 通过 RPC 把日志推给 Follower,Follower 侧统一进入 receive_log_。该函数先根据消息 proposal_id 更新本地任期,再在校验状态后把日志交给滑动窗口。为了处理日志分叉,这里采用了”两次调用滑动窗口”的策略:第一次允许生成截断信息,如果确实存在分叉,则加写锁清理或截断,之后再正式写入。
int PalfHandleImpl::receive_log_(const common::ObAddr &server,
const PushLogType push_log_type,
const int64_t &msg_proposal_id,
const LSN &prev_lsn,
const int64_t &prev_log_proposal_id,
const LSN &lsn,
const char *buf,
const int64_t buf_len)
{
int ret = OB_SUCCESS;
TruncateLogInfo truncate_log_info;
if (IS_NOT_INIT) {
ret = OB_NOT_INIT;
} else if (OB_FAIL(try_update_proposal_id_(server, msg_proposal_id))) {
PALF_LOG(WARN, "try_update_proposal_id_ failed", K(ret), KPC(this), K(server), K(msg_proposal_id));
} else {
RLockGuard guard(lock_);
if (false == palf_env_impl_->check_disk_space_enough()) {
ret = OB_LOG_OUTOF_DISK_SPACE;
} else if (!state_mgr_.can_receive_log(msg_proposal_id)) {
ret = OB_STATE_NOT_MATCH;
} else if (OB_FAIL(sw_.receive_log(server, push_log_type, prev_lsn, prev_log_proposal_id,
lsn, buf, buf_len, true, truncate_log_info))) {
// 第一次:允许生成 truncate 信息
} else {
PALF_LOG(TRACE, "receive_log success", K(ret), KPC(this), K(server), K(msg_proposal_id), K(lsn));
}
}
// 若存在分叉,先截断/清理,再第二次写入
if (OB_SUCC(ret) && (INVALID_TRUNCATE_TYPE != truncate_log_info.truncate_type_)) {
WLockGuard guard(lock_);
// ... 执行 TRUNCATE_CACHED_LOG_TASK 或 TRUNCATE_LOG ...
if (OB_FAIL(sw_.receive_log(server, push_log_type, prev_lsn, prev_log_proposal_id,
lsn, buf, buf_len, false, truncate_log_info))) {
PALF_LOG(WARN, "sw_ receive_log failed", K(ret), KPC(this), K(server), K(msg_proposal_id), K(lsn));
}
}
return ret;
}
对于落后较多的副本,Leader 会走批量路径 receive_batch_log:把多条日志打包成一个 buffer,Follower 用 MemoryStorage 和迭代器拆分为单条后再逐条写入。这种方式减少了 RPC 往返,是追日志的主要加速手段。
关键流程四:ACK 与提交推进
Follower 收到并处理日志后,向 Leader 发送 ack_log。Leader 在 ack_log 中校验 proposal_id,再交给滑动窗口统计多数派。当一条日志被多数派确认后,sw_ 会推进 committed_end_lsn,事务层据此认为日志已提交。
int PalfHandleImpl::ack_log(const common::ObAddr &server,
const int64_t &proposal_id,
const LSN &log_end_lsn)
{
int ret = OB_SUCCESS;
RLockGuard guard(lock_);
if (IS_NOT_INIT) {
ret = OB_NOT_INIT;
} else if (!server.is_valid() || INVALID_PROPOSAL_ID == proposal_id || !log_end_lsn.is_valid()) {
ret = OB_INVALID_ARGUMENT;
} else if (!state_mgr_.can_receive_log_ack(proposal_id)) {
// cannot handle log ack, skip
} else if (OB_FAIL(sw_.ack_log(server, log_end_lsn))) {
PALF_LOG(WARN, "ack_log failed", K(ret), KPC(this), K(server), K(proposal_id), K(log_end_lsn));
} else {
PALF_LOG(TRACE, "ack_log success", K(ret), KPC(this), K(server), K(proposal_id), K(log_end_lsn));
}
return ret;
}
小结
PalfHandleImpl 是 PALF 单日志流的”总指挥”。它通过读写锁保护句柄级状态,把具体协议行为委派给滑动窗口、状态机、配置管理、选举和日志引擎:
- submit_log 是业务写入口,负责准入检查并把日志交给滑动窗口排序;
- inner_append_log / inner_after_flush_log 构成 WAL 落盘与回调闭环,保证先持久化再可见;
- receive_log_ / receive_batch_log 处理 Follower 侧复制,用两次滑动窗口调用解决日志分叉;
- ack_log 汇总多数派确认,推动提交点前进。
理解 PalfHandleImpl 的关键,不在于记住每个函数,而在于看清”控制器 + 滑动窗口”的分层:控制器做状态检查与资源协调,滑动窗口做日志顺序与提交语义。下一篇将深入选举实现,看 PALF 如何在多副本间完成 Leader 选举与任期推进。