Skip to content

并发控制

本文介绍了在分布式系统中,CRDT(Conflict-free Replicated Data Types)如何解决数据冲突。

1. CRDT

1.1 概述

CRDT,全称Conflict-free Replicated Data Type,是一种在可以在不同副本间合并、复制的特殊数据结构,具有以下特征:

  • 不同副本上的CRDT可以独立修改,不需要和其他副本协调;
  • CRDT可以自动解决不一致的情况;
  • 尽管CRDT在某个时间会不同,但是保证最终会一致;

例如现在有一个日程管理程序,我在手机上修改了日程名称,在电脑上修改了同一个日程的时间,注意,以上修改都发生在本地。当网络连通后,两个节点的数据进行同步,最终两个节点的日程一致,都是修改后的名称和时间,自动完成了合并。如下:

image-20260802132523711

图源:distributed system notes by Martin Kleppmann

CRDT又可以视为Convergent or Commutative Replicated Data Type,这是根据其实现方式来划分的:

  • Convergent Replicated Data Type:简称CvRDT,也称为state-based CRDT,即节点之间同步的是状态;
  • Commutative Replicated Data Type:简称CmRDT,也称为operation-based CRDT,即节点之间同步的是操作;

1.2 原子、对象和操作

原子(atom)是不可变的基础数据类型,根据其值唯一确定(也就是值相同就视为原子相同),原子可以在不同节点之间复制。原子可以是整数、字符串、集合等,以及对它们的不可变操作。

对象(object)是可变、可复制的数据类型,对象通常由标识符、内容(payload)、初始状态以及操作组成。内容由任意数量的原子或对象组成。有相同标识符的两个对象,但是处于不同节点,那这两个节点就是副本。

CRDT就是一系列特殊定义的object,定义了其内容、初始状态、操作。

state-based CRDT,节点之间同步的是状态,也就是内容(payload);operation-based CRDT,节点之间同步的是操作。

在分布式系统中,客户端有多个,它们通过调用对象接口中的操作来查询或修改对象状态。客户端会选择一个副本(称为源副本 source replica)进行操作。

查询操作在本地执行,也就是说,只在单个副本上完成。

更新操作分为两个阶段:

  1. 客户端首先在源副本上调用更新操作,源副本可以进行一些初始处理;
  2. 然后,该更新操作会被异步发送到所有副本,这部分称为下游处理(downstream);

1.3 state-based CRDT

state-based CRDT(CvRDT),在节点之间同步的是状态,例如,现在有一个集合的CRDT,在节点N1上集合为{A,B},在节点N2上集合为{A,C},当发生同步时,节点N1是把自身的状态{A,B}全部发送出去。用形式化的语言表述如下:

txt
-- 初始化阶段
payload Payload type; instantiated at all replicas
 	initial Initial value
 	
-- 查询操作,不会修改状态
query Query (arguments) : returns
  pre Precondition  -- 只有满足了相关条件才会返回状态
  let Evaluate synchronously, no side effects
  
-- 更新操作
update Source-local operation (arguments) : returns
  pre Precondition  -- 只有满足了相关条件才会修改状态
  let Evaluate at source, synchronously  -- 在本地同步修改状态
  Side-effects at source to execute synchronously  -- 异步同步状态
  
-- 合并状态
merge (value1, value2) : payload mergedValue
  LUB merge of value1 and value2, at any replica

merge 操作本质上就是计算两个状态的最小上界(Least Upper Bound,LUB ),何为最小上界?即合并结果同时包含 value1value2 的所有信息,并且不存在比它更小的状态也能同时包含两者。例如:

txt
value1 = {A, B}

value2 = {B, C}

merge(value1, value2) = {A, B, C}

{A, B, C} 是两者的上界,因为它包含两个集合;同时它又是最小的,因为不存在更小的集合同时包含 {A,B} 和 {B,C}。

由于 LUB 天然满足交换律、幂等性和结合律,所以无论状态以何种顺序、重复多少次进行合并,最终都会收敛到同一个结果。

  • 交换律:merge(value1, value2)merge(value2, value1)的结果相同;
  • 幂等性:merge(value, value)最终结果为value
  • 结合律:merge(merge(value1, value2), value3)merge(value1, merge(value2, value3))的结果相同;

CvRDT 不要求网络可靠、有序(也就是说CvRDT可以使用unreliable best-effort-broadcast),只要求最终所有副本之间能够不断交换状态,并且通信网络最终连通。因为 merge 满足交换律、结合律和幂等性,即使状态乱序、重复传播,所有副本最终也会收敛。

CvRDT的缺点是,如果状态过大,那么广播消息会很大。

1.4 operation-based CRDT

operation-based CRDT(CmRDT),在节点之间同步的是操作,例如,现在有一个集合的CRDT,初始化都是空集合{},节点N1接收到操作add(A),然后节点N1异步将这个操作add(A)广播出去,其他节点N2接收到这个操作后,将其应用到本地状态实现同步。用形式化的语言表述如下:

