数据库基础体系 · 第 74/139 篇。文章以各产品官方稳定版本的公开语义为准;示例会明确引擎、事务与部署边界。

SQL 批处理与分页:游标、Keyset、批量写、限速和断点恢复

批处理通常不是“每次取几千行,然后循环执行 SQL”这么简单。只要数据量达到一定规模,就必须同时回答以下问题:

  • 如何稳定地找到下一批数据?
  • 如何避免 OFFSET 越翻越慢,或因为并发修改而漏读、重复读?
  • 一批数据应该读取多少、写入多少、提交多少?
  • 读取和写入是否处于同一个事务、同一个数据库,甚至同一个系统?
  • 如何限制数据库压力,而不是只限制应用线程速度?
  • 进程在提交前或提交后崩溃时,下一次从哪里继续?
  • 重试一批数据时,如何避免重复写入或产生错误副作用?

本文围绕这些问题,分别讨论游标、Keyset 分页、批量写、限速和断点恢复。示例主要使用 PostgreSQL;涉及 MySQL 时会明确说明差异。


一、先明确“分页”与“批处理”不是同一个概念

1. 分页是定位数据的方式

分页解决的是:

当前已经处理了前一批数据,如何找到下一批数据?

常见实现有:

  • OFFSET ... LIMIT ...
  • 游标
  • Keyset 分页,也叫 seek pagination
  • 按时间、ID 或其他范围分段

分页本身不决定事务如何提交,也不决定数据如何写入。

2. 批处理是执行和提交的方式

批处理解决的是:

一次取多少行、一次写多少行、什么时候提交、失败后如何重试?

例如,一次处理 1000 行可能包括:

  1. 查询 1000 行;
  2. 在应用中转换;
  3. 批量写入目标表;
  4. 更新断点;
  5. 提交事务。

因此,“分页方式”和“批处理策略”可以组合:

  • OFFSET + 每批提交;
  • Keyset + 每批提交;
  • 游标 + 流式读取;
  • 固定主键范围 + 批量写;
  • 游标读取 + 目标库批量写入。

一个可靠的方案必须同时设计这两个层面。


二、为什么简单的 OFFSET 分页会失效

最直观的分页写法是:

SELECT id, payload
FROM source_table
ORDER BY id
LIMIT 1000 OFFSET 0;

下一页把 OFFSET 改成 1000:

SELECT id, payload
FROM source_table
ORDER BY id
LIMIT 1000 OFFSET 1000;

1. OFFSET 的逻辑含义

数据库必须先按照 ORDER BY id 确定有序结果,然后跳过前 OFFSET 行,再返回后面的 LIMIT 行。

当偏移量为 k 时,数据库通常至少需要处理前面的 k 行。即使优化器能够利用索引,随着 k 增大,扫描和丢弃的数据量通常也会增大。

这带来两个问题:

  1. 性能问题:越靠后的页,定位成本可能越高;
  2. 并发一致性问题:页与页之间如果有插入或删除,偏移量代表的“位置”会改变。

2. 插入导致重复或漏读

假设第一次查询得到:

id
1
2
3

此时使用:

LIMIT 3 OFFSET 0

在第二次查询前插入 id = 0,第二次执行:

LIMIT 3 OFFSET 3

新的有序结果为:

位置 id
1 0
2 1
3 2
4 3
5 4
6 5

跳过前三行后,第二页从 id = 3 开始。这里看似没有重复,但如果第一页和第二页的边界、排序条件或写入过程不同,插入和删除可能导致某些行被跳过或重复。

例如第一页原本为 1, 2, 3,随后删除 id = 1,第二页仍然使用 OFFSET 3,此时会跳过 2, 3, 4,从 5 开始,4 被漏掉。

3. 没有唯一排序的 OFFSET 更不可靠

下面的排序不是稳定排序:

SELECT *
FROM orders
ORDER BY created_at
LIMIT 1000 OFFSET 1000;

如果多个订单拥有相同的 created_at,数据库没有义务在这些相同值的行之间维持固定顺序。即使数据完全不变,不同执行也可能在边界处出现不同结果。

至少应补充唯一键:

ORDER BY created_at, id

但这只能使排序边界可定义,不能解决 OFFSET 在大偏移量和跨请求并发变化下的根本问题。


三、Keyset 分页:用“最后一行的键”定位下一批

1. 基本思想

Keyset 分页不记录“已经跳过了多少行”,而记录上一批最后一行的排序键。

第一批:

SELECT id, payload
FROM source_table
WHERE ...
ORDER BY id
LIMIT 3;

假设返回:

id
10
20
30

下一批不是使用 OFFSET 3,而是:

SELECT id, payload
FROM source_table
WHERE id > 30
ORDER BY id
LIMIT 3;

这里的 30 是断点,也就是上一批最后一行的键。

如果 id 是递增且不可变的唯一键,那么新插入的 id = 35 会出现在后续扫描中,而已经处理过的 id <= 30 不会因为前面插入新数据而被重新计算位置。

2. Keyset 的形式化条件

设排序键为 KK,批大小为 BB,数据按照严格递增顺序排列:

K1<K2<<KnK_1 < K_2 < \cdots < K_n

第一批返回:

K1,,KBK_1, \ldots, K_B

断点记录为:

C=KBC = K_B

下一批查询:

K>CK > C

并取前 BB 行。

要使这个过程可靠,至少需要以下条件:

  1. 排序关系是全序:任意两行都能确定先后;
  2. 排序键唯一,或通过唯一键补足唯一性
  3. 用于继续扫描的键不会向后移动
  4. 查询条件在各批之间有明确语义
  5. 断点只在对应批次成功提交后推进

如果排序键不是唯一的,例如只有 created_at,条件:

WHERE created_at > :last_created_at
ORDER BY created_at

可能漏掉与最后一行拥有相同时间戳的其他行。

3. 复合 Keyset 条件

正确做法是使用 (created_at, id) 作为复合排序键:

SELECT id, created_at, payload
FROM source_table
WHERE (created_at, id) > (:last_created_at, :last_id)
ORDER BY created_at, id
LIMIT 1000;

这表示:

(created_at,id)>(t,i)(created\_at, id) > (t, i)

等价于:

WHERE created_at > :last_created_at
   OR (
        created_at = :last_created_at
        AND id > :last_id
      )

排序:

ORDER BY created_at, id

必须和比较条件保持完全一致。

PostgreSQL 示例

CREATE INDEX source_table_created_id_idx
ON source_table (created_at, id);

MySQL 示例

CREATE INDEX source_table_created_id_idx
ON source_table (created_at, id);

两者都需要保证索引顺序与查询的排序和范围条件相匹配。具体执行计划仍应通过 EXPLAINEXPLAIN ANALYZE 验证,而不能仅凭索引名称判断已经优化。

4. NULL 会破坏简单的 Keyset 逻辑

SQL 中的 NULL 不是普通值,比较结果可能是未知:

created_at > NULL

结果不是 TRUE,而是 UNKNOWN,在 WHERE 中不会被选中。

因此排序列最好定义为 NOT NULL。如果不能做到,必须明确规定空值位置,并让排序和分页条件使用同一规则。例如 PostgreSQL 可以写:

ORDER BY created_at NULLS LAST, id

但对应的分页条件不能简单地只写:

WHERE (created_at, id) > (...)

需要把 NULL 分支单独设计,否则会出现漏行。

5. Keyset 不能自动解决数据修改问题

假设按 id 递增处理,但某行的 id 不变,其他字段可以修改:

  • 已经处理的行被更新:不会重新出现;
  • 尚未处理的行被更新:只要 id 不变,仍会在原位置出现;
  • 行被删除:查询时自然不可见;
  • 新增的行且 id 大于当前断点:可能在后续批次出现;
  • 新增的行且 id 小于等于当前断点:不会出现。

因此 Keyset 的语义通常是:

扫描某个动态数据集在某个键方向上的前进过程,而不是对所有未来新增数据提供全量快照。

如果要求“任务开始时存在的所有行必须恰好处理一次”,需要额外定义边界,例如任务开始时取得最大 ID:

SELECT max(id) FROM source_table;

得到 high_watermark = 1000000 后,所有批次都使用:

WHERE id > :last_id
  AND id <= :high_watermark
ORDER BY id
LIMIT 1000;

这样,任务不会把任务开始后新增且 id > high_watermark 的行混入本次迁移。

但这仍然依赖一个重要前提:id 的生成和排序符合预期,且待处理行不会通过更新改变其筛选条件。


四、游标:让数据库维护结果集的位置

1. 游标的基本概念

游标是一种有状态的数据库对象。它通常包含:

  • 一个查询结果或结果生成过程;
  • 当前读取位置;
  • FETCHMOVE 等移动操作;
  • 与事务或连接相关的生命周期。

在 PostgreSQL 中,游标可以这样使用:

BEGIN;

DECLARE user_cursor CURSOR FOR
SELECT id, payload
FROM source_table
ORDER BY id;

FETCH FORWARD 1000 FROM user_cursor;
FETCH FORWARD 1000 FROM user_cursor;

CLOSE user_cursor;

COMMIT;

每次 FETCH 会从游标当前位置继续读取。

2. PostgreSQL 游标的事务边界

默认情况下,PostgreSQL 游标与事务绑定:

  • DECLARE 必须在事务中执行;
  • 事务提交或回滚后,普通游标失效;
  • 游标绑定在创建它的数据库连接上;
  • 连接池不能把 DECLARE 放在一个连接、把 FETCH 随意交给另一个连接。

这意味着应用必须固定连接:

获取连接
  -> BEGIN
  -> DECLARE
  -> 多次 FETCH
  -> CLOSE
  -> COMMIT/ROLLBACK
归还同一个连接

如果连接池在每次 SQL 执行时随机分配连接,游标操作会失败,或者后续请求根本找不到游标。

3. PostgreSQL 的 WITH HOLD

PostgreSQL 支持:

BEGIN;

DECLARE user_cursor CURSOR WITH HOLD FOR
SELECT id, payload
FROM source_table
ORDER BY id;

COMMIT;

FETCH FORWARD 1000 FROM user_cursor;

WITH HOLD 允许游标跨越创建它的事务继续使用。不过这不等于“游标自动获得无限期、低成本的快照”。数据库需要保存可供后续读取的结果,创建或首次使用时可能产生额外 I/O、内存或临时文件开销。

因此,WITH HOLD 适合需要在事务结束后继续消费某个结果集的场景,但不应把它误解成适合长期迁移的通用断点机制。

4. 游标与一致性

