OceanBase源码解读:Palf日志句柄 palf_handle_impl.cpp

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 选举与任期推进。

发表回复

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