txt
-- 初始化阶段,和CvRDT相同
payload Payload type; instantiated at all replicas
  initial Initial value

-- 查询操作,和CvRDT相同
query Source-local operation (arguments) : returns
  pre Precondition
  let Execute at source, synchronously, no side effects

-- 更新操作
update Global update (arguments) : returns
	-- 第一阶段 在本地生成操作,不会修改状态
  atSource (arguments) : returns
	  pre Precondition at source
    let 1st phase: synchronous, at source, no side effects
  -- 第二阶段 修改本地状态,异步广播操作
  downstream (arguments passed downstream)
    pre Precondition against downstream state
    2nd phase, asynchronous, side-effects to downstream state

在CmRDT中,同步操作有两个问题:

  • 消息重复:假设现在CRDT是一个计数器,如果操作是add(1),假设这条操作消息因为网络问题,在同步时重复发送,接收方如果没有去重机制,就会导致最终状态不一致。
  • 消息顺序:假设现在CRDT是文本序列,如果现在在同一个节点上有两个操作insert(hello )insert(world),如果两条消息在其余节点上交付顺序不同,那么就可能造成两种结果hello worldworldhello ,也会导致最终状态不一致。

因此,CmRDT实现要求底层广播模型为causal broadcast,并且需要设计消息去重机制(可以通过底层广播模型,也可以在CRDT实现层设计去重机制)。

并且,对于没有因果关系的操作,允许可交换,也就是说,操作1和操作2在不同节点上,允许执行顺序不同,对最终结果无影响。

CmRDT相比CvRDT的优点在于传输的消息体比较小,但是缺点在于需要因果广播系统,并且需要设计消息去重机制,不能容忍消息丢失。

1.5 CRDTs

本小结介绍一些常见的CRDT。

1.5.1 Counters

Counters(计数器)是一种整数CRDT,可以对其进行incrementdecrement操作来更新计数器,也可以进行value操作来查询计数值。在实际应用中,计数器可以用来统计当前在线人数等。

根据支持的功能,计数器有不同类型:G-Counter、PN-Counter。

G-Counter:Grow-only Counter,只增计数器,即计数值只会增加不会减少

state-base G-Counter实现如下:

txt
-- 初始化阶段,每个副本都在本地保存一个整数
-- 在实现中,可以使用k-v结构,key为节点ID,value为计数值
payload integer[n] P
 	initial [0, 0, . . . , 0]

-- 增加计数值
update increment ()
  let g = myID()   -- 获取当前节点ID
  P[g] := P[g] + 1  -- 将当前节点ID对应的计数值+1
  
-- 查询计数值 就是把每个节点的计数值加起来
query value () : integer v
	integer v = 0
	for p in P
		v = v + p

-- 合并 X为本地计数值列表 Y为接收到的计数值列表
-- 取两者的较大值作为合并结果
merge (X, Y) : payload Z
  let ∀i ∈ [0, n − 1] : Z.P[i] = max(X.P[i], Y.P[i])

PN-Counter:Positive-Negative Counter,可加可减的计数器

state-based PN-Counter实现如下:

txt
-- 每个副本维护两个整形值:P表示加的值,N表示减的值
payload integer[n] P , integer[n] N 
  initial [0, 0, . . . , 0], [0, 0, . . . , 0]

update increment ()
  let g = myID() 
  P[g] := P[g] + 1

update decrement ()
  let g = myID()
  N[g] := N[g] + 1
  
query value () : integer v
  let v = Pi P[i] − Pi N[i]
 
merge (X, Y ) : payload Z
  let ∀i ∈ [0, n − 1] : Z.P [i] = max(X.P [i], Y.P [i])
  let ∀i ∈ [0, n − 1] : Z.N[i] = max(X.N[i], Y.N[i])

1.5.2 Registers

Registers(寄存器)是一种支持更新(assign)操作和查询(value)操作的CRDT。

