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 将分布式系统中最困难的问题 —— 并发操作下的数据一致性 —— 编码为一个优雅的数学结构。但理论到生产还有很长的路要走:

  1. 选择合适的数据结构:计数器用 PN-Counter,集合用 OR-Set,文本用 RGA/YATA
  2. 设计好传输层:同步协议决定了实际延迟和带宽开销
  3. 处理 GC 和 Tombstones:这是大多数自制 CRDT 系统最终失败的原因
  4. 混合架构不可避免:纯 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/
点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部