在 PostgreSQL 中,如果游标创建和持续读取处于一个长事务中,普通游标通常会沿用该事务的可见性规则。这样可以使读取结果更接近一个固定视图,但代价是:

  • 长事务会延长旧版本保留时间;
  • 可能阻碍 VACUUM 回收;
  • 查询和事务资源长期占用;
  • 事务中途失败时需要重新处理或重新建立游标。

如果每批单独提交,资源占用较小,但各批次看到的数据可能不同。此时应使用 Keyset 或范围断点,并接受“跨批次动态数据集”的明确语义。

5. MySQL 中的游标边界

MySQL 的 SQL 游标主要用于存储过程和存储函数等服务器端程序环境,不等同于所有客户端驱动都自动提供的“服务端流式结果集”。

客户端是否支持:

  • 服务端游标;
  • 流式读取;
  • 读取时是否缓冲全部结果;
  • 游标是否占用连接直到结果消费完成;

取决于具体驱动和 API。不能仅凭“数据库支持游标”就假设应用代码可以跨请求、跨连接使用同一个游标。

在 MySQL 应用中,若需要可恢复的长任务,通常更容易控制的方案是:

  • Keyset 查询;
  • 固定范围分片;
  • 应用层保存断点;
  • 每批独立事务。

五、游标和 Keyset 如何选择

两者解决的问题相近,但状态位置不同。

游标保存在哪里

游标状态保存在数据库连接和数据库会话中:

数据库连接 -> 游标 -> 当前读取位置

优点:

  • 读取语义自然;
  • 不需要应用自己拼接“最后一个键”的条件;
  • 适合连续消费一个结果集。

缺点:

  • 依赖同一连接;
  • 连接断开后通常需要重新开始或重新定位;
  • 长事务和长时间占用连接的成本明显;
  • 很难直接跨进程、跨机器恢复。

Keyset 保存在哪里

Keyset 的状态通常保存在应用或数据库表中:

checkpoint = (last_created_at, last_id)

下一次查询通过断点重新定位:

WHERE (created_at, id) > (:last_created_at, :last_id)

优点:

  • 易于持久化;
  • 可以跨连接、跨进程、跨机器恢复;
  • 每批可以独立提交;
  • 更容易和限速、重试、监控结合。

缺点:

  • 必须严格设计排序键;
  • 必须处理删除、更新、并发插入和边界条件;
  • 如果业务条件变化,断点语义也可能变化。

对于需要长时间运行、可暂停、可重试、可迁移的生产任务,Keyset 或固定范围通常比单一长生命周期游标更容易恢复。


六、批量读取、批量写入和事务提交

1. 为什么不能逐行写入

下面的逻辑很容易实现:

for each row:
    INSERT ...
    COMMIT

但每一行都可能产生:

  • 一次网络往返;
  • 一次解析或执行;
  • 一次日志刷盘或提交协调;
  • 一次锁获取和释放;
  • 一次连接占用。

即使数据库本身很快,网络和提交次数也会把吞吐量压低。

批处理的核心是把多行操作组合起来:

读取一批
  -> 转换一批
  -> 写入一批
  -> 提交一次

2. PostgreSQL 多行 INSERT

INSERT INTO target_table (id, payload)
VALUES
    (1, 'a'),
    (2, 'b'),
    (3, 'c');

这减少了多次单行 SQL 的通信开销。

如果数据来自另一个 PostgreSQL 查询,也可以使用 COPY。例如客户端通常可以执行:

COPY target_table (id, payload)
FROM STDIN WITH (FORMAT csv);

COPY 是 PostgreSQL 的高吞吐数据装载机制,但它不是“自动具备断点恢复”的迁移方案。应用仍需设计:

  • 当前批次如何确认成功;
  • 失败后整批重试还是从行内恢复;
  • 目标端重复数据如何处理;
  • 源端断点何时推进。

3. MySQL 多行 INSERT

INSERT INTO target_table (id, payload)
VALUES
    (1, 'a'),
    (2, 'b'),
    (3, 'c');

MySQL 也支持多行插入。批量大小受多个因素影响,例如:

  • max_allowed_packet
  • 行大小;
  • 事务日志;
  • 锁持有时间;
  • 复制延迟;
  • 客户端驱动的参数绑定方式。

不能把“能放进一个 SQL 字符串”当成合理批量大小。应使用参数化语句,并观察实际执行时间、锁等待和日志压力。

4. 批量 UPDATE 和 DELETE

批量写不只包括 INSERT。

PostgreSQL:按 Keyset 批量更新

WITH batch AS (
    SELECT id
    FROM source_table
    WHERE id > :last_id
      AND id <= :high_watermark
      AND processed = false
    ORDER BY id
    LIMIT 1000
)
UPDATE source_table AS s
SET processed = true,
    processed_at = clock_timestamp()
FROM batch
WHERE s.id = batch.id;

这里有一个并发边界:如果其他事务同时修改 processed 或删除行,实际更新行数可能少于 1000。应用不能只根据“查询时选中了多少行”推断“最终成功更新了多少行”,应检查实际影响行数和业务结果。

PostgreSQL:批量 DELETE

WITH batch AS (
    SELECT id
    FROM event_log
    WHERE id > :last_id
      AND id <= :high_watermark
    ORDER BY id
    LIMIT 1000
)
DELETE FROM event_log AS e
USING batch
WHERE e.id = batch.id;

删除旧数据时还要考虑外键、触发器、级联删除和锁等待。批量删除的每批大小不能只由查询性能决定,还取决于锁和事务日志压力。

5. 批量 Upsert 不是自动幂等

PostgreSQL:

INSERT INTO target_table (id, payload)
VALUES (:id, :payload)
ON CONFLICT (id) DO UPDATE
SET payload = EXCLUDED.payload;

MySQL 8.4:

INSERT INTO target_table (id, payload)
VALUES (?, ?)
ON DUPLICATE KEY UPDATE
    payload = VALUES(payload);

这类语句可以让重复执行同一行时不产生重复主键,但“幂等”仍有边界:

  • 是否所有写入列都被一致地覆盖?
  • 是否会更新时间戳、递增计数器等非幂等字段?
  • 是否触发副作用,例如触发器、消息发布、审计记录?
  • 冲突键是否真的是业务唯一键?
  • 两次批处理是否可能以不同版本的数据覆盖彼此?

因此,数据库层面的唯一约束和 Upsert 只解决了一部分重复写问题。


七、一个可恢复的 PostgreSQL Keyset 批处理示例

下面构造一个“同一个 PostgreSQL 数据库内,将源表转换后写入目标表”的示例。

1. 表结构

CREATE TABLE source_user (
    id bigint PRIMARY KEY,
    email text NOT NULL,
    updated_at timestamptz NOT NULL,
    status text NOT NULL
);

CREATE TABLE target_user (
    id bigint PRIMARY KEY,
    normalized_email text NOT NULL,
    source_updated_at timestamptz NOT NULL
);

CREATE TABLE job_checkpoint (
    job_name text PRIMARY KEY,
    last_id bigint NOT NULL,
    high_watermark bigint NOT NULL,
    updated_at timestamptz NOT NULL
);

准备任务时取得上界:

INSERT INTO job_checkpoint (
    job_name,
    last_id,
    high_watermark,
    updated_at
)
SELECT
    'normalize_users',
    0,
    COALESCE(max(id), 0),
    clock_timestamp()
FROM source_user
ON CONFLICT (job_name) DO NOTHING;

这里的 high_watermark 表示本次任务要处理的最大 ID。last_id = 0 只是因为示例假设 ID 为正数;真实系统应根据主键域选择合适的初始值。

2. 每一批的事务流程

应用循环执行以下逻辑:

BEGIN;

-- 读取本次断点
SELECT last_id, high_watermark
FROM job_checkpoint
WHERE job_name = 'normalize_users'
FOR UPDATE;

然后读取下一批:

SELECT id, lower(trim(email)) AS normalized_email, updated_at
FROM source_user
WHERE id > :last_id
  AND id <= :high_watermark
  AND status = 'active'
ORDER BY id
LIMIT 1000;

假设本批返回:

id normalized_email updated_at
101 alice@example.com ...
102 bob@example.com ...
103 carol@example.com ...

应用得到 new_last_id = 103,然后批量写入:

INSERT INTO target_user (
    id,
    normalized_email,
    source_updated_at
)
VALUES
    (101, 'alice@example.com', :updated_at_101),
    (102, 'bob@example.com', :updated_at_102),
    (103, 'carol@example.com', :updated_at_103)
ON CONFLICT (id) DO UPDATE
SET normalized_email = EXCLUDED.normalized_email,
    source_updated_at = EXCLUDED.source_updated_at;

最后推进断点:

UPDATE job_checkpoint
SET last_id = :new_last_id,
    updated_at = clock_timestamp()
WHERE job_name = 'normalize_users';

COMMIT;

3. 为什么断点必须最后提交

正确的提交顺序是:

写入目标
  -> 写入成功
  -> 更新断点
  -> 一次事务提交

如果目标写入和断点在同一个数据库、同一个事务中:

  • 目标写入失败:事务回滚,断点不变;
  • 更新断点失败:事务回滚,目标写入也回滚;
  • 提交成功:目标和断点同时可见。

于是,批次重试最多导致同一批 SQL 再执行一次,但由于目标表使用主键和 Upsert,结果可以保持一致。

4. 这个例子的实际边界

这个例子有一个容易忽略的问题:

AND status = 'active'

如果某行在任务开始后从 active 变成 inactive,它可能被查询时跳过;如果某行先是 inactive,后来变成 active,而它的 ID 已经低于当前断点,也不会被重新扫描。

因此,断点只对固定筛选条件具有明确语义。若业务要求捕获任务期间发生的状态变化,需要改用:

  • updated_at 的增量扫描;
  • 变更数据捕获;
  • 事件表;
  • 任务开始和结束之间的版本号;
  • 或重新执行全量校验。

不能仅靠 last_id 推导出“所有状态变化都被处理”。


八、跨数据库迁移时,断点不再能与写入原子提交

上一个示例的关键条件是:

源读取、目标写入、断点更新在同一个数据库事务中

如果源是 PostgreSQL、目标是 MySQL,或者目标是外部 HTTP 服务,则通常无法通过一个普通本地事务同时提交三者。

典型失败路径如下:

1. 从源库读取 1000 行
2. 写入目标库成功
3. 进程在更新断点前崩溃
4. 重启后重新读取这 1000 行

这会导致目标端重复写入。解决方案不是简单地“先更新断点”,因为反过来会出现:

1. 断点先提交
2. 目标写入失败
3. 重启后从断点之后继续
4. 这一批数据永久漏写

1. 跨系统批处理的安全顺序

通常选择:

先写目标
  -> 确认目标成功
  -> 再提交断点

这样失败时倾向于重复处理,而不是漏处理。前提是目标写入必须可重试且具备幂等性。

2. 目标端幂等键

目标表可以使用源系统的稳定业务键:

CREATE TABLE target_user (
    source_system text NOT NULL,
    source_id bigint NOT NULL,
    normalized_email text NOT NULL,
    PRIMARY KEY (source_system, source_id)
);

重复批次写入时:

INSERT INTO target_user (
    source_system,
    source_id,
    normalized_email
)
VALUES
    ('crm', 101, 'alice@example.com')
ON CONFLICT (source_system, source_id) DO UPDATE
SET normalized_email = EXCLUDED.normalized_email;

这样,崩溃后的重试不会插入第二条同一源记录。

3. 仍然不能忽略版本问题

如果同一源记录会被更新,目标端不能只依据 source_id Upsert。还需要源版本或更新时间:

INSERT INTO target_user (
    source_id,
    normalized_email,
    source_updated_at
)
VALUES (...)
ON CONFLICT (source_id) DO UPDATE
SET normalized_email = EXCLUDED.normalized_email,
    source_updated_at = EXCLUDED.source_updated_at
WHERE target_user.source_updated_at < EXCLUDED.source_updated_at;

这能避免较旧的重试结果覆盖较新的目标数据,但要求:

  • source_updated_at 的语义可靠;
  • 时间精度足够;
  • 不同写入路径不会使用相同或倒退的时间;
  • 更理想时使用单调递增的源版本号。

九、断点恢复的状态模型

1. 一个断点至少包含什么

最简单的断点:

last_id

但生产任务通常还需要:

job_name
partition
last_key
high_watermark
status
attempt
updated_at
error_message

例如按租户分片时,last_id 不能脱离租户:

(job_name = 'migrate_orders', tenant_id = 42, last_order_id = 10000)

否则不同分片之间会互相覆盖进度。

2. 断点推进的状态转移

可以把每一批抽象成状态:

READY
  -> READ
  -> WRITTEN
  -> CHECKPOINTED
  -> COMMITTED

真正需要持久化的是提交后的状态,而不是“应用已经执行过某个函数”。

对于一个批次 [a, b]

  • READ:只代表读到数据;
  • WRITTEN:只代表目标端报告成功,但可能尚未提交;
  • CHECKPOINTED:只代表断点更新语句执行过;
  • COMMITTED:代表事务提交成功,或者跨系统流程已经记录了成功结果。

进程可能在任意步骤崩溃。恢复逻辑必须按最坏情况设计,而不是假设“执行到下一行代码就一定成功”。

3. 是否保存整个批次的 ID 列表

只保存 last_id 的前提是:

  • 查询结果按稳定、唯一、不可后退的键排序;
  • 批次边界由该键唯一确定;
  • 重试整个批次是安全的。

如果批处理过程中存在复杂过滤、外部副作用或不稳定排序,可以额外保存:

batch_start_key
batch_end_key
row_count
checksum

甚至保存批次成员表:

CREATE TABLE job_batch_item (
    job_name text NOT NULL,
    batch_no bigint NOT NULL,
    source_id bigint NOT NULL,
    state text NOT NULL,
    PRIMARY KEY (job_name, batch_no, source_id)
);

但这会增加存储和协调成本。只有当“按边界重算”不再可靠时,才值得记录更细粒度状态。


十、分批事务与长事务的取舍

1. 一个长事务

BEGIN
  读取全部数据
  转换全部数据
  写入全部目标
COMMIT

优点:

  • 一个一致性视图;
  • 失败时可以整体回滚;
  • 不需要处理中间断点。

缺点:

  • 锁和快照持续时间长;
  • PostgreSQL 可能积累旧版本,影响 VACUUM
  • MySQL InnoDB 可能产生长时间事务历史和锁压力;
  • 重试成本大;
  • 连接长期被占用;
  • 事务日志和复制延迟可能出现峰值。

2. 每批一个事务

BEGIN
  读取 1000 行
  写入 1000 行
  更新断点
COMMIT

优点:

  • 失败重试成本小;
  • 锁持有时间较短;
  • 断点清晰;
  • 可以在批次之间限速或暂停。

缺点:

  • 不同批次可能看到不同数据;
  • 需要处理重复和并发修改;
  • 如果源和目标跨系统,无法自然获得整体原子性;
  • 批次过小会增加事务和网络开销。

通常批处理采用每批事务,但这不是无条件的规则。应先定义一致性目标:

  • 是要求事务级快照?
  • 是要求任务开始时存在的数据集?
  • 还是允许处理过程中新增数据?
  • 是允许重复但不允许漏处理,还是必须严格一次?

不同答案会导向不同设计。


十一、批大小不是越大越好

批大小 BB 影响多个量:

  • 单次网络和 SQL 开销;
  • 单批事务持续时间;
  • 锁持有时间;
  • 内存占用;
  • 重试成本;
  • WAL/redo 产生速度;
  • 复制延迟;
  • 失败时重复工作量。

可以把单批耗时粗略分解为:

T(B)=T固定+BT+T(B)+T日志(B)T(B) = T_{\text{固定}} + B \cdot T_{\text{行}} + T_{\text{锁}}(B) + T_{\text{日志}}(B)

