系统设计面试笔记 第 27 章

第 27 章 数字钱包

引言

支付平台通常会提供钱包服务(Wallet Service),允许客户在应用内存放资金,之后可以随时提取。

你还可以用它购买商品和服务,或向同样使用该数字钱包(Digital Wallet)服务的其他用户转账。这比走常规支付通道更快、更便宜。

数字钱包

第 1 步:理解问题并确定设计范围

  • 候选人:我们只关注数字钱包之间的转账吗?还需要支持其他操作吗?
  • 面试官:目前先关注数字钱包之间的转账。
  • 候选人:系统需要支持每秒多少笔交易?
  • 面试官:假设是 100 万 TPS。
  • 候选人:数字钱包对正确性有严格要求。我们可以认为事务保证就足够了吗?
  • 面试官:可以。
  • 候选人:我们需要证明正确性吗?
  • 面试官:可以通过对账(Reconciliation)来做,但对账只能发现差异,无法告诉我们差异的根本原因。我们希望能够从头重放数据,以重建历史。
  • 候选人:可以假设可用性要求是 99.99% 吗?
  • 面试官:可以。
  • 候选人:需要考虑外汇兑换吗?
  • 面试官:不需要,这不在范围内。

总结一下,我们需要支持:

  • 支持两个账户之间的余额转账
  • 支持 100 万 TPS
  • 可靠性达到 99.99%
  • 支持事务
  • 支持可复现性(Reproducibility)

粗略估算

在云上部署的传统关系型数据库大约可以支持 1000 TPS。

要达到 100 万 TPS,我们需要 1000 个数据库节点。但如果每笔转账涉及两条分录(一出一入),那么实际上需要支持 200 万 TPS。

我们的设计目标之一是提高单个节点能处理的 TPS,从而减少数据库节点的数量。

单节点 TPS 节点数量
100 20,000
1,000 2,000
10,000 200

第 2 步:提出高层设计并获得认可

API 设计

在本次面试中,我们只需要支持一个接口:

POST /v1/wallet/balance_transfer - 将余额从一个钱包转到另一个钱包

请求参数:from_account、to_account、amount(用字符串表示,以免丢失精度)、currency、transaction_id(幂等键)。

响应示例:

{
    "status": "success"
    "transaction_id": "01589980-2664-11ec-9621-0242ac130002"
}

基于内存的分片方案

我们的钱包应用为每个用户账户维护账户余额。

表示它的一个好的数据结构是 map<user_id, balance>,可以用内存型的 Redis 存储来实现。

由于单个 Redis 节点无法承受 100 万 TPS,我们需要把 Redis 集群划分到多个节点上。

分区算法示例:

String accountID = "A";
Int partitionNumber = 7;
Int myPartition = accountID.hashCode() % partitionNumber;

可以用 Zookeeper 存储分区数量和 Redis 节点的地址,因为它是一个高可用的配置存储。

最后,钱包服务是一个无状态服务,负责执行转账操作。它可以轻松地水平扩展:

钱包服务

虽然这个方案解决了可扩展性问题,但它无法让我们以原子方式执行余额转账。

分布式事务

处理事务的一种方法,是在标准的、已分片的关系型数据库之上使用两阶段提交(Two-Phase Commit)协议:

基于关系型数据库的分布式事务

两阶段提交(2PC)协议的工作方式如下:

2PC 协议
  • 协调者(钱包服务)像平常一样在多个数据库上执行读写操作
  • 当应用准备提交事务时,协调者请求所有数据库对事务进行准备(prepare)
  • 如果所有数据库都回复“yes”,协调者就请求这些数据库提交事务
  • 否则,请求所有数据库中止事务

2PC 方法的缺点:

  • 由于锁竞争,性能不佳
  • 协调者是单点故障

使用 Try-Confirm/Cancel(TC/C)实现分布式事务

TC/C 是 2PC 协议的一种变体,它借助补偿事务(Compensating Transaction)来工作:

  • 协调者请求所有数据库为事务预留资源
  • 协调者收集各数据库的回复:如果都是 yes,就请求数据库执行 try-confirm;如果有 no,就请求数据库执行 try-cancel。

