Riak 的 CRDT 数据类型一共包含五种:计数器(counter)、集合(set)、映射(map)、寄存器(register)和标志位(flag)。每一种类型都自带冲突合并规则,当多个副本同时被修改时,Riak 会按照这些规则自动合并,而不是简单地用最新写入覆盖旧值。这意味着在 Node.js 中操作这些类型时,不能沿用传统键值存储的读取、修改、写回思路,而要围绕上下文(context)和类型专用的更新操作来设计代码。理解这一点,是顺利使用 Riak 数据类型的关键一步。

一、CRDT 与普通键值存储的本质区别
普通键值存储通常采用 last-write-wins 策略,服务器用时间戳判断谁是最新写入,旧数据会被覆盖。这在单节点环境下问题不大,但在分布式集群中,不同节点之间的时钟可能不一致,网络分区期间多个客户端可能分别对同一键写入不同值。分区恢复后,必须选择其中一个值,另一个值就丢失了。Riak 的 CRDT 数据类型通过数学上可交换、可结合、可幂等的合并操作来解决这个问题,例如计数器使用加法和减法,集合使用并集和差集,这些操作无论以什么顺序执行,最终结果都一致。
Node.js 中操作 CRDT 和操作普通键值还有一个重要差异:CRDT 的更新不是通过覆盖整个值来完成的。你不能先读取一个集合,在本地添加元素后再把整个集合写回去,因为这样做会退化成普通键值覆盖,丢失并发更新。正确做法是发送增量更新指令,比如向集合添加某个元素、从集合移除某个元素、让计数器增加 3,由 Riak 在服务端根据类型规则合并。更新请求中还可以携带读取时获得的 context,让 Riak 判断是否发生过并发修改,从而实现条件更新。
const http = require('http');
function riakRequest(method, path, body, callback) {
const data = body ? JSON.stringify(body) : null;
const options = {
hostname: '127.0.0.1',
port: 8098,
path: path,
method: method,
headers: {
'Content-Type': 'application/json',
'Content-Length': data ? Buffer.byteLength(data) : 0
}
};
const req = http.request(options, (res) => {
let chunks = '';
res.on('data', (chunk) => { chunks += chunk; });
res.on('end', () => {
callback(null, res.statusCode, chunks ? JSON.parse(chunks) : null);
});
});
req.on('error', callback);
if (data) { req.write(data); }
req.end();
}
上面的封装函数足以完成本文所有示例。它使用 Node.js 内置 http 模块向 Riak 的 HTTP API 发送请求,路径和请求体由调用方决定。Riak 默认监听 8098 端口,数据以 JSON 格式传输。接下来会看到不同数据类型在路径、请求体和响应结构上的差异。
二、计数器与标志位的操作
计数器是最容易理解的 CRDT 类型。在 Riak 中,计数器只能增加或减少,不能直接设置为某个绝对值。即使两个客户端同时增加同一个计数器,Riak 也会把增量合并,最终结果等于正确总数。获取计数器时,路径格式为 /types/counters/buckets/统计桶名/datatypes/键名,返回的 JSON 中包含 value 字段,表示当前计数值。
更新计数器需要向同一路径发送 POST 请求,请求体中指定 increment 或 decrement。例如将页面访问量增加 5,可以这样写:
riakRequest('POST', '/types/counters/buckets/visits/datatypes/page1?returnbody=true', { increment: 5 }, (err, status, body) => {
if (err) { console.error(err); return; }
console.log('当前计数:', body.value);
});
标志位(flag)只有启用和禁用两种状态,底层用布尔值表示。与计数器只能增量不同,flag 的更新指令是 enable 或 disable,分别将标志设为 true 或 false。读取 flag 类型的路径以 /types/flags 开头,更新时请求体中使用 enable 或 disable 字段。需要注意的是,flag 并不支持取反操作,必须在客户端明确指定要设置成哪个状态,否则在并发场景下最后一次无上下文的赋值会覆盖对方的修改。
riakRequest('POST', '/types/flags/buckets/features/datatypes/new_ui?returnbody=true', { enable: true }, (err, status, body) => {
if (err) { console.error(err); return; }
console.log('标志位状态:', body.value);
});
三、集合与寄存器的更新语义
集合(set)可以存储多个不重复的元素,其合并规则是并集和差集。当两个副本都添加了不同元素时,合并后得到全部元素;当某个副本移除了元素时,合并后该元素被删除。移除操作优先于添加操作,这是为了保证删除意图不会被并发添加覆盖。在 Node.js 中,添加元素使用 add_all 字段,移除元素使用 remove_all 字段。两个字段可以同时出现在同一个请求体中。
集合非常适合用来维护标签、关注列表、设备在线状态等场景。更新时不需要先读取整个集合,只需告诉 Riak 要添加或移除哪些元素。例如给文章添加标签可以像这样:
riakRequest('POST', '/types/sets/buckets/articles/datatypes/post1?returnbody=true', { add_all: ['nodejs', 'crdt'], remove_all: ['oldtag'] }, (err, status, body) => {
if (err) { console.error(err); return; }
console.log('集合内容:', body.value);
});
寄存器(register)用于存储一个不可拆分的值,比如用户状态、配置项、当前价格等。它不支持像计数器那样的增量合并,因为值本身没有可交换的合并规则。寄存器采用最后写入优先的策略,但配合 context 可以避免无意识的覆盖。读取寄存器时,响应中会带有一个 context 字段,更新时将该 context 原样提交,Riak 会检查该 context 是否还是最新版本。如果期间发生了其他更新,当前提交就会被拒绝。这比无条件的最后写入覆盖要安全得多。
riakRequest('GET', '/types/registers/buckets/config/datatypes/status', null, (err, status, body) => {
if (err) { console.error(err); return; }
const ctx = body.context;
riakRequest('POST', '/types/registers/buckets/config/datatypes/status?returnbody=true', { assign: 'active', context: ctx }, (err2, status2, body2) => {
if (err2) { console.error(err2); return; }
console.log('更新后的值:', body2.value);
});
});
上面代码先读取寄存器的当前值和 context,再提交新的值。如果两个客户端同时读取到同一个 context,只有第一个提交会成功,第二个会收到冲突提示,然后可以重新读取并重试。这种方式牺牲了一点简单的直接覆盖,但避免了并发更新时后提交者把先提交者的修改覆盖掉。
四、嵌套 Map 数据类型的操作
Map 类型允许在一个键下组织多个不同类型的字段,每个字段可以是一个计数器、集合、寄存器、标志位,或者另一个 Map。这种层次结构适合表示复杂的实体,例如用户资料中既有关注数(计数器)、标签(集合),又有昵称(寄存器)和是否启用通知(标志位)。Map 的更新请求体中使用 update 字段,再在内部按类型名称包裹各个子字段的更新指令。
更新 Map 时,不需要把整个 Map 的内容重新提交,只需要指定要更新的嵌套字段及其增量操作。例如给用户增加 2 个关注者、添加一个标签,同时修改昵称:
riakRequest('POST', '/types/maps/buckets/users/datatypes/alice?returnbody=true', {
update: {
counters: { followers: { increment: 2 } },
sets: { tags: { add_all: ['developer'] } },
registers: { nickname: { assign: 'Alice_Dev' } }
}
}, (err, status, body) => {
if (err) { console.error(err); return; }
console.log('Map 内容:', JSON.stringify(body.value));
});
Map 中的每个嵌套类型都有独立的合并规则:计数器继续增量合并,集合继续并集处理,寄存器则使用最后写入优先。因此一个 Map 可以同时包含适合增量更新的字段和适合覆盖更新的字段,Riak 会分别处理,不会互相干扰。读取 Map 时也会返回 context,如果需要对寄存器等容易冲突的字段做条件更新,应该带上 context。
五、context 与并发安全的最佳实践
context 是 Riak CRDT 操作中不可忽略的概念。它实际上是一个不透明字符串,Riak 用它追踪数据版本。对于计数器、集合这类只依赖增量合并的类型,不携带 context 一般不会造成数据丢失,因为合并规则本身能保证并发更新被正确累积。但对于寄存器和 Map 中寄存器的赋值操作,不携带 context 就等同于无条件覆盖,可能会覆盖掉其他客户端的并发修改。
一个常见错误是开发者为了图省事,每次都直接发送不带 context 的更新。这样在低并发场景下可能运行正常,一旦两个请求同时到达,后处理的请求就会把前一个请求的修改清掉。正确做法是:如果业务上读取和写入之间有明确的先后依赖,就带上 context 做条件更新;如果拿不到可靠的 context,就评估业务是否真的需要寄存器,还是可以改用集合或计数器来表示状态。很多本可以用集合解决的问题,因为最初选用了寄存器,导致并发冲突处理变得复杂。
此外,所有更新请求都可以加上 returnbody=true 参数,让 Riak 在响应中直接返回更新后的完整数据。这在调试和快速验证时很方便,但在生产环境中如果数据体较大,建议根据实际需要决定是否拉取完整值,避免不必要的网络开销。Node.js 中处理 Riak HTTP 响应时要注意 JSON 解析失败的情况,可以封装一个安全的解析函数。总体而言,理解类型语义、使用增量更新、合理携带 context,是在 Node.js 中正确使用 Riak 数据类型的三个核心原则。