其中:

  • T固定T_{\text{固定}}:建立查询、网络往返、事务提交等固定开销;
  • TT_{\text{行}}:每行处理成本;
  • T(B)T_{\text{锁}}(B):批量增大后锁竞争可能增加;
  • T日志(B)T_{\text{日志}}(B):日志刷写、复制和检查点压力。

BB 太小,固定成本占比高;当 BB 太大,锁、日志、内存和重试成本可能主导。

因此批大小应通过实际观测调整,而不是直接套用某个固定数字。观察指标至少包括:

  • 单批查询和写入耗时;
  • p95/p99 延迟;
  • 锁等待;
  • 数据库 CPU、I/O;
  • WAL/redo 生成;
  • 复制延迟;
  • 连接池等待时间;
  • 失败后的重试耗时。

十二、限速:控制的不只是应用循环速度

1. 限速的目标

批处理可能影响:

  • 在线请求延迟;
  • 数据库 CPU 和磁盘;
  • 锁等待;
  • 连接池容量;
  • 复制链路;
  • 下游服务速率限制。

如果应用只是写:

处理完一批
sleep(100ms)

它只限制了循环频率,不一定限制数据库压力。因为单批可能持续 10 秒,或者多个工作线程同时执行,最终压力仍然很高。

2. 固定间隔限速

最简单的方式是每批提交后休眠:

处理批次
提交事务
等待 200 ms
处理下一批

适合低复杂度、低并发任务,但它不能适应数据库当前负载。

3. 按目标速率限速

如果目标是每秒处理 RR 行,某批实际处理 nn 行,批次耗时为 tt 秒,则理论上可等待:

w=max(0,nRt)w = \max(0, \frac{n}{R} - t)

例如:

  • 目标速率 R=1000R = 1000 行/秒;
  • 本批 n=500n = 500 行;
  • 批次耗时 t=0.2t = 0.2 秒。

则:

w=max(0,0.50.2)=0.3 秒w = \max(0, 0.5 - 0.2) = 0.3\text{ 秒}

这样总周期约为 0.5 秒,平均速率接近 1000 行/秒。

这只是应用层速率控制,仍需结合数据库反馈。

4. Token Bucket

令牌桶用固定速率产生令牌,桶容量为 CC,产生速率为 rr 个令牌/秒。每处理一行消耗一个令牌;没有足够令牌时等待。

它允许短时间突发,但长期平均速率受 rr 限制。

如果按批消耗令牌:

需要处理 n 行
  -> 等待获得 n 个令牌
  -> 执行批次

桶容量不能小于允许的最大批次,否则批量操作可能频繁等待。

5. 自适应限速

更可靠的限速通常根据反馈调整:

  • 锁等待增加:降低并发或批大小;
  • 复制延迟增加:暂停或降低速率;
  • 连接池等待增加:减少工作线程;
  • 在线请求 p95 升高:降低迁移速率;
  • 数据库资源恢复:逐步增加速率。

这相当于闭环控制:

执行一批
  -> 采集数据库和业务指标
  -> 判断是否过载
  -> 调整批大小、并发或等待时间
  -> 执行下一批

不能把数据库错误简单归因于“批大小太大”。例如连接池耗尽可能是并发过高或连接泄漏,锁等待可能是访问顺序不一致,复制延迟可能是写入日志速度超过从库消费速度。

6. 限速与连接池

如果有 WW 个工作线程,每个线程都持有一个数据库连接,那么批处理最多可能占用 WW 个连接。在线流量需要的连接数不能被批处理完全挤占。

如果池大小为 PP,在线业务保留 PonlineP_{\text{online}} 个连接,则批处理并发不应只看 CPU,而应受:

WPPonlineW \leq P - P_{\text{online}}

约束。

实际还要考虑连接池中的排队时间、数据库最大连接数、每个任务是否同时占用源库和目标库连接等因素。


十三、批处理中的锁和并发读取

1. 普通读取不等于排他领取

下面的查询只是读取:

SELECT id
FROM work_item
WHERE status = 'ready'
ORDER BY id
LIMIT 100;

两个 worker 可能同时读到同一批数据,然后重复处理。

如果任务是“多个消费者领取待处理任务”,应使用数据库提供的行锁语义。PostgreSQL 和 MySQL InnoDB 都支持类似模式:

BEGIN;

SELECT id, payload
FROM work_item
WHERE status = 'ready'
ORDER BY id
FOR UPDATE SKIP LOCKED
LIMIT 100;

-- 更新这些行为 processing
UPDATE work_item
SET status = 'processing'
WHERE id IN (...);

COMMIT;

SKIP LOCKED 的含义是遇到已被其他事务锁住的行时跳过,而不是等待。它适合工作队列,但不适合被误用为“全量扫描的一致分页机制”。

原因包括:

  • 不同 worker 看到的行集合会动态变化;
  • 行可能被跳过后由其他 worker 处理;
  • 结果不再是按一个全局连续顺序分页;
  • 必须设计处理失败、超时回收和重复执行。

2. 领取与处理是否放在同一事务

如果在事务中锁住行,然后进行长时间外部处理:

BEGIN
  SELECT ... FOR UPDATE
  调用外部服务 30 秒
  UPDATE ...
COMMIT

锁会持续 30 秒,可能阻塞其他业务。

更常见的方式是:

  1. 短事务领取任务并写入 processing 状态;
  2. 提交释放锁;
  3. 应用执行处理;
  4. 成功后短事务写入 done
  5. 失败时写入 failed 或恢复为 ready

