智能体数据怎么增量同步?用 SQLite 核对版本、重跑与回滚

以 AI 样本复核元数据为例,用固定 SQLite 写入工具处理增量版本,验证重跑不增行、迟到旧版不覆盖新版,以及失败批次回滚。

本例同步 AI 样本的复核元数据:同一个 sample_id 可能从 pending 更新到 verified,重跑任务与迟到的旧批次都不能把新状态冲掉。同步由智能体外围的固定数据工具执行,模型不生成 SQL,也不决定某条训练样本是否已经通过复核。

样本编号、状态和版本均为人为构造,没有接入真实智能体或由模型审核样本。程序只应用可信上游已经确认的元数据。输入的 version 必须是该样本由上游定义的递增版本,不是模型猜测的数字或服务器当前时间。

智能体数据怎么增量同步?用 SQLite 核对版本、重跑与回滚

先定义这四种输入如何处理

上游记录相对当前记录 工具动作 理由
ID 尚不存在 插入 保持稳定样本身份
版本更高 更新状态和版本 应用后续已确认复核结果
同版本、同内容 保持原状 允许同批次安全重跑
版本更低 保持原状 迟到旧消息不能覆盖新版
同版本、内容不同 拒绝并回滚整批 版本契约冲突,需要回源确认

数据库的键是 dataset_id 与 sample_id 的组合;DATASET 由可信应用固定为 ai-sample-demo。真实系统若同时处理多个数据集,应由身份与任务上下文决定范围,不能让模型任意传入租户、表名或同步范围。

在 SQLite 里执行固定写入

示例已在 Windows、Python 3.11.15、SQLite 运行库 3.53.1 执行,使用 Python 自带 sqlite3,无需安装额外包。UPSERT 需要 SQLite 3.24.0 或之后版本,可先运行 python -c “import sqlite3; print(sqlite3.sqlite_version)” 检查实际运行库。

把代码保存为 incremental_sync.py,执行 python incremental_sync.py。它创建一次性内存演示库,不接触真实训练数据。

import json
import sqlite3

# dataset_id is fixed by the trusted application, not chosen by an LLM.
DATASET = 'ai-sample-demo'
con = sqlite3.connect(':memory:', isolation_level=None)
con.execute('''CREATE TABLE sample_metadata (
    dataset_id TEXT NOT NULL, sample_id TEXT NOT NULL,
    version INTEGER NOT NULL CHECK(version > 0), payload TEXT NOT NULL,
    PRIMARY KEY(dataset_id, sample_id)
)''')

def snapshot():
    return con.execute('''SELECT sample_id, version, payload FROM sample_metadata
        WHERE dataset_id=? ORDER BY sample_id''', (DATASET,)).fetchall()

def sync_batch(rows, fail_after=None):
    prepared = []
    ids = set()
    for sample_id, version, status in rows:
        if not isinstance(sample_id, str) or not sample_id.strip() or sample_id in ids:
            raise ValueError('invalid or duplicate sample_id')
        if type(version) is not int or version <= 0:
            raise ValueError('version must be a positive integer')
        if status not in {'pending', 'verified', 'rejected'}:
            raise ValueError('unknown review_status')
        ids.add(sample_id)
        payload = json.dumps({'review_status': status}, sort_keys=True, separators=(',', ':'))
        prepared.append((sample_id, version, payload))
    changed = 0
    con.execute('BEGIN IMMEDIATE')
    try:
        for position, (sample_id, version, payload) in enumerate(prepared, 1):
            old = con.execute('''SELECT version,payload FROM sample_metadata
                WHERE dataset_id=? AND sample_id=?''', (DATASET, sample_id)).fetchone()
            if old and old[0] == version and old[1] != payload:
                raise ValueError('same version has different content')
            cursor = con.execute('''INSERT INTO sample_metadata VALUES(?,?,?,?)
                ON CONFLICT(dataset_id,sample_id) DO UPDATE
                SET version=excluded.version,payload=excluded.payload
                WHERE excluded.version > sample_metadata.version''',
                (DATASET, sample_id, version, payload))
            changed += cursor.rowcount
            if position == fail_after:
                raise RuntimeError('injected failure for rollback verification')
        con.execute('COMMIT')
    except Exception:
        con.execute('ROLLBACK')
        raise
    return changed

