CMU15445数据库系统通关指南

前言

关于Vibe Coding

其实动手写这个作业,一方面原因是当初大学就没认真学数据库,虽然这些年在工作中积累了一些经验,但还是对底层细节缺乏深刻的理解,因此对这门CMU神课种草许久,可惜难以找出整块的时间来完成它。另一部分原因,也是经过了一年多的职场gap之后,想借助这个机会,感受一下目前AI在编程领域能够起到怎样的辅助效果。体验得出的结论是,AI在context较小、思考难度较深的实际工业问题上,表现出的性能令人震惊,整体上给人感受像是记忆力有限的编程天才,这也与其在ICPC决赛中AK所有题目的能力相符。但是一旦context变大,消耗的token急剧增多,其给出结论的精准性也明显下降。

在内存池部分的代码中,我为了优化并发性能做了细粒度锁分片,这部分代码极其难写对,涉及到复杂的嵌套锁、二次检查以及try-lock逻辑,但在AI的辅助下,基本能够抓到所有的并发bug。而在B+树的相关编程中,涉及一些非常繁琐的边界处理、offset计算,这部分代码在AI的辅助下可以很顺利地完成并验证正确性,节省了大量的调试时间。

比较神奇的一点是,如果你怀疑哪部分代码有某种“坏味道”,隐隐约约有些模糊的想法,那么直接拿去问AI,它几乎都能理解问题,并给你构造出一个合适的、触发racing condition的例子,从而帮助你快速确认/否认想法。由此可见,AI是真实理解了多线程编程技术,并能够创造性地将其应用在真实场景中。但尽管如此,当前的AI也并未摆脱“讨好型人格”的问题——如果你排查bug时完全怀疑错了方向,那么AI并不能给你一个肯定的“这部分代码没毛病”的结论,而是会尽其所能地胡编乱造。

编码建议

  • 良好的编程风格和接口设计有助于降低思维负担,并减少出错的可能
  • 不要搞骚操作!如果你的代码有一些“我寻思”的迹象,一般就说明哪里踩坑了
  • 对于并发竞态问题,如果感觉代码哪里有不对劲的“坏味道”,可以直接拿去问AI,不用自己一直拍脑袋
  • 多线程场景下要想清楚对象生命周期,免得析构函数出什么奇怪的漏子
  • 在macbook air上,benchmark表现比较奇怪,只能作为参考,正经的性能数据还是要看gradescope上面的结果