如果没有特殊处理,并发更新是不可交换的,也就是说,两个节点同时更新:节点1将寄存器的值更新为A(assign(A)),节点2将寄存器的值更新为(assign(B)

两个更新的顺序不同,会造成结果不同:

txt
assign(A) assign(B) --> B
assign(B) assign(A) --> A

有两种方案解决并发更新:LWW-Register和MV-Register

LWW-Register,全称为Last-Writer-Wins Register ,最后写入值胜出,该方案的实现是给每个操作一个“时间戳”(这里的时间戳并不是指物理时间,也不是单纯指逻辑时间),可以通过这些时间戳给更新操作排出唯一序列,之后就能决定哪个操作是最后写入值了。

例如,我们可以把时间戳定义为{Lamport Clock,NodeId},首先比较Lamport时钟,如果Lamport时钟相同,则比较NodeId。

state-based LWW-Register实现如下:

txt
-- 初始化阶段 X表示某种类型,例如Integer\String\Set
payload X x, timestamp t 
 initial ⊥, 0

-- 更新操作 now()并不是获取当前物理时间戳,而是获取唯一的逻辑时间戳,
-- 例如当前为{Lamport Clock,NodeId},那now()返回{Lamport Clock + 1,NodeId}
update assign (X w)
  x, t := w, now() -- Timestamp, consistent with causality

-- 查询
query value () : X w
  let w = x

-- 合并操作 比较时间戳,谁的时间戳更大就用谁的值
merge (R, R′) : payload R′′
  if R.t ≤ R′.t then R′′.x, R′′.t = R′.x, R′.t
  else R′′.x, R′′.t = R.x, R.t

operation-based LWW-Register实现如下:

txt
-- 初始化阶段,与state-based 实现相同
payload X x, timestamp t 
	initial ⊥, 0
	
-- 查询
query value () : X w
	let w = x
	
-- 更新
update assign (X x′)
	atSource () t′
		let t′ = now() -- 获取时间戳,例如当前为{Lamport Clock,NodeId},那now()返回{Lamport Clock + 1,NodeId}
	downstream (x′, t′) 
		if t < t′ then   -- 只有t'大于当前本地时间戳,才更新值
			x, t := x′, t′

在operation-based 实现中,并没有体现广播消息以及接收到广播消息后的操作,其实downstream(x', t')就可以视为接收到广播消息后的操作,那反推广播消息时,构造的消息体就是(x', t'),所以完整的assign应该如下:

txt
update assign (X x′)
	atSource () t′
		let t′ = now() -- 获取时间戳,例如当前为{Lamport Clock,NodeId},那now()返回{Lamport Clock + 1,NodeId}
	downstream (x′, t′) 
		if t < t′ then   -- 只有t'大于当前本地时间戳,才更新值
			x, t := x′, t′
			broadcast(x, t)  -- x,t已经是更新后的值了

MV-Register,全称为Multi-Value Register,针对并发更新操作,选择同时保留并发值。

要判断两个操作是否并发,可以使用向量逻辑时钟(vector clock)。因此,内容payload是一个集合:

S={(x,V)}

其中:

  • x:当前寄存器保存的值
  • V:这个值对应的版本向量

stated-based MV-Register实现如下:

txt
-- 初始化
payload set S ⊲ set of (x, V ) pairs; x ∈ X; V its version vector
	initial {(⊥, [0, . . . , 0])}
	
-- 生成新的版本向量
query incVV () : integer[n] V ′
	let g = myID()
	let V = {V |∃x : (x, V ) ∈ S}  -- 从当前状态 S 中提取所有版本向量。
	let V ′ = [ max(V [j]) ] j != g  -- 除了自己节点 g:每一位取最大值。
	let V ′[g] = max(V [g]) + 1  -- 自己节点,取最大值 +1
	
-- 赋值操作,R是集合
update assign (set R) ⊲ set of elements of type X
	let V = incVV ()  -- 计算新版本
	S := R × {V }   -- 计算新值:笛卡尔积
	
-- 查询操作,返回当前节点的值
query value () : set S′
	let S′ = S

-- 合并操作,基于版本向量进行偏序比较,去除被覆盖的旧版本,只保留并发或最新的版本集合
merge (A, B) : payload C
	-- 对于集合A中的值,只要该值的版本没有被集合B中的版本覆盖,就保留 
	let A′ = {(x, V ) ∈ A|∀(y, W) ∈ B : V k W ∨ V ≥ W}  
	-- 对于集合B中的值,只要该值的版本没有被集合A中的版本覆盖,就保留 
	let B′ = {(y, W) ∈ B|∀(x, V ) ∈ A : W k V ∨ W ≥ V }
	-- 合并保留的值
	let C = A′ ∪ B′

举一个例子来理解,假设寄存器值类型为整型,有两个节点N1和N2。

初始化:

txt
N1: {(null, [0,0])}
N2: {(null, [0,0])}

赋值:N1写入值1,N2写入值2

txt
N1: {(1, [1,0])}
N2: {(2, [0,1])}

然后进行合并:由于[1,0]和[0,1]互不覆盖,所以都保留:

txt
N1: {(1, [1,0]), (2,[0,1])}
N2: {(1, [1,0]), (2,[0,1])}

之后N1再次写入新值3:

txt
N1: {(3, [2,1])}
N2: {(1, [1,0]), (2,[0,1])}

再次合并,[2,1]同时覆盖[1,0]和[0,1],最终结果:

txt
N1: {(3, [2,1])}
N2: {(3, [2,1])}

其他数据类型可参考资料[1]和[3]。

参考资料

[1] https://inria.hal.science/inria-00555588/document

[2] https://en.wikipedia.org/wiki/Conflict-free_replicated_data_type

[3] https://docs.yjs.dev/