CRDT:冲突无关数据类型 — 实时协作与分布式系统的底层革命
当你打开 Google Docs 与千里之外的同事实时协作编辑时,当你在 Notion 中看到朋友的光标在页面上来回移动时,当 Redis 集群在异地多活架构中保持最终一致时——背后运作的,都是一种名为 CRDT(Conflict-free Replicated Data Type,冲突无关数据类型) 的数学结构。
一、为什么我们需要 CRDT?
在分布式系统中,CAP 定理告诉我们:当一个网络分区发生时,必须在一致性和可用性之间做出选择。传统系统(如关系型数据库)倾向于 CP —— 宁可拒绝写入也要保证强一致;但越来越多的现代应用(在线文档、多人游戏、协同白板)选择 AP —— 允许局部写入,通过某种机制最终收敛。
CRDT 正是 AP 阵营的理论基石。它保证:无论节点以何种顺序、何种时序接收到何种操作,最终所有副本都会收敛到完全相同的状态,无需中心化协调器、无需锁、无需冲突解决。
这听起来像是魔法,但它不是 —— 它是格论(Lattice Theory)在计算机科学中最优雅的应用之一。
二、数学基石:半格(Semilattice)与单调性
CRDT 的数学基础建立在 连接半格(Join Semilattice) 之上。
一个半格是一个偏序集 (S, ≤),其中任意两个元素都存在一个最小上界(join,写作 ∨)。最常见的例子就是幂集上的并集运算:{a} ∨ {b} = {a, b}。
CRDT 利用半格的三个关键性质:
- 交换律:a ∨ b = b ∨ a —— 合并顺序无关
- 结合律:(a ∨ b) ∨ c = a ∨ b ∨ c —— 分组不影响结果
- 幂等律:a ∨ a = a —— 重复操作无害
这意味着,只要所有节点最终收到了所有更新,无论顺序如何,它们都会收敛到相同的状态 —— 这就是 "冲突无关" 的本质。
三、两大范式:State-based vs Operation-based
CRDT 家族分为两大流派,各有千秋。
3.1 State-based CRDT(CvRDT)
基于状态的 CRDT 将整个数据结构作为状态来同步。每个副本维护自己的状态,定期将完整状态发送给其他副本,接收方通过 merge 操作合并。
核心要求: 状态空间构成一个半格,merge 操作对应 join。
class GCounter:
"""基于状态的 G-Counter(只增计数器)"""
def __init__(self, num_nodes):
self.node_id = id(self)
self.payload = [0] * num_nodes # 每个节点维护一个槽位
def increment(self):
self.payload[self.node_id] += 1
def value(self):
return sum(self.payload)
def merge(self, other):
"""Join 操作:逐元素取最大值"""
for i in range(len(self.payload)):
self.payload[i] = max(self.payload[i], other.payload[i])
def query(self):
return sum(self.payload)
3.2 Operation-based CRDT(CmRDT)
基于操作的 CRDT 不传输完整状态,而是传输操作本身。每个操作为因果广播(causal broadcast),保证送达后每个副本只执行一次。
核心要求: 操作必须可交换(commute),或在所有因果历史上一致顺序下等价。
class LWWRegister:
"""基于操作的 Last-Write-Wins 寄存器"""
def __init__(self, node_id):
self.node_id = node_id
self.timestamp = (0, node_id) # (counter, node_id) 构成全序
self.value = None
def set(self, value, counter):
op = {
'counter': counter,
'node_id': self.node_id,
'value': value,
'type': 'set'
}
self.apply(op)
return op # 广播给其他节点
def apply(self, op):
remote_ts = (op['counter'], op['node_id'])
if remote_ts > self.timestamp:
self.timestamp = remote_ts
self.value = op['value']
def read(self):
return self.value
3.3 何时选择哪种?
| 维度 | State-based | Operation-based |
|---|---|---|
| 网络开销 | 传输完整状态 | 仅传输增量操作 |
| 一致性保证 | 无需额外机制 | 需要因果广播基础设施 |
| 实现复杂度 | 简单直接 | 需保证 exactly-once投递 |
| 适用场景 | 节点数少、状态简短 | 高频增量变更、节点多 |
| 容错性 | 对丢包容忍(最终补齐) | 操作丢失即为永久丢失 |
在实际工程中,两者常混合使用:用 State-based 协议做周期性全量同步兜底,用 Operation-based 做实时增量传输。
四、经典数据结构:从计数器到序列
4.1 PN-Counter(增减计数器)
G-Counter 只能递增,PN-Counter 通过组合两个 G-Counter 实现增减:
class PNCounter:
def __init__(self, num_nodes):
self.inc = GCounter(num_nodes) # 递增计数
self.dec = GCounter(num_nodes) # 递减计数
def increment(self):
self.inc.increment()
def decrement(self):
self.dec.increment()
def value(self):
return self.inc.query() - self.dec.query()
def merge(self, other):
self.inc.merge(other.inc)
self.dec.merge(other.dec)
4.2 OR-Set(观察移除集合)
经典的 Set 在分布式环境下会遇到"删除胜出还是添加胜出"的问题。OR-Set 通过为每个元素分配唯一标签(unique tag)巧妙解决:
class ORSet:
"""Observed-Remove Set:正确处理并发添加/删除"""
def __init__(self):
self.adds = {} # {element: set(tags)}
self.removes = {} # {element: set(tags)} —— 仅记录已观察到的添加
def add(self, elem, tag):
if elem not in self.adds:
self.adds[elem] = set()
self.adds[elem].add(tag)
def remove(self, elem):
"""只移除当前已观察到的所有标签"""
if elem in self.adds:
self.removes[elem] = self.adds[elem].copy()
def contains(self, elem):
"""元素存在当且仅当有添加标签未被移除"""
if elem not in self.adds:
return False
effective = self.adds[elem] - self.removes.get(elem, set())
return len(effective) > 0
def merge(self, other):
# 合并 adds:并集
for elem, tags in other.adds.items():
if elem not in self.adds:
self.adds[elem] = set()
self.adds[elem] |= tags
# 合并 removes:并集(只增不减)
for elem, tags in other.removes.items():
if elem not in self.removes:
self.removes[elem] = set()
self.removes[elem] |= tags
4.3 RGA(Replicated Growable Array)
用于实现实时协作编辑的文本序列。RGA (Replicated Growable Array) 通过为每个字符分配全局唯一的 ID(基于 Lamport 时钟 + 节点 ID),保证所有插入操作最终收敛为相同序列。
class RGANode:
def __init__(self, id, value, deleted=False, ts=None):
self.id = id # (counter, node_id) 全局唯一
self.value = value # 字符
self.deleted = deleted
self.ts = ts # 用于双向链表排序
class RGASequence:
"""用于实时文本协作的 RGA 序列"""
def __init__(self, node_id):
self.node_id = node_id
self.counter = 0
self.sequence = [] # 有序的 RGANode 列表
def generate_id(self):
self.counter += 1
return (self.counter, self.node_id)
def local_insert(self, pos, char):
"""在本地位置 pos 处插入字符"""
new_id = self.generate_id()
node = RGANode(id=new_id, value=char)
if pos < len(self.sequence):
# 在已有节点间插入:继承前驱的时序戳
prev_ts = self.sequence[pos - 1].id if pos > 0 else (0, 0)
node.ts = (prev_ts[0] + 1, self.node_id)
else:
node.ts = (new_id[0], self.node_id)
self.sequence.insert(pos, node)
return {'op': 'insert', 'node': node}
def local_delete(self, pos):
"""标记删除:逻辑删除,不物理移除"""
node = self.sequence[pos]
node.deleted = True
return {'op': 'delete', 'id': node.id}
def apply_remote(self, op):
"""应用远程操作"""
if op['op'] == 'insert':
remote_node = op['node']
# 基于 timestamp 的确定性排序插入
insert_pos = self._find_insert_position(remote_node)
self.sequence.insert(insert_pos, remote_node)
elif op['op'] == 'delete':
for node in self.sequence:
if node.id == op['id']:
node.deleted = True
break
def _find_insert_position(self, new_node):
"""确定性插入位置查找"""
for i, node in enumerate(self.sequence):
if node.id == new_node.id:
return i # 已存在
if self._compare_timestamp(node.ts, new_node.ts) > 0:
return i
return len(self.sequence)
def _compare_timestamp(self, ts1, ts2):
"""全序比较"""
if ts1[0] != ts2[0]:
return -1 if ts1[0] < ts2[0] else 1
return -1 if ts1[1] < ts2[2-1] else (1 if ts1[1] > ts2[1] else 0)
def read(self):
return ''.join(n.value for n in self.sequence if not n.deleted)
五、生产级实现:Yjs 与 Automerge
了解了理论基础后,让我们看看真正跑在生产环境中的 CRDT 库。
5.1 Yjs:性能之王
Yjs 是用 JavaScript 编写的 CRDT 库,号称 "最快最成熟的 CRDT 实现"。其核心数据结构 Y.Array、Y.Map 基于 YATA (Yet Another Transformation Approach) 算法。
import * as Y from 'yjs'
// 创建文档
const doc = new Y.Doc()
const yMap = doc.getMap('content')
// 本地写入
yMap.set('title', 'CRDT 实战')
yMap.set('author', 'yebinbing')
// 模拟分布式节点
const remoteDoc = new Y.Doc()
const remoteMap = remoteDoc.getMap('content')
// 同步:交换 state vectors 获取 diff
function sync(docA, docB) {
const stateVectorA = Y.encodeStateVector(docA)
const diffB = Y.encodeStateAsUpdate(docB, stateVectorA)
Y.applyUpdate(docA, diffB)
}
// 本地修改
yMap.set('title', 'CRDT 深度实战')
// 远程修改(并发)
remoteMap.set('tag', 'distributed-systems')
// 双向同步
sync(doc, remoteDoc)
sync(remoteDoc, doc)
// 两个文档内容完全一致
console.log(yMap.get('title')) // "CRDT 深度实战"
console.log(yMap.get('tag')) // "distributed-systems"
5.2 Automerge:Elegant API
Automerge 由 Ink & Switch 团队开发,更接近 Git 风格的变更模型:
let doc1 = Automerge.init()
doc1 = Automerge.change(doc1, '初始化', doc => {
doc.items = ['apple', 'banana']
})
// 创建分支(fork)
let doc2 = Automerge.init()
doc2 = Automerge.merge(doc2, doc1) // 初始合并
// 并发修改
doc1 = Automerge.change(doc1, '添加 cherry', doc => {
doc.items.push('cherry')
})
doc2 = Automerge.change(doc2, '添加 date', doc => {
doc.items.push('date')
})
// 合并:自动收敛
const merged = Automerge.merge(doc1, doc2)
console.log(merged.items) // ['apple', 'banana', 'cherry', 'date']
5.3 性能对比:谁更快?
| 指标 | Yjs | Automerge |
|---|---|---|
| 首次插入延迟 | < 1ms | ~3ms |
| 合并 1000 个并发操作 | ~2ms | ~150ms |
| 编码后体积(1000 ops) | ~8KB | ~25KB |
| 文档历史追踪 | 不支持 | 原生支持 |
| 浏览器兼容性 | 优秀 | 良好 |
Yjs 的优势源于 YATA 算法的扁平化编码,以及 lib0 底层的高效二进制编解码。Automerge 因需要保留完整变更历史(类似 Git),在体积和延迟上有所妥协。
六、实战:构建一个多人在线计数器
让我们从零构建一个最小可用(MVP)的实时协作计数器,演示 CRDT 在真实场景中的落地。
6.1 后端:基于 WebSocket 的增量广播
import asyncio
import websockets
import json
from crdt_gcounter import GCounter # 使用前面定义的 G-Counter
class CRDTServer:
def __init__(self):
self.clients = set()
self.counter = GCounter(num_nodes=0) # 动态扩展
self.node_counter = 0
async def register(self, websocket):
self.clients.add(websocket)
node_id = self.node_counter
self.node_counter += 1
# 为现有 counter 扩展新节点位
self.counter.payload.append(0)
# 发送当前快照
await websocket.send(json.dumps({
'type': 'snapshot',
'node_id': node_id,
'payload': self.counter.payload
}))
async def unregister(self, websocket):
self.clients.discard(websocket)
async def broadcast(self, sender, message):
"""广播增量操作到所有其他客户端"""
response = json.dumps(message)
await asyncio.gather(
*[client.send(response) for client in self.clients if client != sender],
return_exceptions=True
)
async def handler(self, websocket, path):
await self.register(websocket)
try:
async for raw in websocket:
msg = json.loads(raw)
if msg['type'] == 'increment':
node_id = msg['node_id']
self.counter.payload[node_id] += 1
await self.broadcast(websocket, msg)
elif msg['type'] == 'merge_request':
await websocket.send(json.dumps({
'type': 'snapshot',
'payload': self.counter.payload
}))
finally:
await self.unregister(websocket)
server = CRDTServer()
start_server = websockets.serve(server.handler, 'localhost', 8765)
asyncio.get_event_loop().run_until_complete(start_server)
asyncio.get_event_loop().run_forever()
6.2 前端:React 集成
import React, { useState, useEffect } from 'react';
function CRDTMonitor() {
const [payload, setPayload] = useState([]);
const [value, setValue] = useState(0);
const [nodeId, setNodeId] = useState(null);
const [ws, setWs] = useState(null);
useEffect(() => {
const socket = new WebSocket('ws://localhost:8765');
socket.onopen = () => setWs(socket);
socket.onmessage = (event) => {
const msg = JSON.parse(event.data);
if (msg.type === 'snapshot') {
if (msg.node_id !== undefined) {
setNodeId(msg.node_id);
}
setPayload(prev => {
const merged = prev.length > 0
? Math.max(...prev)
: 0;
// 逐元素取最大值
const newPayload = msg.payload.map(
(v, i) => Math.max(v, prev[i] || 0)
);
setValue(newPayload.reduce((a, b) => a + b, 0));
return newPayload;
});
} else if (msg.type === 'increment') {
setPayload(prev => {
const i = msg.node_id;
if (i !== undefined && prev[i] !== undefined) {
prev[i] = Math.max(prev[i], msg.count || prev[i] + 1);
}
setValue(prev.reduce((a, b) => a + b, 0));
return [...prev];
});
}
};
return () => socket.close();
}, []);
const handleClick = () => {
if (ws && nodeId !== null) {
ws.send(JSON.stringify({
type: 'increment',
node_id: nodeId
}));
}
};
return (
<div>
<h2>CRDT Global Counter: {value}</h2>
<p>Payload: [{payload.join(', ')}]</p>
<p>Node ID: {nodeId}</p>
<button onClick={handleClick}>Increment</button>
</div>
);
}
七、挑战与局限性
CRDT 虽然强大,但并非银弹。
7.1 垃圾回收问题
OR-Set 中的 add 标签一旦写入就永远保留。随着时间推移,内存会无限增长。解决方案包括:
- 墓碑机制(Tombstones):保留删除标记以处理延迟到达的并发添加,给 GC 带来挑战
- 因果一致性(Causal Stability):当所有副本都已确认收到某个操作后,才能清理相关数据
- 有界 CRDT(Bounded CRDT):如 Add-Wins Set 的变体,限制标签集合大小
7.2 大文档性能
当文档包含数万字符时,基于 Y.Text(RGA)的全文遍历变得缓慢。优化手段包括:
- 将大文档分块为多个 Y.Text 子节点
- 使用索引结构(如基于 Y.Map 的 B-Tree)辅助查找
- 懒加载:只下载可视区域内的数据
7.3 语义约束
CRDT 只能处理在数学上可表示的操作。当你需要业务级语义("库存不能为负"、"转账总额守恒")时,纯 CRDT 无能为力。此时需要:
- 混合方案:CRDT 处理并发写入 + 应用层校验处理业务规则
- Hybrid Logical Clocks(HLC):提供近似全局时钟,辅助业务层决策
- RGA + Effect Schemes:部分系统通过"操作语义标记"编码业务意图
八、生态全景:CRDT 在真实产品中
| 产品/项目 | 使用的 CRDT | 应用场景 |
|---|---|---|
| Redis (CRDTs) | 注册/计数器/有序集合副本 | 多活 Redis 集群 |
| Riak | OR-Set / LWW | 分布式 KV 存储 |
| SoundCloud (Roshi) | LWW-Element-Set | 时间线系统 |
| Apple Notes | 协作编辑 | 实时文档协作 |
| Figma | 自定义 CRDT | 图形编辑协同 |
| Notion | 混合 CRDT + OT | 块编辑器协作 |
| Automerge | 自定义 JSON CRDT | 通用 JSON 数据协作 |
| Yjs-based (Liveblocks) | YATA CRDT | 协同白板/文档 |
| Legend State (React) | Observable-CRDT | 前端状态管理 |
尤其值得关注的是 Figma 的做法:他们将 CRDT 用于图形编辑器的核心数据结构,每个图层(Layer)的属性都是独立的 CRDT 操作 —— 但最终布局约束仍然由中心化服务器计算。这种 "CRDT + Centralized Authority" 的混合架构很可能是未来协作系统的范式。
九、总结:从理论到生产
CRDT 将分布式系统中最困难的问题 —— 并发操作下的数据一致性 —— 编码为一个优雅的数学结构。但理论到生产还有很长的路要走:
- 选择合适的数据结构:计数器用 PN-Counter,集合用 OR-Set,文本用 RGA/YATA
- 设计好传输层:同步协议决定了实际延迟和带宽开销
- 处理 GC 和 Tombstones:这是大多数自制 CRDT 系统最终失败的原因
- 混合架构不可避免:纯 CRDT 不适用所有场景,需要与 OT、中心化协调、HLC 结合
当你的团队下次争论 "该用 CRDT 还是 OT(Operational Transformation)" 时,记住 CRDT 的数学保证比 OT 更强:它从数据结构层面消除了冲突可能,而不是在算法层面解决冲突。当然,这也意味着 CRDT 的设计约束更严格 —— 并非所有操作都能表示为 join-semilattice 上的单调函数。
选择 CRDT,你付出的代价是设计的灵活度;你获得的,是无论网络如何分裂、无论消息以何种诡异顺序到达,系统最终都会优雅收敛的数学保证。
参考资料:
- Shapiro, M., et al. (2011). "A comprehensive study of Convergent and Commutative Replicated Data Types."
- Yjs Documentation: https://docs.yjs.dev/
- Automerge Repo: https://github.com/automerge/automerge
- CRDT.tech: https://crdt.tech/

发表评论 取消回复