通关感受

  • 项目整体代码框架非常精巧,尤其是PageGuard以及TupleMeta的设计,基本上不太容易踩到恶性的死锁或者内存异常bug
  • Benchmark的验证似乎比较弱,leaderboard上面有一些明显超出正常性能范围的hack,这样就让排名没意义了;而且最后的leaderboard task甚至有bug导致优化无法生效,比如作业4实现实时GC的时候(这个问题已经被修复
  • 作业与作业之间的unit test不完全兼容,需要想一些tricky的方式绕过,对于后续回头做优化刷benchmark不太友好

关于调试

首先要把基于VSCode的调试环境配置好,这可以为后期节省大量时间精力。为了配合CMake Tools,在macOS上,launch.json大致内容是这样:

{
  "version": "0.2.0",
  "configurations": [
    {
      "name": "(lldb) Launch",
      "type": "cppdbg",
      "request": "launch",
      "preLaunchTask": "CMake: build",
      "program": "${command:cmake.launchTargetPath}",
      "cwd": "${workspaceFolder}/build",
      "args": [],
      "stopAtEntry": false,
      "environment": [],
      "externalConsole": true,
      "MIMode": "lldb",
    },
  ]
}

在整个作业的编码过程中,我基本上很少有遇到难以调试的死锁问题。一个比较高效的调试方法是直接把程序运行到死锁位置处,然后暂停调试器,检查各个线程的调用栈,看看有没有锁重入的问题,一般都能很直观地找到bug。

另一方面,对于一些运行时异常,也可以采取直接无脑调试器运行,走到异常抛出位置之后会自动暂停,此时再检查调用栈,回溯各个上下文的变量,尝试还原错误现场来找出问题所在。

性能分析利器:火焰图

火焰图是分析程序性能的有效工具,尤其是这种跑分代码都已经帮你写好了的场景。macOS上使用如下脚本可以用dtrace来收集作业1的benchmark性能数据并可视化:

sudo echo
./bin/bustub-bpm-bench --duration 1000 --latency 1 &
PID=$!
sudo dtrace -x ustackframes=100 \
  -n 'profile-2000 /pid == $target/ { @[ustack()] = count(); }' \
  -p $PID > out.stacks
stackcollapse.pl out.stacks > out.folded
flamegraph.pl out.folded > out.svg
open out.svg

右键在新页面中打开以下图片,可以用Chrome进行交互式查看。从火焰图中我们可以看到,对于这个测试case的主要性能瓶颈集中在磁盘IO上,因此可以阅读官方代码的IO模拟逻辑,并进行针对性优化。

Flame

Buffer Pool Manager

先上leaderboard跑分:

排名最靠前的几个明显是做了某种hack,后面大家都差不太多,这个分数算是还可以。

LRU-K算法

首先官方指南上来就给你挖了个坑:LRU-K算法的复杂度其实是O(logN)的,之后在跑B+树benchmark的时候非常吃亏,所以后面需要你改写成复杂度O(1)的实现。一个简单的做法是直接按是否有K次访问来分成两个LRU队列即可,然后你可以加一个诸如LRUKReplacer::enable_fast_mode_的变量,来切换想用哪个算法,也便于保证逻辑的一致性。

Page IO优化

显然我们不应该在持有全局锁的时候做IO,而是要把IO请求丢给线程池异步处理,不然性能会很差,那么就需要想清楚一些细节问题:

  • 当我们写回一个dirty frame的时候,需要一直占用着这个frame吗?
    • 并不一定,我们可以将其数据存进buffer,此时这个frame就已经可用了,然后用线程池将buffer异步写入磁盘对应的page_id
  • 那么问题来了:当我们试图加载一个page_id的时候,发现它正在被异步写入磁盘,这时候应该怎么做?
    • 应该用条件变量挂起当前调用线程,直到这个page_id上的IO操作完成,这时候我们才能合法读取这个page

由此,我们可以抽象出如下接口,用来管理page_id上的IO状态:

// class BufferPoolManager....
std::unordered_set<page_id_t> io_pages_;
std::condition_variable io_cv_;

void BufferPoolManager::WaitForPageReady(std::unique_lock<std::mutex> &lck, page_id_t page_id) {
  io_cv_.wait(lck, [page_id, this]() { return io_pages_.count(page_id) == 0; });
}

void BufferPoolManager::MarkPageBusy(page_id_t page_id) { io_pages_.insert(page_id); }

void BufferPoolManager::MarkPageReady(page_id_t page_id) {
  BUSTUB_ENSURE(io_pages_.count(page_id), "page_id not in IO status");
  io_pages_.erase(page_id);
  io_cv_.notify_all();
}

新的问题来了:如果我们试图ReadPage的时候,系统里既没有空闲frame,又没有可以淘汰的frame,此时应该如何处理?

答案是同样用条件变量陷入等待,直到有新的frame被归还给系统——但是注意,如果你直接这样实现,那么EvictableTest测试会无法通过,因为它就是期望测试Evict()失败不得不返回nullopt的场景。有个两全其美的方案,就是稍微修改一下接口,例如:

// class BufferPoolManager....
CheckedReadPage(page_id_t page_id, AccessType access_type, bool wait_for_page=false)

auto BasicBufferPoolManager::ReadPage(page_id_t page_id, AccessType access_type) -> ReadPageGuard {
  auto guard_opt = CheckedReadPage(page_id, access_type, true);
  if (!guard_opt.has_value()) {
    fmt::println(stderr, "\n`CheckedReadPage` failed to bring in page {}\n", page_id);
    std::abort();
  }
  return std::move(guard_opt).value();
}

// when allocating a new frame...
if (wait_for_page) {
  frames_cv_->wait_for(lock, std::chrono::seconds(3),
                       [this]() { return !free_frames_.empty() || replacer_->Size() > 0; });
}

这样的话我们就区分了CheckedReadPageReadPage的语义:调用前者时可以选择是否等待可用frame,并且允许返回空值,这样就满足了unit test需求;而后者则是期望一定要拿到一个frame,并且允许调用线程挂起等待,直到系统中某个frame可用。

最后一个优化思路是关于磁盘IO。在DiskManagerUnlimitedMemory里可以看到对于IO latency的模拟,简而言之就是连续的磁盘IO,或者落在同一个block内的IO会有更低的latency。我们可以手动在这部分代码里加入一些统计,计算有多少比例的IO操作被“惩罚”了高延迟。一种优化思路是对于线程池里的每一个线程,它监听连续的一系列page_id(例如16个为一组),这样每次尽可能攒出一个batch,让IO请求尽可能连续,再调用实际的磁盘操作,最终可以把IO连续命中率从纯随机访问的50%优化到75%左右。

锁分片

这部分内容不在官方实验指南中,是我自己在性能优化方面所做的一些尝试。非常不建议做这个优化,调试起来简直是地狱难度,而且对于leaderboard workload也并没有太高的收益。这里详细讲讲锁分片难在哪里。

1. 为什么需要锁分片

锁分片(Lock Striping)本质就是把原本保护所有资源的一把大锁bpm_latch_,拆成了更细粒度的锁,从而在高负载环境下降低锁竞争。如果内存池所有操作都用一个全局锁,调用ReadPage/WritePage时线程会频繁阻塞在上面,系统吞吐量低。将所有的page按照page_id % NUM_PARTITIONS方式分片,每个分片独立加锁,不同分片上的操作可以并发进行,大大减少锁竞争,充分利用多核优势。

// class BufferPoolManager...
constexpr static size_t NUM_PARTITIONS = 32;
std::array<std::unordered_map<page_id_t, frame_id_t>, NUM_PARTITIONS> page_table_;
std::array<std::mutex, NUM_PARTITIONS> latches_;
// data structure for pages that pending IO
std::array<std::unordered_set<page_id_t>, NUM_PARTITIONS> io_pages_;
std::array<std::condition_variable, NUM_PARTITIONS> io_cv_;

2. 为什么需要局部锁与全局锁的协作

有些数据结构(如page_table_io_pages_)只影响单个分片,用局部锁保护即可。但像 free_frames_replacer_ 这类资源是全局的——因为frames是整个系统中共享的,因此必须用全局锁保护。当操作既涉及局部资源又涉及全局资源时,需要先加局部锁,再加全局锁,保证一致性和避免死锁。这样既保证了高并发,又能安全地管理全局状态,但却给代码实现带来了很大的挑战。例如,当我们检查page_table_,发现page_id已经被影射到某个frame上,此时需要把这个frame给pin住——这个步骤其实非常复杂,我们在只持有局部锁的情况下,无法避免其他线程并发淘汰(Evict())并重新分配这个frame,因为frame的生命周期是由全局逻辑管理的。我们必须先升级到全局锁,然后二次检查frame的page_id_pin_count_,确保frame还映射到当前page且pin_count_仍为0,此时再去pin住这个frame才是安全的。如果没有全局锁保护,可能会出现如下bug:

  • 线程A和B都看到pin_count_==0,A准备加pin,B持有全局锁把frame给Evict掉并分配给新page;线程A随后pin_count_++,导致新page的frame被错误地加pin,数据错乱
  • 也可能在线程A准备加pin前,线程B已经把这个frame回收并放入空闲链表,于是之后会出现两个page共用一个frame的场景

3. 为什么需要CAS操作

对于增加pin count这个步骤,除了上面的锁升级,还存在一个极其隐蔽的竞态条件。如果我们发现pin_count_>0,说明frame已经被pin住,此时是不是可以直接pin_count_++并返回FrameHeader了呢?并不是!虽然原子变量的自增是绝对安全,不会导致计数错误,但是在条件检查和自增之间并不是原子的,这里可能产生竞态:在此间隙期间,别的线程可能已经把pin_count_又变成了0,这时候已经不再满足先前的判断条件,要做的操作也不一样了。因此,这里我们必须要引入CAS操作,保证条件检查和增加计数的原子性,否则就直接重试:

while (true) {
  auto frame_hdr = GetMappedFrameLocked(page_id);
  if (frame_hdr) {
    auto old_pin_count = frame_hdr->pin_count_.load();
    if (old_pin_count > 0) {
      if (frame_hdr->pin_count_.compare_exchange_strong(old_pin_count, old_pin_count + 1)) {
        // ...
        return frame_hdr;
      }
    } else {
      std::scoped_lock global(*bpm_latch_);
      if (frame_hdr->page_id_ == page_id && frame_hdr->pin_count_ == 0) {
        replacer_->SetEvictable(frame_hdr->frame_id_, false);
        frame_hdr->pin_count_++;
        // ...
        return frame_hdr;
      }
    }
  } else {
    break;
  }
}

一个值得思考的问题:为什么这里做CAS时,不会因为ABA问题而导致错误?

4. 为什么需要 try-lock

正常情况下,获取锁必须遵循统一的先局部后全局的顺序,否则会导致死锁。但是在特殊情况下,比如淘汰页面时,我们需要先持有全局锁,用Evict()拿到一个FrameHeader,再去获取它对应page_id的局部锁来清除现有的映射关系。此时,我们就必须使用try-lock而不是阻塞锁,因为如果拿不到锁的话,我们不能挂起,而是直接归还这个frame,再尝试重新获取一个新的frame。类似地,如果我们拿到的frame所对应的page_id正处于IO状态,也不得不归还它。

5. 为什么需要做各种二次检查

二次检查是为了防止在检查条件和拿到锁的间隙产生竞态:因为在你检查条件之后、拿到锁前,其他线程可能已经更新了状态,只有二次检查才能确保你的操作基于最新的、正确的状态,避免并发下的逻辑漏洞。例如当我们为page_id分配一个新的frame时并不需要分区锁的保护,等到准备更新page_table_中的映射关系时才会去拿锁,此时我们可能会发现已经有其他线程把这个page给映射好了,于是我们就必须归还这个新拿到的frame,然后从头开始去读取映射,否则就会覆盖已有的page<=>frame映射关系,导致数据错误。

类似地,在析构PageGuard的时候,我们会先检查pin_count_是否为0,但如果符合条件并升级到全局锁之后,还是要再做一遍检查,以确保没有其他线程又pin住了这个页面:

if (--frame_->pin_count_ == 0) {
  std::scoped_lock global(*bpm_latch_);
  if (frame_->pin_count_ == 0) {
    replacer_->SetEvictable(frame_->frame_id_, true);
    // ...
  }
}

DeletePage时要先拿分区锁,如果我们发现page_id已经被映射到frame上,并且pin_count_==0,那么此时对frame的任何后续操作都必须要在全局锁的保护下进行,并且升级到全局锁之后,也要做二次检查:

auto BufferPoolManager::DeletePage(page_id_t page_id) -> bool {
  auto frame_hdr = ...;
  if (frame_hdr->pin_count_ > 0) return false;
  std::scoped_lock global(*bpm_latch_);
  BUSTUB_ENSURE(frame_hdr->page_id_ == page_id, "Unexpected: frame mapped to different page_id");
  if (frame_hdr->pin_count_ > 0) return false;
  if (frame_hdr->is_dirty_) {
    // async flush data to disk...
  }
  // ...
}

因为在检查的这一刻,该frame完全没有被pin住,会有两种可能的情况出现:

  • 在拿到全局锁之前的间隙,这个frame就已经被Evict()了,但并不会产生任何修改例如重新分配给其他page,因为当前page_id所在的分片还是被锁住的,所以这种情况下是安全的
  • 也可能在同一瞬间,有并发线程调用WritePage重新pin住了这个frame,此时当我们拿到全局锁并二次检查发现pin_count_>0就应该立刻返回false,否则就会产生竞态条件:frame里的数据一边在向磁盘写回,一边被并发更新,导致数据错误

Database Index

还是先上leaderboard跑分:

B+树这部分的重点在于接口设计,好的设计可以大幅度降低编码的心智负担,并且有助于管理循环不变量,让代码的正确性更为显然。可以先把接口写好,整个算法逻辑捋顺了之后再填充细节。这部分编码比较建议直接vibe coding,因为都是一些具有强对称性的代码,AI补全的准确率相当高。

接口设计

分享一下我的B+树辅助函数接口:

class BPlusTree {
private:
  auto AcquireSiblingPage(Context &ctx, int choose) -> std::optional<WritePageGuard>;
  auto DeleteNodeFromParent(Context &ctx, WritePageGuard sibling_guard_to_del, int choose) -> void;
  auto FindLeafPageOptimistic(const KeyType &key, bool) -> WritePageGuard;
  auto FindLeafPage(const std::optional<KeyType> &key) -> ReadPageGuard;
  auto FindLeafPage(Context &ctx, const KeyType &key, std::function<bool(const InternalPage *)>, bool)
      -> WritePageGuard;
  auto BorrowFromRight(InternalPage *f, InternalPage *i, InternalPage *right, int offset) -> void;
  auto BorrowFromLeft(InternalPage *f, InternalPage *i, InternalPage *left, int offset) -> void;
  auto MergeIntoLeft(InternalPage *f, InternalPage *i, InternalPage *left, int offset) -> void;
  auto MergeFromRight(InternalPage *f, InternalPage *i, InternalPage *right, int offset) -> void;
  // ...
};

class BPlusTreeInternalPage : public BPlusTreePage {
public:
  std::pair<page_id_t, int> FindChildPage(const KeyType &key, const KeyComparator &less) const;
  void Insert(const KeyType &key, page_id_t page_id, const KeyComparator &comp);
  void InsertAt(int offset, const KeyType &key, page_id_t page_id);
  KeyType SplitInto(BPlusTreeInternalPage *other);
  page_id_t EraseAt(int offset);
  // ...
};
class BPlusTreeLeafPage : public BPlusTreePage {
public:
  auto KeyAt(int index) const -> const KeyType &;
  auto ValueAt(int index) const -> const ValueType &;
  auto Find(const KeyType &key, const KeyComparator &comp, int *offset_result) const -> std::optional<ValueType>;
  void Insert(const KeyType &key, const ValueType &value, const KeyComparator &comp);
  void InsertAt(int offset, const KeyType &key, const ValueType &value);
  void Erase(const KeyType &key, const KeyComparator &comp);
  void EraseAt(int offset);
  void SplitInto(BPlusTreeLeafPage *other);
  void AppendInto(BPlusTreeLeafPage *other);
  // ...
};

这里面BorrowFromRightBorrowFromLeftMergeIntoLeftMergeFromRight实现起来需要想清楚细节,比较建议构造一个具体的例子写在注释里,然后对照着进行编码,很有助于理顺逻辑。例如:

INDEX_TEMPLATE_ARGUMENTS
auto BPLUSTREE_TYPE::MergeIntoLeft(InternalPage *f, InternalPage *i, InternalPage *left, int offset) -> void {
  // Consider the following B+ tree structure:
  // Father:
  //
  //     FK1 [FK2] FK3 FK4
  // FP0 FP1 (FP2) FP3 FP4
  //      |   |
  //      L   I
  //
  // Left:
  //
  //     LK1 LK2 LK3 [FK2] IK1 IK2 IK3
  // LP0 LP1 LP2 LP3 [IP0] IP1 IP2 IP3
  //
  // Intermediate:
  //
  //       IK1 IK2 IK3
  // [IP0] IP1 IP2 IP3
  // This is to make sure: FK2 <= Key < IK1 should be routed to IP0
  left->InsertAt(left->GetSize(), f->KeyAt(offset), i->ValueAt(0));
  for (int j = 1; j < i->GetSize(); j++) {
    left->InsertAt(left->GetSize(), i->KeyAt(j), i->ValueAt(j));
  }
};

Leaderboard优化

性能优化方面,按照官方指南实现乐观加锁(FindLeafPageOptimistic)其实比较简单,但是tombstone deletion是一个比较tricky的策略,经历了上面做锁分片的折磨之后,我决定不踩这个大坑。其实benchmark的workload比较单一,只是一个稍微变体的线性插入,所以只需要在叶子节点分裂上面做一个简单的优化,让新节点尽可能留空即可,这样就能让顺序插入尽量不触发新的分裂:

INDEX_TEMPLATE_ARGUMENTS
void B_PLUS_TREE_LEAF_PAGE_TYPE::SplitInto(BPlusTreeLeafPage *other) {
  other->Init();
  auto remain = GetSize() / 2 + GetSize() % 2;
  // note: this is an allowed dedicated optimization for B+ tree sequential insert performance
  if (this->GetMaxSize() > 100) {
    remain = GetSize() * 0.9;
  }
  for (int i = remain; i < GetSize(); i++) {
    other->key_array_[i - remain] = key_array_[i];
    other->rid_array_[i - remain] = rid_array_[i];
  }
  other->SetSize(GetSize() - remain);
  this->SetSize(remain);
}

最后用火焰图再找找热点,发现LRU-K里的加锁也是个小瓶颈,把里面的std::mutex换成手搓自旋锁之后,还能再换来10%性能提升。

class Spinlock {
 private:
  std::atomic_flag flag = ATOMIC_FLAG_INIT;
 public:
  void lock() noexcept {
    while (flag.test_and_set(std::memory_order_acquire)) {
      // do not yield here for better performance
      // std::this_thread::yield();
    }
  }
  void unlock() noexcept { flag.clear(std::memory_order_release); }
  bool try_lock() noexcept { return !flag.test_and_set(std::memory_order_acquire); }
  Spinlock() = default;
  Spinlock(const Spinlock &) = delete;
  Spinlock &operator=(const Spinlock &) = delete;
};

Query Execution

老规矩先上跑分:

在这个作业里,实现各种query executors本身比较乏善可陈,基本上就是要理解对数据库“行”的抽象:SchemaColumnTupleValue,以及整个AbstractExecutor框架是怎么回事,然后把课上讲过的算法动手敲一遍,最复杂的无非就是External Merge Sort,但也没有什么坑。不过leaderboard task里面两个query的优化pass却是最大亮点——如果只是针对性地处理跑分用的特例倒是不难,但如果想要写得严谨通用还是很费功夫的。

Init()Next()语义

Executor的构造函数和Init()是比较容易被混淆的地方,需要意识到前者是用于初始化对象,只调用一次;但是Init()方法有可能被其father调用多次!例如对于NestedLoopJoinExecutor,每处理左表的一个元素,都要遍历一遍右表,于是就必须重新调用一遍right_executor_->Init()来重置内部状态,这样才能重新开始读取。

Optimizing SeqScan to IndexScan

首先要理解Optimizer的整体运行流程:不在原本的表达式树上做任何修改,而是如果命中优化条件,就直接创建新的AbstractPlanNodeRef并返回,而旧的子树会被RAII自动清理掉。

对于这个优化pass,逻辑还是比较简单的,我们只需要递归遍历表达式树,确保它里面的比较都是“等值查询”,即Column=ConstantValue的形式,并提取出所有比较表达式中的列和常量用于构造IndexScanPlanNode。注意由于我们的B+树只能单点查询,因此无法支持复合索引,在这个pass中只能做简化处理,即假设只会在一个列上做查询。

Optimizing NestedLoopJoin to HashJoin

与上一个pass非常相似,同样是遍历表达式树,检查它是全都由AND连接,并且要保证等号一侧的ColumnValue必须都来自同一个table,例如下面这条SQL就无法被优化成HashJoin

bustub> EXPLAIN (o) SELECT * FROM test_1 t1, test_2 t2 WHERE t1.colB + t2.colA = t2.colC;
=== OPTIMIZER ===
NestedLoopJoin { type=Inner, predicate=((#0.1+#1.0)=#1.2) }
  SeqScan { table=test_1 }
  SeqScan { table=test_2 }

External Merge Sort

这又是一个考察设计模式的题目,重点在于把所有临时pages都放进MergeSortRun里面来管理,写入时调用MergeSortRun::Add(const Tuple& tuple),然后用Iterator来处理元素的遍历读取(实现operator++()时要想清楚),这样就可以避免在merge sort算法里管理这些琐碎的边界计算细节以及pages的生命周期。虽然官方指南只要求你实现简单的2路归并,但其实写个通用的K路归并也不复杂,建议一步到位。

另外,建议想清楚在归并排序时,WritePageGuardReadPageGuard应该分别放在哪里进行管理?由于归并排序的上述读写特性,这里其实可以做一个小的优化,尽量重用Read/WritePageGuard以降低调用BufferPoolManager的不必要开销(因为涉及到拿锁)。具体来说:

  • IteratorReadPageGuard每次只锁定当前访问的页面,遍历到下一个页面时再lazy切换guard。
  • 每次MergeSortRun::Add只在新建页面或切换页面时才lazy调用bpm_->WritePage,避免每插入一个Tuple都去申请/释放页面锁。

Query优化1:谓词提取和下推

考虑如下SQL:

SELECT * FROM t4, t5, t6
  WHERE (t4.x = t5.x) AND (t5.y = t6.y) AND (t4.y >= 1000000)
    AND (t4.y < 1500000) AND (t6.x >= 100000) AND (t6.x < 150000);

这个优化思路看起来非常直观,就是把WHERE后面的等值连接下推到HashJoin子节点中,以及把范围过滤条件放进对应的SeqScan,但如果想写出一个通用的、能够处理任意个表连接的递归实现会麻烦得多,以至于回头再看自己写的代码都感觉难以理解,有太多的边界细节需要考虑,没有什么好的设计模式可以简化它。

事实上,bustub里的这种“手工递归遍历+类型分支”的写法,更多是关注教学和原理演示,实际工程中会引入更高层次的抽象。工业界数据库优化器通过“规则引擎+模式匹配+统一抽象+辅助工具”,允许开发者用声明式方式描述“模式-替换”规则,从而让优化pass的开发变得声明化、模块化、易组合、易维护,远比手写递归和类型判断要高效和健壮。

对于这个pass,具体思路如下:

  • 谓词分类与分离:遍历传入的谓词列表,判断哪些谓词可以作为Hash Join的等值连接条件(如t1.a = t2.b),然后分别存放。
  • 先下推谓词到扫描节点:对NLJ的每个子计划(child plan),如果是SeqScanMockScan,则将能下推的谓词给下推到扫描节点,提升过滤效率。这个过程中注意要对谓词表达式中的column index进行递归修正——直接设置为0即可。
  • 然后递归处理嵌套的NLJ:对于child plan list中仍为NestedLoopJoin的节点,先递归修正谓词表达式中的列索引,使其适配子NLJ的左右表结构,再递归调用自己进行处理。
  • 优化终止条件:如果还有剩余不能下推的谓词,说明无法完全优化,直接返回nullptr,表示本次重写失败。
  • 处理完子节点之后,创建新的NLJ节点并返回:
    • 如果当前有可用的Hash Join谓词,则提取出左右表的join key,重写为HashJoinPlanNode,以提升连接效率。
    • 如果此时已经没有可用的Hash Join谓词,则新建一个NLJ节点,谓词设为恒true(即无条件连接)。

这样,我们的优化器就可以完美处理更加复杂的情况:

bustub> EXPLAIN (o) SELECT * FROM t4, t5, t6, t7, t8
...   WHERE (t4.x + t4.y = t5.x + 1) AND (t5.y = t6.y) AND (t5.x = t7.y) AND (t4.x = t7.x) AND (t8.y = 1)
...   AND (t4.y >= 1000000) AND (t4.y < 1500000) AND (t6.x >= 100000) AND (t6.x < 150000) AND (t7.y < 100);
=== OPTIMIZER ===
NestedLoopJoin { type=Inner, predicate=true }
  HashJoin { type=Inner, left_key=["#0.2", "#0.0"], right_key=["#1.1", "#1.0"] }
    HashJoin { type=Inner, left_key=["#0.3"], right_key=["#1.1"] }
      HashJoin { type=Inner, left_key=["(#0.0+#0.1)"], right_key=["(#1.0+1)"] }
        SeqScan { table=t4, filter=((#0.1>=1000000)and(#0.1<1500000)) }
        SeqScan { table=t5 }
      SeqScan { table=t6, filter=((#0.0>=100000)and(#0.0<150000)) }
    SeqScan { table=t7, filter=(#0.1<100) }
  SeqScan { table=t8, filter=(#0.1=1) }

Query优化2:Column Pruning + Common Expression Elimination

考虑如下SQL:

SELECT v, d1, d2 FROM (
  SELECT v,
         MAX(v1) AS d1, MIN(v1), MAX(v2), MIN(v2),
         MAX(v1) + MIN(v1), MAX(v2) + MIN(v2),
         MAX(v1) + MAX(v1) + MAX(v2) AS d2
    FROM t7 LEFT JOIN (SELECT v4 FROM t8 WHERE 1 == 2) ON v < v4
    GROUP BY v
);

我们可以针对性地递归优化当前节点是Projection,并且其子节点也是Projection,且子Projection的子节点是Aggregation的情况。具体思路如下:

1st Pass:裁剪Projection只保留必要列

  • 统计顶层Projection所需的列索引:即子Projection里哪些输出column会被上层用到。
  • 用这些索引裁剪掉子Projection的表达式,只保留必要的column。
  • 记得同时裁剪子Projection的输出Schema

2nd Pass:裁剪Aggregation只保留必要聚合

  • 与上一个pass类似,统计子Projection还需要哪些子Aggregation的输出column。
  • 只保留这些Aggregation表达式和类型,记录其所在的index,裁剪掉没用到的部分。
  • 注意AggregationPlanNode的输出Schema的前几列是GroupBy的结果,因此要做一下特殊处理,减去这个偏移。

3rd Pass:重复Aggregation表达式消除

  • 检查Aggregation表达式列表里是否存在重复(比如SUM(x), SUM(x)),只保留一份,避免重复计算。
  • 维护一个mapping,记录原始表达式index到新index的映射。
  • 用mapping修正上层Projection表达式中的column index引用。
  • 重新构造Aggregation节点的输出Schema,只保留用到的聚合表达式的对应column。

最后,用裁剪后的Aggregation表达式列表和新的Schema,重建Projection → Projection → Aggregation结构,返回最终的ProjectionPlanNode

一个复杂的例子

实现了以上逻辑之后,就可以完美地化简如下这种非常复杂冗余的查询:

bustub> explain (o) select v, d1 + d2, d3 + d4 - d5 from (
...     select
...         v, max(v1) as d1, min(v2) as d2,
...         max(v1) + min(v2) as d3,
...         max(v1) + max(v2) + min(v2) as d4,
...         max(v1 + v2 + v4) as d5,
...         min(v1), max(v2), min(v2), max(v1) + min(v1), max(v2) + min(v2), min(v1), max(v2), min(v2), max(v1) + min(v1), max(v2) + min(v2), min(v1), max(v2), min(v2), max(v1) + min(v1), max(v2) + min(v2), min(v1), max(v2), min(v2), max(v1) + min(v1), max(v2) + min(v2), min(v1), max(v2), min(v2), max(v1) + min(v1), max(v2) + min(v2), min(v1), max(v2), min(v2), max(v1) + min(v1), max(v2) + min(v2), min(v1), max(v2), min(v2), max(v1) + min(v1), max(v2) + min(v2), min(v1), max(v2), min(v2), max(v1) + min(v1), max(v2) + min(v2)
...     from __mock_t7 left join (select v4 from __mock_t8 where 1 == 2) on v < v4 group by v
... );
=== OPTIMIZER ===
Projection { exprs=["#0.0", "(#0.1+#0.2)", "((#0.3+#0.4)-#0.5)"] }
  Projection { exprs=["#0.0", "#0.1", "#0.2", "(#0.1+#0.2)", "((#0.1+#0.3)+#0.2)", "#0.4"] }
    Agg { types=["max", "min", "max", "max"], aggregates=["#0.1", "#0.2", "#0.2", "((#0.1+#0.2)+#0.3)"], group_by=["#0.0"] }
      NestedLoopJoin { type=Left, predicate=(#0.0<#1.0) }
        MockScan { table=__mock_t7 }
        Values { rows=0 }

可以看到,在Aggregation这一层,只计算了几个必要的聚合结果,然后交给上层的Projection进行重新组合。

Concurrency Control

最终的跑分如下:

动手之前的一些注意事项:

  • 编码时注意尽可能封装并重用一些相似的逻辑,这样可以为后续减轻工作负担。具体来说,Tuple Update和Delete这两个操作会在多处被引用到,因此最好早点放进common。
  • 时刻记得要正确维护好write set,这样在最后实现serializable isolation的时候可以避免回过头来修bug

Watermark

忽略掉官方指南中提及的O(1)算法,它指的是平均复杂度,意义不大。直接用std::set实现一个最差O(logN)的就足够用了。

MVCC与Version Chain

首先要深刻理解UndoLog的语义:如果apply了这个log,我可以把状态revert到log.ts_,也就是读取到这个历史时刻的Tuple。在CollectUndoLogs()里,对于一个给定的read_ts,我们只需要找到第一个log并满足log.ts_<=read_ts即可。如果version chain上的所有UndoLog都不能使得这个条件满足,则说明在read_ts时刻,这个Tuple还不存在。注意理清楚其他判断Tuple不存在的条件。

如果从TupleMeta读取到一个临时timestamp (ts >= TXN_START_ID),这个特殊值相当于会挡住所有对同一Tuple的并发更新,从而保证线程安全。这也是为什么UndoLog不可能包含临时timestamp:因为其他线程/事务只要看到它就自动abort了,没有机会对这个Tuple进行修改并往version chain里写入新的UndoLog

注意理解为什么要分别实现GenerateNewUndoLog()GenerateUpdatedUndoLog()这两个方法?这是因为如果在同一个事务内部,如果我们连续修改同一个Tuple,则不应该将新的UndoLog写入version chain!这里要维护的循环不变量是:每个事务对每个Tuple只产生一个新的UndoLog,这样不仅节省空间,而且有利于其他功能的正确实现。

CAS语义下的并发安全

在更新/删除Tuple时,我们都需要调用UpdateTupleAndUndoLink()以保证原子地同时更新TupleMetaTupleUndoLink。但这里的“原子性”基于的是CAS语义,所以必须传入一个检查函数,它的本质是“乐观并发控制”下的冲突检测器,以确保你要修改的Tuple在你读到和你写入之间没有被其他事务更改过。只有校验通过,才能安全地进行物理更新,否则就要abort当前事务。

这个检查函数其实非常简单:

auto check_func = [&meta](const TupleMeta &cur_meta, const Tuple &tuple, RID rid, std::optional<UndoLink>) {
  return meta == cur_meta;
};

这是基于如下两点考虑:

  • 显然ts_在读取和写回之间不能有变化,不然说明别的事务肯定改了这个Tuple
  • 对于INSERT操作,可能会插入一个已经被软删除的Tuple,因此要保证is_deleted_也没有发生变化

更新Primary Key

首先,要理解B+树上的RID一旦被插入就不能再更新,这是并发安全的来源。因为UndoLink是与RID唯一关联的,哪怕我们去查找一个当前已经被删除的key,只要它曾经在树上存在过,就必须能够访问到它的全部历史,才能符合MVCC的语义。在实际工业实现中,如果某个RID对应的所有历史版本都已经“过期”,即对所有活跃事务都不可见,且不会再被任何快照或回滚访问到,那么就可以安全地将该RID从B+树索引中物理删除了——只不过bustub在这里做了简化,不关心这部分空间浪费。

对于Primary Key的更新,我们不得不把它实现为通过TupleMeta软删除旧key,再重新插入新key。那么如果B+树上找不到新key所对应的RID怎么办?往TableHeap里再插入一个新的Tuple就好了,然后再把新的RID插入B+树。这里面有个坑点在于我们要正确abort掉例如UPDATE key SET key = 1这样的SQL:它其实是个写-写冲突,所以要想清楚怎么检测出来。

注意这里有一个小的优化:如果我们检查到Primary Key在Tuple更新前后没有发生变化,那么直接把新Tuple插入进RID即可,不需要应用上述先删除再插入的复杂逻辑。

Delete Executor的特殊处理

在实现DELETE的时候,官方指南会提到如果某个Tuple经历了INSERT => UPDATE => DELETE这种特殊sequence,需要将其meta.ts_设置为零,使其对事务外部不可见。这是因为在最后实现Serializable Verification时,我们需要遍历事务的write_set,检查是否存在读写冲突,此时就必须忽略掉这个“外部不可见”的Tuple。但是在此处,可以先忽略掉这个提示,以免把它实现在不合适的位置,后续还得改来改去浪费时间。我们最后可以把它放进tuple delete的公共逻辑里。

实时GC与ABORT逻辑

在实现Abort的时候,官方给出了两种方式,其中第二种似乎对GC更友好,但其实如果要实现real time GC的话,两者都存在同一种微妙的并发问题,都不得不退化为只有当整个系统里不存在running txn的时候才能GC。

首先,如果没有GC,那么Index Scan和Abort逻辑之间是并发安全的——CollectUndoLogs()永远不会访问到悬空的UndoLink。而一旦引入GC之后,如果Abort逻辑将undo log从version chain中移除,在某个瞬间就会出现这样一种竞态条件:

  • Index Scan里读取到的还是临时timestamp (ts >= TXN_START_ID) 以及指向要被abort的txn内部的undo link,因此,它需要回溯version chain来拿到一个对自己可见的Tuple
  • GC线程读到的则是已经被Abort逻辑更新完成的tuple meta,在回溯version chain的时候,并没有看到对当前txn的引用,就会认为它可以被GC
  • 由此,aborted txn的undo logs被清理,而Index Scan里还拿着一个已经失效的undo link,就会导致错误

那么如果我们采取另一种“丑陋”一点的实现方式,把aborted txn的undo log还留在version chain上,这样是不是可以避免GC时错误将其回收呢?很可惜,仍然是有问题的,考虑这样一种情况:当GC进行检查时,此时如果读取到watermark >= meta.ts_,那么直接就触发清理,根本没有去回溯version chain。但由于Index Scan里读取到的还是临时timestamp,不得不访问undo log,这就导致了悬空指针。

本质上这是由于读取tuple meta以及后续操作不是原子的,在多线程之间存在不一致而导致的问题。正确的解决方案是对undo logs进行原子引用计数,从而避免悬空指针问题,但在bustub的代码框架下不支持这样的修改。

Serializable Verification

由于MVCC本身已经保证了snapshot级别isolation,因此这部分只是让你实现一个基于乐观并发控制(Optimistic Concurrency Control, OCC)的事务校验逻辑来避免“写倾斜”(write skew)。只需要把官方指南的算法描述给翻译成代码即可,本质上就是检查一下当前事务所读取到的Tuple是不是被后续其他事务修改/插入/删除了,如果是的话就需要abort。记得尽量复用一下ReconstructTuple()的逻辑,以及注意一个小坑:如果在某个Tuple上找不到UndoLog代表了什么情况?应该如何处理?

Pessimistic Concurrency Control

本质上,PCC就是在写入/修改/删除某个Tuple时提前做一下检查,看自己要处理的Tuple是否命中了某个正在运行的事务的读取条件(read predicate)——如果是的话,就提前abort掉,以避免事务继续执行下去所带来的时间浪费。糙快猛的实现方案是在TransactionManager里直接存下这些predicates,每次查询时遍历一遍做求值,当然这样显然效率不高——在工业实现里会用一些更复杂的数据结构来优化这个查找过程。

PCC实现起来其实比官方指南里的OCC更简单,但是加上之后会过不了unit test,因为会有一些case要求你一定要在Commit阶段才能abort。那么如何在悲观和乐观并发控制之间自动切换呢?我们可以实时统计一下abort rate,在PCC检查函数入口处如果看到冲突率很低就直接跳过检查,反正最后由OCC兜底;冲突率高才会继续采用悲观并发控制策略,例如运行leaderboard benchmark 2的时候。


(全文完)🎉

PS:如果你有耐心看到这里,并对我的实现感兴趣,可以email找我索取代码。仅供参考学习,请勿公开发表!

comments powered by Disqus
Published:
2025-12-31
分类:
Tag: