Skip to content

[To DLDB] 为dldb增加按bucket批量partial update能力 #24

Description

@hsballoon

背景

当前 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 设计,但需要满足以下语义:

  1. 每条记录可以有不同的 update values;
  2. 只修改指定字段,未指定字段保持原值;
  3. 默认只更新已存在的记录,不自动插入缺失 ID;
  4. 同一 bucket 内的一批 updates 尽量只产生一次 Lance commit;
  5. 输入跨多个 HASH bucket 时,按 bucket 分组,每个 bucket 独立提交;
  6. 同一批次出现重复 ID 或冲突 patch 时应明确报错;
  7. 支持返回 updated/not-found 的结果;
  8. 对 concurrent writer conflict 使用有限重试和退避,超过预算后快速返回明确的 conflict 错误;
  9. 不应要求调用方先读取并提交包含完整宽列的整行;
  10. 应保留现有分区字段、主键字段等不可变字段的保护语义。

实现方式

希望 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 的执行结果和失败信息。

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions