背景
当前 landing 表使用 HASH(job_id),共 128 个 bucket。SAfactory 的 reward 提交通常只需要修改少量字段,例如:
reward
step_reward
is_trainable
is_session_completed
source_updated_at
但每条记录的 reward 可能不同,当前 dldb 的 Session.update() 只能对所有匹配行应用同一组 values:
session.update(
table_name,
where="id IN (...)",
values={"reward": 0.8},
partition=partition,
)
因此对于不同记录的不同 reward,只能逐条调用 update,产生多次 Lance commit,容易造成同一 bucket 上的 writer contention、重试和 queue 堆积。
当前 Session.upsert() 底层使用 merge_insert(...).when_matched_update_all().when_not_matched_insert_all()。该方式要求提交完整行,属于整行替换/插入,不适合作为 reward partial update:
- landing 行包含较大的 JSON/payload 字段;
- 为更新一个 reward 需要读取和重新提交完整宽行;
- 可能覆盖其他并发更新;
- 不存在的 ID 会被插入,语义不符合 reward patch。
期望能力
希望增加一个批量 partial update 接口,例如:
session.batch_update(
table_name="landing_test",
updates=[
{
"id": "record-1",
"values": {
"reward": 0.8,
"step_reward": 0.8,
},
},
{
"id": "record-2",
"values": {
"reward": 0.6,
"step_reward": 0.6,
},
},
],
partition=34,
)
具体 API 名称可以由 dldb team 设计,但需要满足以下语义:
- 每条记录可以有不同的 update values;
- 只修改指定字段,未指定字段保持原值;
- 默认只更新已存在的记录,不自动插入缺失 ID;
- 同一 bucket 内的一批 updates 尽量只产生一次 Lance commit;
- 输入跨多个 HASH bucket 时,按 bucket 分组,每个 bucket 独立提交;
- 同一批次出现重复 ID 或冲突 patch 时应明确报错;
- 支持返回 updated/not-found 的结果;
- 对 concurrent writer conflict 使用有限重试和退避,超过预算后快速返回明确的 conflict 错误;
- 不应要求调用方先读取并提交包含完整宽列的整行;
- 应保留现有分区字段、主键字段等不可变字段的保护语义。
实现方式
希望 dldb 评估基于 native Lance 的 partial update 能力实现,而不是简单循环调用现有 update()。
Native Lance 当前支持 SQL expression update,例如 values_sql。对于 reward 等标量字段,可以考虑由 dldb 安全地生成按 ID 匹配的 CASE WHEN 表达式,或者使用其他底层批量 update 原语。
不建议仅通过完整行 merge_insert/upsert 模拟该接口,因为这会带来宽列读写放大和 stale row 覆盖风险。
并发场景
该接口主要用于以下场景:
Gateway 持续 append telemetry
RewardCommitter 批量更新同一 HASH bucket 中多条记录的 reward
目标是在降低 commit 次数的同时,减少同一物理 bucket 上的 writer conflict 和重试等待。
需要注意:跨 bucket 的批量操作不要求全局原子性;应明确返回每个 bucket 的执行结果和失败信息。
背景
当前 landing 表使用
HASH(job_id),共 128 个 bucket。SAfactory 的 reward 提交通常只需要修改少量字段,例如:但每条记录的 reward 可能不同,当前 dldb 的
Session.update()只能对所有匹配行应用同一组 values:因此对于不同记录的不同 reward,只能逐条调用 update,产生多次 Lance commit,容易造成同一 bucket 上的 writer contention、重试和 queue 堆积。
当前
Session.upsert()底层使用merge_insert(...).when_matched_update_all().when_not_matched_insert_all()。该方式要求提交完整行,属于整行替换/插入,不适合作为 reward partial update:期望能力
希望增加一个批量 partial update 接口,例如:
具体 API 名称可以由 dldb team 设计,但需要满足以下语义:
实现方式
希望 dldb 评估基于 native Lance 的 partial update 能力实现,而不是简单循环调用现有
update()。Native Lance 当前支持 SQL expression update,例如
values_sql。对于 reward 等标量字段,可以考虑由 dldb 安全地生成按 ID 匹配的CASE WHEN表达式,或者使用其他底层批量 update 原语。不建议仅通过完整行
merge_insert/upsert模拟该接口,因为这会带来宽列读写放大和 stale row 覆盖风险。并发场景
该接口主要用于以下场景:
目标是在降低 commit 次数的同时,减少同一物理 bucket 上的 writer conflict 和重试等待。
需要注意:跨 bucket 的批量操作不要求全局原子性;应明确返回每个 bucket 的执行结果和失败信息。