TC/C 与 2PC 的一个重要区别是:2PC 执行的是单个事务,而 TC/C 中是两个相互独立的事务。

TC/C 按阶段的执行方式如下:

阶段 操作 A C
1 Try 余额变动:-$1 不做任何事
2 Confirm 不做任何事 余额变动:+$1
Cancel 余额变动:+$1 不做任何事

第 1 阶段:try

try 阶段
  • 协调者在 A 的数据库中启动本地事务,把 A 的余额减少 1$
  • C 的数据库收到一条 NOP 指令,什么也不做

第 2a 阶段:confirm

confirm 阶段
  • 如果两个数据库都回复“yes”,就开始 confirm 阶段。
  • A 的数据库收到 NOP,而 C 的数据库被指示把 C 的余额增加 1$(本地事务)

第 2b 阶段:cancel

cancel 阶段
  • 如果第 1 阶段中的任何操作失败,就开始 cancel 阶段。
  • A 的数据库被指示把 A 的余额增加 1$,C 的数据库收到 NOP

2PC 与 TC/C 的对比如下:

第一阶段 第二阶段:成功 第二阶段:失败
2PC 事务尚未完成 提交/取消所有事务 取消所有事务
TC/C 所有事务都已完成(已提交或已取消) 如有需要,执行新的事务 撤销已经提交的事务

TC/C 也被称为补偿式分布式事务。高层操作在业务逻辑中处理。

TC/C 的其他特性:

  • 与数据库无关,只要数据库支持事务即可
  • 分布式事务的细节和复杂性需要在业务逻辑中处理

TC/C 的故障模式

如果协调者在执行过程中宕机,它需要恢复自己的中间状态。 这可以通过维护阶段状态表(Phase Status Table)来实现,这些表在数据库分片内被原子地更新:

阶段状态表

这张表包含哪些内容:

  • 分布式事务的 ID 和内容
  • try 阶段的状态:未发送、已发送、已收到响应
  • 第二阶段的名称:confirm 或 cancel
  • 第二阶段的状态
  • 乱序标志(稍后解释)

使用 TC/C 的一个注意事项是:当分布式事务正在执行时,会有一个短暂的时刻,各账户的状态彼此不一致:

不平衡状态

只要我们总能从这种状态中恢复,并且用户无法利用这个中间状态(例如把钱花掉),这就没有问题。 通过始终先执行扣减、后执行增加,就能保证这一点。

Try 阶段的选择 账户 A 账户 C
选择 1 -$1 NOP
选择 2(无效) NOP +$1
选择 3(无效) -$1 +$1

注意,上表中的选择 3 是无效的,因为不依赖 2PC,我们就无法保证跨不同数据库的事务以原子方式执行。

需要处理的一个边界情况是乱序执行:

乱序执行

数据库有可能在收到 try 之前先收到 cancel 操作。这种边界情况可以通过在阶段状态表中添加乱序标志来处理。 当我们收到 try 操作时,先检查乱序标志是否已设置,如果已设置,就返回失败。

使用 Saga 实现分布式事务

另一种流行的方法是使用 Saga,它是在微服务架构中实现分布式事务的一种标准做法。

它的工作方式如下:

  • 所有操作按顺序排列。每个操作在各自的数据库中都是独立的。
  • 操作从第一个到最后一个依次执行
  • 当某个操作失败时,整个流程开始回滚,通过补偿操作一直回退到起点
saga

我们如何协调这个工作流?有两种方法可选:

  • 协同(Choreography):参与 saga 的所有服务都订阅相关事件,各自完成 saga 中属于自己的那部分
  • 编排(Orchestration):由单个协调者指示所有服务按正确的顺序完成各自的工作

使用协同方式的难点在于,业务逻辑被拆分到多个以异步方式通信的服务中。 编排方式能很好地处理复杂性,因此在数字钱包系统中通常是首选方法。

TC/C 与 Saga 的对比如下:

TC/C Saga
补偿动作 在 Cancel 阶段 在回滚阶段
中心化协调 是 是(编排模式)
操作执行顺序 任意 线性
能否并行执行 能 不能(线性执行)
是否可能看到部分不一致的状态 是 是
应用逻辑还是数据库逻辑 应用 应用