initial = [('s01', 1, 'pending'), ('s02', 1, 'pending')]
print('initial_changed:', sync_batch(initial))
before_replay = snapshot()
print('replay_changed:', sync_batch(initial))
assert snapshot() == before_replay
print('new_version_changed:', sync_batch([('s01', 2, 'verified')]))
print('late_old_version_changed:', sync_batch([('s01', 1, 'pending')]))
assert snapshot()[0][1:] == (2, '{"review_status":"verified"}')
before_failure = snapshot()
try:
    sync_batch([('s02', 2, 'verified'), ('s03', 1, 'pending')], fail_after=1)
except RuntimeError:
    print('failed_batch_rolled_back:', snapshot() == before_failure)
else:
    raise AssertionError('Failure was not injected')
assert snapshot() == before_failure
try:
    sync_batch([('s02', 2, 'verified'), ('s01', 2, 'rejected')])
except ValueError:
    print('same_version_conflict_rejected:', True)
else:
    raise AssertionError('Version conflict was accepted')
assert snapshot() == before_failure
print('final_rows:', snapshot())
con.close()

为什么失败不会留下半批状态

连接设置 isolation_level=None,示例显式执行 BEGIN IMMEDIATE、COMMIT 和 ROLLBACK。每条记录先在同一事务里读当前版本,然后使用固定参数化 UPSERT;WHERE 条件只允许更高版本覆盖。中间任一异常进入回滚,避免第一条更新成功、后续失败后仍被误报整批完成。

故意在第一条写入后抛错,整批回滚;同版本异内容冲突也会回滚该批已执行的其他写入。第二个验证先尝试把 s02 升到版本 2,再发现 s01 的版本 2 内容冲突,最终 s02 也必须恢复版本 1。

对照实际回读结果

initial_changed: 2
replay_changed: 0
new_version_changed: 1
late_old_version_changed: 0
failed_batch_rolled_back: True
same_version_conflict_rejected: True
final_rows: [('s01', 2, '{"review_status":"verified"}'), ('s02', 1, '{"review_status":"pending"}')]

初次写入改变 2 行,原批次重跑改变 0 行;s01 的版本 2 写入后,迟到版本 1 改变 0 行,仍保留 verified。最终数据库只有 s01 和 s02:s01 为版本 2/verified,s02 为版本 1/pending,没有失败批次中的 s03。

这些判断通过完整排序快照和逐字段断言核对,而不是只看一条“同步成功”日志。changed=0 可能是重复消息,也可能是旧版本;需要审计时应另记每条原因,不把它一律称为新数据已更新。

接入真实上游时补齐边界

上游负责给出增量记录、稳定编号及可信版本;本工具只负责应用。不要将同步状态与模型训练是否完成混为一谈:元数据写入成功后,还需按样本编号核对训练输入和标签版本。增量批次里没有某个 ID,不表示这个样本已删除;删除需要单独的可信上游事件与规则。

内存演示库结束进程后消失。若要保留状态,应连接经过批准的新测试库文件,并将初始化演示数据与同步函数分开,再验证重启后的同批重跑。不要直接把本例的 CREATE TABLE 和人工数据块用于现有业务库。

BEGIN IMMEDIATE 可能因另一写事务返回 SQLITE_BUSY;这时应让任务进入可重试队列并回读状态,不能把锁失败报告为成功。本例串行执行,没有做跨进程并发测试。

本例保证的是一个 SQLite 库内、一个事务里的写入一致性,不能据此声称外部 API、消息队列和多个数据库实现 exactly-once。如果同步后还要调用外部服务,应另设计提交后消息与恢复核对,不能在事务中随意发送不可回滚动作。

官方资料

以下资料实际读取于 2026 年 10 月 1 日;接口或界面范围按正文注明的版本核对。

Ai菜鸟网。发布者:AI小管家,转载请注明出处:https://www.alyyhw.com/32730.html

赞 (0)
AI小管家的头像AI小管家
工业 AI 数据融合怎么对齐传感器时间?只匹配已发生的读数
上一篇 1小时前
AI 表格分析工具怎么选?用同一份 CSV 比较 pandas、Polars 与 DuckDB
下一篇 1小时前

相关推荐

联系我们

联系我们

1

在线咨询: QQ交谈

邮件:admin@example.com

工作时间:周一至周五,9:30-18:30,节假日休息

关注微信
关注微信
分享本页
返回顶部