这又引入了租约、超时和重复处理问题,因此处理结果必须尽量幂等。


十四、常见失败路径与恢复策略

1. 批次提交前崩溃

流程:

读取批次
写入目标
更新断点
进程在 COMMIT 前崩溃

如果都在同一事务中,数据库会回滚,重启后重新处理该批次。

如果目标和断点在不同系统中,可能出现:

  • 目标已提交;
  • 断点未提交。

此时必须依赖目标幂等写入。

2. 数据库返回超时,但事务结果未知

客户端收到超时,不代表数据库一定回滚了。可能的情况是:

  • SQL 尚未执行;
  • SQL 已执行但响应丢失;
  • 事务已提交但客户端未收到确认;
  • 连接断开后数据库正在回滚。

因此遇到超时不能盲目认为“失败,可以直接重试”。正确做法包括:

  • 使用稳定幂等键;
  • 查询目标端确认结果;
  • 查询事务或批次状态表;
  • 将该批次标记为 unknown,由恢复流程核查。

3. 只更新断点,不记录批次状态

如果只保存:

last_id = 103

而没有任务版本、上界、过滤条件版本等信息,未来可能无法判断这个断点属于哪次任务。

例如任务条件从:

status = 'active'

改为:

status IN ('active', 'pending')

同一个 last_id 在两种规则下含义不同。断点应和任务配置、任务版本或快照边界关联。

4. 把异常行直接跳过

如果一批中有一行格式错误,应用可能采取:

记录错误
继续推进 last_id

这样会保证任务继续,但该行成为永久漏处理。

更明确的做法是建立错误表:

CREATE TABLE job_error (
    job_name text NOT NULL,
    source_id bigint NOT NULL,
    error_code text NOT NULL,
    error_message text NOT NULL,
    first_seen_at timestamptz NOT NULL,
    last_seen_at timestamptz NOT NULL,
    retry_count integer NOT NULL,
    PRIMARY KEY (job_name, source_id)
);

然后明确选择:

  • 错误行阻塞整个任务;
  • 错误行进入隔离表,主任务继续;
  • 错误达到阈值后暂停任务;
  • 修复数据后重新处理错误行。

“继续”不是免费的,它改变了任务的完整性语义。


十五、按范围分片:比单一游标更适合并行处理

如果主键具有可比较的范围,可以把任务划分为:

[1, 100000]
[100001, 200000]
[200001, 300000]

每个分片内部使用 Keyset:

SELECT id, payload
FROM source_table
WHERE id > :last_id
  AND id <= :partition_end
ORDER BY id
LIMIT 1000;

分片状态可以独立保存:

partition_start
partition_end
last_id
status

优点:

  • 可以并行;
  • 单个分片失败不会阻塞全部任务;
  • 断点更容易定位;
  • 可以把不同分片分配到不同 worker。

风险:

  • 分片边界必须互斥且覆盖完整;
  • 并行写目标时可能互相竞争;
  • 如果目标唯一键跨分片冲突,需要协调;
  • 分片数量过多会增加调度和状态管理成本;
  • 主键分布不均时,按 ID 范围不代表按行数均衡。

例如 ID 有大量空洞:

[1, 1000000] 只有 10 行
[1000001, 2000000] 有 1,000,000 行

此时范围分片负载严重不均,需要根据实际数据密度划分,或使用预先生成的分片边界。


十六、增量同步不应只依赖分页

如果任务不是一次性迁移,而是持续同步,单纯依赖:

WHERE id > :last_id

通常不够,因为更新旧行不会改变 id

一种简单的增量条件是:

WHERE (updated_at, id) > (:last_updated_at, :last_id)
ORDER BY updated_at, id
LIMIT 1000;

但这要求:

  • updated_at 每次业务更新都会变化;
  • 时间精度足够;
  • 同一时间戳下有唯一的 id 作为 tie-breaker;
  • 数据库时钟、应用时钟和更新逻辑不会产生倒退;
  • 删除事件有单独记录。

如果一行被删除,直接查询源表已经看不到它,目标端无法知道需要删除。此时需要:

  • 软删除标志;
  • 删除日志;
  • 变更表;
  • 数据库日志捕获;
  • 事件流。

Keyset 负责定位“源表中当前可见的数据”,并不自动提供完整的变更数据捕获能力。


十七、可执行的诊断方法

1. 检查排序和分页是否使用索引

PostgreSQL:

EXPLAIN (ANALYZE, BUFFERS)
SELECT id, payload
FROM source_table
WHERE id > 100000
ORDER BY id
LIMIT 1000;

重点观察:

  • 是否使用预期索引;
  • 是否出现大范围 Sort
  • 实际扫描行数是否远大于返回行数;
  • 是否发生临时文件或大量缓冲区读取。

MySQL:

EXPLAIN ANALYZE
SELECT id, payload
FROM source_table
WHERE id > 100000
ORDER BY id
LIMIT 1000;

关注:

  • 使用的索引;
  • 估算行数与实际行数差异;
  • 是否发生额外排序;
  • 扫描行数是否随断点变化异常。

2. 检查批次是否重复或漏行

迁移任务应保留可验证的统计信息:

读取行数
成功写入行数
实际影响行数
跳过行数
错误行数
重复重试次数
当前断点
最后成功提交时间

对于固定上界的 Keyset 任务,可以检查:

SELECT count(*)
FROM source_table
WHERE id > :last_id
  AND id <= :high_watermark;