主要区别在于 TC/C 可以并行执行,因此我们的决策取决于延迟要求:如果需要实现低延迟,就应该选择 TC/C 方法。

无论采用哪种方法,我们仍然需要支持审计和重放历史,以便从失败状态中恢复。

事件溯源

在现实中,数字钱包应用可能会接受审计,我们必须回答某些问题:

  • 我们是否知道任意时刻的账户余额?
  • 我们如何知道历史余额和当前余额是正确的?
  • 在代码变更之后,我们如何证明系统逻辑是正确的?

事件溯源(Event Sourcing)是一种能帮助我们回答这些问题的技术。

它由四个概念组成:

  • 命令(Command):来自现实世界的预期动作,例如从账户 A 向 B 转账 1$。命令需要有全局顺序,因此它们被放入一个 FIFO 队列。
    • 与事件不同,命令可能失败,并且由于 IO 或无效状态等原因带有一定的随机性。
    • 命令可以产生零个或多个事件
    • 事件生成可能包含外部 IO 之类的随机性。稍后会再讨论这一点
  • 事件(Event):系统中已发生事件的历史事实,例如“从 A 向 B 转账了 1$”。
    • 与命令不同,事件是在我们系统中已经发生的事实
    • 与命令类似,事件也需要有序,因此它们被放入一个 FIFO 队列
  • 状态(State):事件导致的变化。例如一个记录账户与其余额对应关系的键值存储。
  • 状态机(State Machine):驱动事件溯源过程。它主要负责校验命令,并应用事件来更新系统状态。
    • 状态机应当是确定性的,因此它不应读取外部 IO,也不应依赖随机性。
事件溯源

下面是事件溯源的动态视图:

事件溯源动态视图

对于我们的钱包服务,命令就是余额转账请求。我们可以把它们放入一个 FIFO 队列,例如 Kafka:

命令队列

完整的架构如下:

钱包服务状态机
  • 状态机从命令队列中读取命令
  • 从数据库中读取余额状态
  • 校验命令。如果有效,就为两个账户各生成一个事件
  • 读取下一个事件,并通过更新数据库中的余额(状态)来应用它

使用事件溯源的主要优势是它的可复现性。在这个设计中,所有状态更新操作都被保存为所有余额变动的不可变历史。

通过从头重放事件,总能重建历史余额。 由于事件列表是不可变的,且状态机是确定性的,我们可以保证成功重放出任何一个中间状态。

历史状态

本节开头提出的所有审计相关问题,都可以借助事件溯源来解答:

  • 我们是否知道任意时刻的账户余额?可以从头开始重放事件,直到我们关心的那个时间点
  • 我们如何知道历史余额和当前余额是正确的?可以通过从头重新计算所有事件来验证正确性
  • 在代码变更之后,我们如何证明系统逻辑是正确的?可以让不同版本的代码处理同一批事件,验证它们的结果是否完全一致

客户查询自己余额的请求可以用 CQRS 架构来处理:可以有多个只读状态机,它们基于不可变的事件列表,负责查询历史状态:

CQRS 架构

第 3 步:深入设计

在本节中,我们将探讨一些性能优化,因为我们仍然需要扩展到 100 万 TPS。

高性能事件溯源

我们要探讨的第一个优化,是把命令和事件保存到本地磁盘存储中,而不是 Kafka 这样的外部存储。

这避免了网络延迟;而且由于我们只做追加操作,这类操作在 HDD 上通常也很快。

下一个优化是把最近的命令和事件缓存在内存中,以节省从磁盘重新加载它们的时间。

在底层,我们可以借助一个叫 mmap 的命令来实现上述优化,它既把数据存储到本地磁盘,又把数据缓存在内存中:

mmap 优化

下一个可以做的优化,是用 SQLite(一种基于文件的本地关系型数据库)把状态也存储在本地文件系统中。RocksDB 也是一个不错的选择。

就我们的用途而言,我们选择 RocksDB,因为它使用日志结构合并树(Log-Structured Merge-Tree,LSM),这种结构针对写操作做了优化。 读性能则通过缓存来优化。

RocksDB 方案

为了优化可复现性,我们可以定期把快照保存到磁盘,这样就不必每次都从头重建某个状态。我们可以把快照作为大型二进制文件存储在分布式文件存储中,例如 HDFS:

