在混合存储架构里,SQLite常被用作端侧或单机持久化组件,Riak KV则承担分布式副本和高可用职责。要让这两类数据库的数据达成最终一致,核心不是网络层的一次性同步,而是设计一套可重放、可冲突消解的变更流。变更流必须记录每一次本地写入的元数据,保证推送过程幂等,并能在并发修改发生时给出明确裁决。

下面的实现从变更日志开始,逐步扩展到推送、拉取和冲突处理。所有代码以Go语言为例,SQLite部分通过标准库database/sql配合纯Go驱动,Riak KV使用官方riak-go-client。
变更日志与版本号:让SQLite每次写入留下可追踪的记录
最终一致性的第一步,是让本地数据库的每一次插入、更新、删除都变成一条可读取的增量记录。直接监听SQLite的写操作不现实,但可以借助触发器把变更自动写入一张独立的变更日志表。这张表只追加、不修改,配合单调递增的版本号,形成类似数据库binlog的结构。
版本号的设计要足够简单:在单机场景下,一个全局自增整数就能保证顺序。如果将来要扩展到多端,可以改为逻辑时间戳,例如毫秒级时间加上节点编号。本文先用自增主键作为版本号,并保留逻辑时间字段,方便后续演进。
-- 业务表:用户资料
CREATE TABLE user_profile (
user_id TEXT PRIMARY KEY,
nickname TEXT NOT NULL,
updated_at INTEGER NOT NULL
);
-- 变更日志表:只追加,记录每次写操作
CREATE TABLE change_log (
log_id INTEGER PRIMARY KEY AUTOINCREMENT,
table_name TEXT NOT NULL,
row_key TEXT NOT NULL,
payload TEXT NOT NULL,
version INTEGER NOT NULL,
is_deleted INTEGER NOT NULL DEFAULT 0,
synced INTEGER NOT NULL DEFAULT 0,
created_at INTEGER NOT NULL
);
-- 写入触发器:插入或更新时自动记录变更
CREATE TRIGGER trg_user_profile_upsert
AFTER INSERT ON user_profile
BEGIN
INSERT INTO change_log (table_name, row_key, payload, version, created_at)
VALUES ('user_profile', NEW.user_id,
json_object('user_id', NEW.user_id, 'nickname', NEW.nickname, 'updated_at', NEW.updated_at),
(SELECT COALESCE(MAX(version), 0) + 1 FROM change_log),
strftime('%s','now') * 1000);
END;
-- 更新触发器
CREATE TRIGGER trg_user_profile_update
AFTER UPDATE ON user_profile
BEGIN
INSERT INTO change_log (table_name, row_key, payload, version, created_at)
VALUES ('user_profile', NEW.user_id,
json_object('user_id', NEW.user_id, 'nickname', NEW.nickname, 'updated_at', NEW.updated_at),
(SELECT COALESCE(MAX(version), 0) + 1 FROM change_log),
strftime('%s','now') * 1000);
END;
-- 删除触发器:写入墓碑标记
CREATE TRIGGER trg_user_profile_delete
AFTER DELETE ON user_profile
BEGIN
INSERT INTO change_log (table_name, row_key, payload, version, is_deleted, created_at)
VALUES ('user_profile', OLD.user_id,
json_object('user_id', OLD.user_id),
(SELECT COALESCE(MAX(version), 0) + 1 FROM change_log),
1,
strftime('%s','now') * 1000);
END;
日志表的payload字段保存JSON格式的完整行快照,这样推送端不需要再回查业务表,只需要读取日志即可生成远端写入。is_deleted字段用于区分物理删除和逻辑删除,避免远程同步时把已经被删除的数据重新写回。
触发器里用子查询取当前最大版本号加一,虽然简单,但在高并发写入下可能产生重复版本号。如果写入频率很高,可以在应用层维护版本计数器,或者改用UUID加时间戳的组合。对于大多数SQLite单写者场景,这个方案足够可靠。
推送与拉取:实现双向同步循环
有了变更日志,下一步是把未同步的记录推送到Riak KV。Riak是一个基于键值对的分布式数据库,写入时需要指定bucket和key,值可以是任意二进制数据。Go代码从SQLite读取synced等于0的日志行,按版本号升序处理,每条记录对应一次Riak的put操作。处理成功后把日志行标记为已同步。
推送过程中必须考虑失败重试。如果一条日志推送成功但标记失败,下次同步会再次推送同一条记录,可能造成远端数据被旧值覆盖。因此,需要在Riak写入时带上版本号作为条件,或者依赖Riak的向量时钟做冲突检测。简单做法是使用版本号作为Riak对象的一个元数据字段,并在应用端检查远端当前版本,避免旧版本覆盖新版本。
package main
import (
"database/sql"
"encoding/json"
"fmt"
"log"
"time"
riak "github.com/basho/riak-go-client"
_ "github.com/mattn/go-sqlite3"
)
type ChangeRow struct {
LogID int64
TableName string
RowKey string
Payload string
Version int64
IsDeleted int
CreatedAt int64
}
func pushUnsynced(db *sql.DB, bucket string) error {
rows, err := db.Query(`SELECT log_id, table_name, row_key, payload, version, is_deleted, created_at
FROM change_log WHERE synced = 0 ORDER BY version ASC`)
if err != nil {
return err
}
defer rows.Close()
client, err := riak.NewClient(&riak.NewClientOptions{
RemoteAddresses: []string{"127.0.0.1:8087"},
})
if err != nil {
return err
}
defer client.Stop()
for rows.Next() {
var cr ChangeRow
if err := rows.Scan(&cr.LogID, &cr.TableName, &cr.RowKey, &cr.Payload, &cr.Version, &cr.IsDeleted, &cr.CreatedAt); err != nil {
return err
}
obj := &riak.Object{
Bucket: bucket,
Key: fmt.Sprintf("%s:%s", cr.TableName, cr.RowKey),
ContentType: "application/json",
Value: []byte(cr.Payload),
}
// 将本地版本号写入用户元数据,供后续冲突判断
obj.UserMeta = []*riak.Pair{
{Key: "local_version", Value: []byte(fmt.Sprintf("%d", cr.Version))},
}
cmd, err := riak.NewStoreValueCommandBuilder().
WithBucket(bucket).
WithObject(obj).
Build()
if err != nil {
return err
}
if err := client.Execute(cmd); err != nil {
return err
}
// 标记为已同步
_, err = db.Exec(`UPDATE change_log SET synced = 1 WHERE log_id = ?`, cr.LogID)
if err != nil {
return err
}
}
return rows.Err()
}
上面代码中的远端地址使用了127.0.0.1,实际部署时应替换为Riak集群的节点地址。同步循环通常放在一个goroutine里,定时执行,例如每5秒或每30秒跑一次。为了减少对SQLite的锁竞争,可以用一个专门的连接池处理同步任务。
拉取远程变更则相反:从Riak读取所有键,或使用Riak的2i索引查询最近修改的对象,然后把远端值应用到本地SQLite。应用本地时需要先检查本地变更日志是否已存在更高版本,如果远端版本落后于本地未推送变更,则丢弃远端值,避免本地新数据被覆盖。这种双向同步需要维护一个统一的版本比较规则。
冲突检测与解决:Riak向量时钟与删除墓碑
Riak KV内置了向量时钟机制,每次读取对象都会返回一个不透明的vclock字符串,写入时如果携带这个vclock,Riak会判断该对象是否被其他副本并发修改过。如果并发修改导致冲突,Riak默认保存多个兄弟值,由客户端决定合并或选择其一。利用这个特性,可以把冲突检测下推到Riak层,而不需要在应用层完全重造轮子。
实际同步时,推送阶段可以不携带vclock,使用最后写入胜出策略;但对于拉取后的本地更新再推送回Riak,必须携带读取时的vclock,否则可能覆盖掉其他客户端刚写入的新数据。下面代码演示从Riak读取对象、修改后带vclock写回,并处理返回的冲突情况。
func fetchAndUpdateRiak(client *riak.Client, bucket, key string) error {
// 读取当前对象及其vclock
fetchCmd, _ := riak.NewFetchValueCommandBuilder().
WithBucket(bucket).
WithKey(key).
Build()
if err := client.Execute(fetchCmd); err != nil {
return err
}
fetchRes := fetchCmd.(*riak.FetchValueCommand).Response
if fetchRes.IsNotFound {
return fmt.Errorf("key %s not found", key)
}
if len(fetchRes.Values) == 0 {
return fmt.Errorf("no value returned")
}
// 假设只有一个值
obj := fetchRes.Values[0]
// 修改值,例如把JSON里的nickname改成大写
var data map[string]interface{}
json.Unmarshal(obj.Value, &data)
data["nickname"] = "UPDATED_NAME"
newValue, _ := json.Marshal(data)
// 构建新对象,必须带上原vclock
newObj := &riak.Object{
Bucket: bucket,
Key: key,
ContentType: "application/json",
Value: newValue,
VClock: obj.VClock, // 关键:携带读到的向量时钟
}
storeCmd, _ := riak.NewStoreValueCommandBuilder().
WithBucket(bucket).
WithObject(newObj).
WithReturnBody(true).
Build()
if err := client.Execute(storeCmd); err != nil {
return err
}
storeRes := storeCmd.(*riak.StoreValueCommand).Response
if len(storeRes.Values) > 1 {
// 出现了兄弟值,说明发生了并发冲突
log.Printf("conflict detected, siblings: %d", len(storeRes.Values))
// 这里可以执行自定义合并逻辑,比如取最新时间戳对应的值
}
return nil
}
对于删除操作,必须使用墓碑标记而不是物理删除。如果推送端直接把Riak对象删除,万一有另一个副本在稍后重新推送旧数据,该键又会复活。正确做法是写入一个特殊值,例如{"deleted":true},并保留这个墓碑一段时间。Riak支持通过设置X-Riak-Deleted头来创建墓碑,也可以直接写一个业务层墓碑对象。所有客户端在读取时发现墓碑值就当作不存在处理。
最终一致性的工程难度不在于单次同步的成功,而在于网络抖动、节点故障、并发写入等异常情况下系统依然能收敛。上述变更日志、幂等推送、向量时钟冲突检测和墓碑删除,构成了一个最小可用的SQLite与Riak KV最终一致性方案。实际项目中还要补充监控和重试队列,但核心思路已经完整。