但这个数量只表示当前可见的剩余行数。如果源表持续变化,不能把它直接当成最终剩余工作量。

对于目标数据,可使用抽样校验、分段计数、校验和或按主键范围比对:

SELECT
    min(id),
    max(id),
    count(*)
FROM target_table
WHERE id BETWEEN :start_id AND :end_id;

计数一致仍不能证明内容一致,所以必要时还需比较字段摘要或逐行校验。

3. 检查锁等待和连接池排队

批处理变慢时,不要只看 SQL 执行时间。还要区分:

等待连接
  -> 等待锁
  -> 执行 SQL
  -> 等待提交
  -> 等待目标系统

如果连接池等待时间持续增加,增加批量线程只会加剧问题。若锁等待增加,应检查:

  • 不同事务是否以相反顺序访问表;
  • 批次是否过大;
  • 是否缺少索引导致锁定范围扩大;
  • 是否存在长事务;
  • 是否错误地在事务内执行外部调用。

十八、一个完整的设计推导

假设需求是:

将 PostgreSQL 中任务开始时存在的 active 用户迁移到目标表;允许进程重启;目标库与源库不同;不能明显影响在线业务。

可以按以下步骤推导。

第一步:确定数据集边界

任务开始时取得:

SELECT max(id)
FROM source_user;

记为 HH。本次任务只处理:

idHid \leq H

避免任务执行期间新增数据无限进入本次任务。

第二步:确定排序键

使用:

id

要求:

  • 唯一;
  • 非空;
  • 不会改变;
  • 可以通过索引高效范围扫描。

第三步:确定下一批条件

断点为 CC,查询:

WHERE id > C
  AND id <= H
  AND status = 'active'
ORDER BY id
LIMIT B

其中 BB 是批大小。

第四步:确定失败语义

源库和目标库不同,不能让目标写入与断点更新进入同一个本地事务。因此选择:

先写目标,再推进断点

这会偏向“可能重复,不应漏写”。

第五步:保证目标幂等

目标使用:

source_id 作为唯一键

写入使用 Upsert,并用源版本防止旧数据覆盖新数据。

第六步:确定限速反馈

初始使用较小批量,逐步增加,同时观测:

  • 源库查询延迟;
  • 目标库写入延迟;
  • 在线请求延迟;
  • 复制延迟;
  • 连接池等待;
  • 锁等待。

出现过载时降低 BB、降低并发或增加批间等待。

第七步:定义重启动作

重启后:

  1. 读取持久化断点 CC
  2. 重新查询 id > C 的批次;
  3. 目标端通过唯一键 Upsert;
  4. 成功后推进断点;
  5. 任务达到 HH 后结束;
  6. 对错误表和统计信息做最终核查。

这套设计的结论不是“绝对一次”,而是:

源端范围内尽量不漏;
目标端允许重复提交;
重复提交通过幂等键收敛;
断点在目标成功后推进。

这是跨系统批处理常见且可验证的语义。


十九、容易混淆的结论

“游标就是断点恢复”

不是。游标位置通常依赖数据库连接和会话。连接丢失后,游标可能无法继续。断点恢复需要把可重建的位置持久化到应用或数据库中。

“Keyset 一定不会重复或漏行”

不是。它要求排序键稳定、唯一,并且查询边界语义明确。键值修改、过滤条件变化、NULL、删除、跨系统提交都可能破坏简单推理。

“批量写就是一次提交全部数据”

不是。批量写和事务大小是两个维度。可以每 1000 行组成一个批量 SQL,但每 10 个批次才提交;也可以每个批次独立提交。两者的锁、恢复和日志行为不同。

“加 sleep 就完成限速”

不是。多个 worker、长事务、锁等待、连接池排队和数据库复制延迟都可能使实际压力与 sleep 时间无关。

“Upsert 后就可以放心重试”

只有当写入内容和副作用都具备幂等语义时才成立。递增计数、发送消息、写审计日志、更新时间戳等操作可能在重试时产生额外效果。

“高水位线可以保证数据完全一致”

高水位线只能限制键范围。它不能冻结行内容,也不能记录范围内的删除和更新。若需要严格快照,应使用数据库提供的快照或一致性导出机制;若需要持续捕获变化,应使用变更记录或 CDC。


二十、最终应明确的批处理契约

一个批处理任务在实现前,至少应写清楚以下契约:

  1. 数据集:处理哪些行,任务开始后新增行是否包含在本次任务中;
  2. 顺序:使用哪个唯一排序键,键是否可变;
  3. 边界:使用 last_id、复合键、时间窗口还是固定范围;
  4. 事务:每批提交还是长事务,源和目标是否同库;
  5. 写入语义:插入、更新、Upsert,冲突如何处理;
  6. 失败语义:允许重复还是允许漏行,外部副作用如何重试;
  7. 断点:何时推进,保存哪些任务版本和范围信息;
  8. 限速:限制行速率、批速率、并发数还是数据库资源占用;
  9. 并发:是否多个 worker,是否需要 FOR UPDATE SKIP LOCKED
  10. 验证:如何检查范围完整性、内容一致性和错误行。

当这些条件都能被明确回答时,游标、Keyset、批量写和限速就不再是孤立的 SQL 技巧,而是同一个可恢复数据处理协议的不同组成部分。


系列导航与关联阅读

官方资料

本文依据数据库官方文档重新梳理;正文、示例与生产检查清单由 WR BLOG 编写。