快照方案

可靠的高性能事件溯源

到目前为止做的所有优化都很好,但它们让我们的服务变成了有状态的。出于可靠性考虑,我们需要引入某种形式的复制。

在此之前,我们应该分析一下系统中哪类数据需要高可靠性:

  • 状态和快照总是可以通过事件列表重新生成。因此,我们只需要保证事件列表的可靠性。
  • 有人可能认为事件列表总能从命令列表重新生成,但事实并非如此,因为命令是非确定性的。
  • 结论是,我们只需要确保事件列表的高可靠性

为了实现事件的高可靠性,我们需要把事件列表复制到多个节点上。我们需要保证:

  • 没有数据丢失
  • 日志文件中数据的相对顺序在各副本之间保持一致

为此,我们可以采用一种共识算法,例如 Raft。

在 Raft 中,有一个处于活跃状态的 leader,以及若干处于被动状态的 follower。如果 leader 宕机,会由某个 follower 接替。 只要超过一半的节点正常运行,系统就能继续运行。

Raft 复制

采用这种方法,所有节点都基于事件列表来更新状态。Raft 确保 leader 和 follower 拥有相同的事件列表。

分布式事件溯源

到目前为止,我们已经设计出了一个单节点性能高且可靠的系统。

我们还需要解决一些局限:

  • 单个 Raft 组的容量是有限的。到了某个时候,我们需要对数据分片并实现分布式事务
  • 在 CQRS 架构中,请求/响应流程很慢。客户端需要定期轮询系统,才能知道自己的钱包何时被更新

轮询不是实时的,因此用户可能要过一段时间才能得知自己余额的更新。另外,如果轮询频率太高,还可能让查询服务过载:

轮询方案

为了减轻系统负载,我们可以引入一个反向代理(Reverse Proxy),由它代表用户发送命令,并代表用户轮询响应:

反向代理

这减轻了系统负载,因为我们可以用一个请求为多个用户获取数据,但它仍然没有解决实时回执的需求。

我们可以做的最后一个改动,是让只读状态机在响应就绪后立即把它推送回反向代理。这能让用户感觉更新是实时发生的:

推送式状态机

最后,为了让系统进一步扩展,我们可以把系统分片成多个 Raft 组,并在它们之上借助协调者,通过 TC/C 或 Saga 实现分布式事务:

分片的 Raft 组

下面是在最终系统中,一个余额转账请求的生命周期示例:

  • 用户 A 向 Saga 协调者发送一个包含两个操作(A-1 和 C+1)的分布式事务。
  • Saga 协调者在阶段状态表中创建一条记录,用于跟踪该事务的状态
  • 协调者确定需要把命令发送到哪些分区。
  • 分区 1 的 Raft leader 收到 A-1 命令,对其进行校验,把它转换为事件,并复制到 Raft 组内的其他节点
  • 事件结果被同步到读状态机,读状态机再把响应推送回协调者
  • 协调者创建一条记录,表示该操作已成功,然后继续执行下一个操作 C+1
  • 下一个操作的执行过程与第一个类似:确定分区、发送命令、执行,然后读状态机推送回响应
  • 协调者创建一条记录,表示操作 2 也已成功,最后把结果通知客户端

第 4 步:总结

我们的设计演进过程如下:

  • 我们从一个使用内存型 Redis 的方案开始。这种方法的问题在于它不是持久化存储。
  • 我们转而使用关系型数据库,并在其上通过 2PC、TC/C 或分布式 saga 执行分布式事务。
  • 接着,我们引入事件溯源,使所有操作都可审计
  • 一开始我们用外部数据库和队列把数据存储到外部存储中,但这种方式性能不佳
  • 于是我们改为把数据存储在本地文件存储中,利用只追加操作的性能优势。我们还使用缓存来优化读路径
  • 上一种方法虽然性能好,但不够持久。因此,我们引入了带复制的 Raft 共识,以避免单点故障
  • 我们还采用了 CQRS,并借助反向代理代表用户管理事务的生命周期
  • 最后,我们把数据分区到多个 Raft 组中,并通过分布式事务机制(TC/C 或分布式 saga)对它们进行协调