在 上篇文章 中,我们主要介绍了 Pump Server 的上线过程、gRPC API 实现、以及下线过程和相关辅助机制,其中反复提到了 Pump Storage 这个实体。本文就将介绍 Pump Storage 的实现,其主要代码在 pump/storage 文件夹中。
Storage interface
WriteBinlog、
GC和
PullCommitBinlog函数,我们将在下文具体介绍。Storage 的接口定义如下:
type Storage interface {
WriteBinlog 写入 binlog 数据到 Storage
WriteBinlog(binlog *pb.Binlog) error
GC 清理 tso 小于指定 ts 的 binlog
GC(ts int64)
GetGCTS 返回最近一次触发 GC 指定的 ts
GetGCTS() int64
AllMatched 返回是否所有的 P-binlog 都和 C-binlog 匹配
AllMatched() bool
MaxCommitTS 返回最大的 CommitTS,在这个 TS 之前的数据已经完备,可以安全的同步给下游
MaxCommitTS() int64
GetBinlog 指定 ts 返回 binlog
GetBinlog(ts int64) (binlog *pb.Binlog, err error)
PullCommitBinlog 按序拉 commitTs > last 的 binlog
PullCommitBinlog(ctx context.Context, last int64) <-chan []byte
Close 安全的关闭 Storage
Close() error
}
Append
初始化
NewAppendWithResolver函数中实现的,首先初始化 Valuelog、goleveldb 等组件,然后启动处理写入 binlog、GC、状态维护等几个 goroutine。
WriteBinlog
WriteBinlog由 Pump Server 调用,用于写入 binlog 到本地的持久化存储中。在 Append 实现的
WirteBinlog函数中,binlog 在编码后被传入到
Append.writeChChannel 由专门的 goroutine 处理:
toKV := append.writeToValueLog(writeCh)
go append.writeToSorter(append.writeToKV(toKV))
Append.writeCh后将按如下顺序流经数个处理流程:

writeToValueLog中:
// valuePointer 定义
type valuePointer struct {
Fid 是 valuelog 文件 Id
Fid uint32
Offset 是 pointer 指向的 valuelog 在文件中的偏移量
Offset int64
}
Append.writeCh读出的 binlog,批量写入到 ValueLog 组件中。我们可以将 ValueLog 组件看作一种由
valuePointer映射到 binlog 的持久化键值存储实现,我们将在下一篇文章详细介绍 ValueLog 组件。
writeBatchToKV中,Append 将 binlog 的 tso 作为 Key,
valuePointer作为 Value 批量写入 Metadata 存储中,在目前的 Pump 实现中,我们采用 goleveldb 作为 Metadata 存储数据库。由于 goleveldb 的底层是数据结构是 LSM-Tree,存储在 Metadata 存储的 binlog 相关信息已经天然按 tso 排好序了。
writeToSorter呢?这里和《TiDB Binlog 架构演进与实现原理》一文提到的 Binlog 工作原理有关:
TiDB 的事务采用 2pc 算法,一个成功的事务会写两条 binlog,包括一条 Prewrite binlog 和 一条 Commit binlog;如果事务失败,会发一条 Rollback binlog。

PullCommitBinlog
GC
gc)。
注:由于生产环境中发现用户有时会关闭了 drainer 却没有使用 binlogctl 将相应 drainer 节点标记为 offline,导致 Pump Storage 的数据一直在膨胀,不能 GC。因此在 v3.0.1、v2.1.15 后无论 Binlog 是否已经同步到下游,都会正常进入 GC 流程。
[0,GcTso]这个范围内的 Metadata,每 1024 个 KVS 作为一批次进行删除:
for iter.Next() && deleteBatch < 100 {
batch.Delete(iter.Key())
deleteNum++
lastKey = iter.Key()
if batch.Len() == 1024 {
err := a.metadata.Write(batch, nil)
if err != nil {
log.Error("write batch failed", zap.Error(err))
}
deletedKv.Add(float64(batch.Len()))
batch.Reset()
deleteBatch++
}
}
if l0Num >= l0Trigger {
log.Info("wait some time to gc cause too many L0 file", zap.Int("files", l0Num))
if iter != nil {
iter.Release()
iter = nil
}
time.Sleep(5 * time.Second)
continue
}
在示例代码的 doGCTS 函数中存在一个 Bug,你发现了么?欢迎留言抢答。
小结
💡 文中划线部分均有跳转,请点击【阅读原文】查看原版
TiDB 源码阅读系列文章
TiDB Binlog(github.com/pingcap/tidb-binlog)组件用于收集 TiDB 的 binlog,并准实时同步给下游,如 TiDB、MySQL 等。该组件在功能上类似于 MySQL 的主从复制,会收集各个 TiDB 实例产生的 binlog,并按事务提交的时间排序,全局有序的将数据同步至下游。利用 TiDB Binlog 可以实现数据准实时同步到其他数据库,以及 TiDB 数据准实时的备份与恢复。我们希望通过《TiDB Binlog 源码阅读系列文章》帮助大家理解和掌握这个项目,也有助于我们和社区共同进行 TiDB Binlog 的设计、开发和测试。





