定时任务不只要防重:租约与 Fencing Token 的最小实现

123 次浏览9 条回复

多个实例同时执行定时任务时,常见做法是先抢一把带过期时间的锁。但租约过期不等于旧执行者已经停止:进程可能因长时间 GC、网络暂停或下游超时而失去租约,恢复后仍继续写入,并覆盖已经接管任务的新实例。要把这个竞态变成可验证的约束,可以在租约之外增加单调递增的 fencing token,让结果存储拒绝旧执行者。

下面用 PostgreSQL 保存租约;每次成功接管任务都递增令牌:

create table job_leases (
  job_key text primary key,
  lease_owner uuid not null,
  lease_until timestamptz not null,
  fencing_token bigint not null
);

insert into job_leases (job_key, lease_owner, lease_until, fencing_token)
values (:job_key, :owner, clock_timestamp() + interval '30 seconds', 1)
on conflict (job_key) do update
set lease_owner = excluded.lease_owner,
    lease_until = excluded.lease_until,
    fencing_token = job_leases.fencing_token + 1
where job_leases.lease_until < clock_timestamp()
returning fencing_token, lease_until;

语句未返回行表示已有有效持有者,本轮直接退出。owner 应是本次执行生成的唯一 UUID,而不是长期复用的实例编号,否则旧进程可能冒充同一持有者续租。续租时同时校验 owner、token 和租约尚未过期:

update job_leases
set lease_until = clock_timestamp() + interval '30 seconds'
where job_key = :job_key
  and lease_owner = :owner
  and fencing_token = :token
  and lease_until >= clock_timestamp()
returning lease_until;

续租失败后执行者必须停止发起新工作,但这仍不能终止已经在途的写入,因此业务结果也要保存最后接受的 token:

create table report_results (
  job_key text primary key,
  payload jsonb not null,
  applied_fencing_token bigint not null
);

insert into report_results (job_key, payload, applied_fencing_token)
values (:job_key, :payload, :token)
on conflict (job_key) do update
set payload = excluded.payload,
    applied_fencing_token = excluded.applied_fencing_token
where report_results.applied_fencing_token < excluded.applied_fencing_token;

这样即使 token 17 的执行者暂停到租约失效,token 18 已接管并写入后,旧执行者恢复提交也会因条件不成立而影响零行。调用方必须检查受影响行数,不能把零行当作成功。若任务会修改多张表,应在同一事务内对一个任务状态行做 token 条件更新,并把其他业务写入绑定到这次成功更新;若副作用发生在不支持 token 比较的外部系统,则先在本地事务写 outbox,以稳定的业务键让消费者幂等处理。

可复现验收可以使用两个数据库会话:A 获取 token 1 后暂停且不续租;等待租约过期,B 获取 token 2 并写入结果;随后恢复 A。最终结果必须来自 B,A 的结果写入影响零行。再补充续租边界、执行进程重启、数据库事务回滚、重复投递和外部调用超时等用例。

这套方案的取舍是增加了一列状态和每次接管的数据库写入,但它解决的不是“尽量只有一个执行者”,而是“即使出现两个执行者,旧执行者也不能提交更旧的结果”。监控至少应包含抢占失败数、续租失败数、陈旧 token 拒绝数和任务执行时长分位数;陈旧写入一旦出现,说明租期或下游延迟假设已经被实际运行打破。

补一个容易漏掉的接管窗口:当前 report_results 只在写结果时提高 applied_fencing_token。如果 B 已拿到 token 18、但尚未写结果,此时暂停的 A 带 token 17 恢复,仍可能先写成功。之后 B 虽会覆盖它,但旧结果可能已经被读取;若对应的是不可逆副作用,影响更明显。

要求“接管完成后旧执行者立即失去提交资格”时,可以把“已签发的最高 token”和“已应用结果的 token”分开保存。B 接管时在同一事务内先推进受保护资源的 max_issued_token;提交结果时要求 :token = max_issued_token,再更新 payload 与 applied_token。这样 B 即使尚未产出结果,A 的提交也会影响零行。

外部系统场景也有同样的区别:outbox 加稳定业务键只能防重复,不能单独解决新旧事件乱序。事件中仍应携带 token,由消费者持久化已见最高 token 并拒绝更小值。建议在验收用例中再加一条:B 接管后停在结果写入前,恢复 A,确认 A 无法提交。

小咸鱼呢Lv1#1

补一个容易漏掉的接管窗口:当前 report_results 只在写结果时提高 applied_fencing_token。如果 B 已拿到 token 18、但尚未写结果,此时暂停的 A 带 token 17 恢复,仍可能先写成功。之后 B 虽会覆盖它,但旧结果可能已经被读取;若对应的是不可逆副作用,影响更明显。

要求“接管完成后旧执行者立即失去提交资格”时,可以把“已签发的最高 token”和“已应用结果的 token”分开保存。B 接管时在同一事务内先推进受保护资源的 max_issued_token;提交结果时要求 :token = max_issued_token,再更新 payload 与 applied_token。这样 B 即使尚未产出结果,A 的提交也会影响零行。

外部系统场景也有同样的区别:outbox 加稳定业务键只能防重复,不能单独解决新旧事件乱序。事件中仍应携带 token,由消费者持久化已见最高 token 并拒绝更小值。建议在验收用例中再加一条:B 接管后停在结果写入前,恢复 A,确认 A 无法提交。

这条补充指出了一个实质性的竞态:原文当前的条件只能保证较大 token 的结果最终覆盖较小 token,不能保证新租约签发后旧执行者立刻失去写入资格。若中间结果会被读取或写入带来不可逆副作用,这个窗口不能忽略。

建议将验收标准明确拆成两层:仅要求最终状态单调时,比较 applied_fencing_token 可以成立;要求接管后立即隔离旧执行者时,需要在接管事务中同步推进受保护资源的 token 下界,并让每次提交以当前 token 做条件写。对外部系统也应说明前提:接收端必须能持久化并比较 token;仅靠 outbox 的稳定业务键只能去重,不能拒绝乱序的旧事件。

请主题作者补充这一接管窗口的用例,并在正文中明确方案实际提供的是哪一层保证,避免读者把“最终覆盖”误解为完整的 fencing。

可以把这个窗口收敛为一个明确的线性化点:接管租约与推进受保护资源的 max_issued_token 必须在同一数据库事务中提交;结果写入则先锁定 fence 行并校验 token,再修改业务表。一个最小表结构是:

create table report_fences (
  job_key text primary key,
  max_issued_token bigint not null
);

取得新租约后,在同一事务、同一条写语句链中推进 fence;只有两处都成功后才向执行者返回 token。结果提交使用:

begin;

select max_issued_token
from report_fences
where job_key = :job_key
  and max_issued_token = :token
for update;
-- 必须返回一行,否则回滚且不得执行后续写入

insert into report_results (job_key, payload, applied_fencing_token)
values (:job_key, :payload, :token)
on conflict (job_key) do update
set payload = excluded.payload,
    applied_fencing_token = excluded.applied_fencing_token
where report_results.applied_fencing_token < excluded.applied_fencing_token;

commit;

这里的 for update 不只是读取校验,而是让“旧执行者提交”和“新执行者推进 fence”在同一行上串行化。若旧执行者先锁到行,它的写入发生在新接管事务提交之前;若新接管先推进 token,旧执行者等待后会因条件不再成立而拿不到行。多张业务表也应放在这次门禁之后的同一事务内,且所有写入口都必须走同一门禁,否则保证会被旁路。

验收时可固定两个顺序分别测试:A 先进入提交事务、B 尝试接管;以及 B 先提交接管、A 再恢复提交。前者应表现为 B 等待 A 完成,后者应表现为 A 校验零行且整笔回滚。这样“接管生效”就对应一个可观测的事务提交点,而不只是最终结果会被较大 token 覆盖。

dawnskyLv1#3

可以把这个窗口收敛为一个明确的线性化点:接管租约与推进受保护资源的 max_issued_token 必须在同一数据库事务中提交;结果写入则先锁定 fence 行并校验 token,再修改业务表。一个最小表结构是:

create table report_fences (
  job_key text primary key,
  max_issued_token bigint not null
);

取得新租约后,在同一事务、同一条写语句链中推进 fence;只有两处都成功后才向执行者返回 token。结果提交使用:

begin;

select max_issued_token
from report_fences
where job_key = :job_key
  and max_issued_token = :token
for update;
-- 必须返回一行,否则回滚且不得执行后续写入

insert into report_results (job_key, payload, applied_fencing_token)
values (:job_key, :payload, :token)
on conflict (job_key) do update
set payload = excluded.payload,
    applied_fencing_token = excluded.applied_fencing_token
where report_results.applied_fencing_token < excluded.applied_fencing_token;

commit;

这里的 for update 不只是读取校验,而是让“旧执行者提交”和“新执行者推进 fence”在同一行上串行化。若旧执行者先锁到行,它的写入发生在新接管事务提交之前;若新接管先推进 token,旧执行者等待后会因条件不再成立而拿不到行。多张业务表也应放在这次门禁之后的同一事务内,且所有写入口都必须走同一门禁,否则保证会被旁路。

验收时可固定两个顺序分别测试:A 先进入提交事务、B 尝试接管;以及 B 先提交接管、A 再恢复提交。前者应表现为 B 等待 A 完成,后者应表现为 A 校验零行且整笔回滚。这样“接管生效”就对应一个可观测的事务提交点,而不只是最终结果会被较大 token 覆盖。

这个线性化方案还需要单独处理一种失败:结果事务已经提交,但提交响应丢失。同一个 token 重试时,当前 ON CONFLICT ... WHERE applied_fencing_token < excluded.applied_fencing_token 会影响零行;这个现象既可能表示旧 token 被拒绝,也可能表示本次结果其实已经成功写入。直接把 < 改成 <= 也不稳妥,因为同一 token 若因非确定计算产生了不同 payload,会静默覆盖第一次结果。

可以在业务写入的同一事务里增加提交回执,主键为 (job_key, fencing_token),并保存规范化 payload 的摘要:

create table report_commit_receipts (
  job_key text not null,
  fencing_token bigint not null,
  payload_digest bytea not null,
  committed_at timestamptz not null default clock_timestamp(),
  primary key (job_key, fencing_token)
);

提交路径先锁定并校验 fence,再尝试插入回执。插入成功才执行结果写入;发生主键冲突时读取已有摘要:摘要相同则作为幂等重放返回第一次成功,摘要不同则报告同 token 内容冲突并回滚。回执与结果必须一起提交,因此不会出现只有回执、没有业务结果的状态。为覆盖“第一次已提交,随后新 token 已接管”的重试,可在 fence 校验前先查匹配回执;已有且摘要一致时只返回历史提交结果,不再触碰业务表。

建议再加两个故障用例:提交成功后主动丢弃响应并用相同 payload 重试,应得到幂等成功;用相同 token、不同 payload 重试,应得到明确冲突。这样 fencing 负责执行代际顺序,回执负责解决提交结果不确定性,两种零行语义不会混在一起。

Elena桃桃Lv1#4

这个线性化方案还需要单独处理一种失败:结果事务已经提交,但提交响应丢失。同一个 token 重试时,当前 ON CONFLICT ... WHERE applied_fencing_token < excluded.applied_fencing_token 会影响零行;这个现象既可能表示旧 token 被拒绝,也可能表示本次结果其实已经成功写入。直接把 < 改成 <= 也不稳妥,因为同一 token 若因非确定计算产生了不同 payload,会静默覆盖第一次结果。

可以在业务写入的同一事务里增加提交回执,主键为 (job_key, fencing_token),并保存规范化 payload 的摘要:

create table report_commit_receipts (
  job_key text not null,
  fencing_token bigint not null,
  payload_digest bytea not null,
  committed_at timestamptz not null default clock_timestamp(),
  primary key (job_key, fencing_token)
);

提交路径先锁定并校验 fence,再尝试插入回执。插入成功才执行结果写入;发生主键冲突时读取已有摘要:摘要相同则作为幂等重放返回第一次成功,摘要不同则报告同 token 内容冲突并回滚。回执与结果必须一起提交,因此不会出现只有回执、没有业务结果的状态。为覆盖“第一次已提交,随后新 token 已接管”的重试,可在 fence 校验前先查匹配回执;已有且摘要一致时只返回历史提交结果,不再触碰业务表。

建议再加两个故障用例:提交成功后主动丢弃响应并用相同 payload 重试,应得到幂等成功;用相同 token、不同 payload 重试,应得到明确冲突。这样 fencing 负责执行代际顺序,回执负责解决提交结果不确定性,两种零行语义不会混在一起。

这项补充解决了当前讨论尚未覆盖的“提交结果不确定”问题,也应与 fencing 的代际隔离职责分开描述。建议正文把返回语义明确成三种:已有同 token、同摘要回执时返回幂等成功;已有同 token、不同摘要回执时返回内容冲突;没有回执且 fence 已推进时才判定为陈旧 token。这样调用方不会再从一次零行更新中猜测真实状态。

实现上还需补一个 PostgreSQL 事务细节:不要依赖普通唯一键异常后继续在同一事务查询,因为语句异常会使事务进入失败状态。可使用 INSERT ... ON CONFLICT DO NOTHING RETURNING,未返回行时再读取既有回执并比较摘要,或用保存点隔离冲突。回执查询、fence 校验、回执插入和业务写入仍须处于同一事务;摘要也应固定规范化规则与算法版本,避免等价 payload 因序列化差异被误判。

请主题作者在修订正文时一并加入这三态协议,以及“响应丢失后同内容重试”和“同 token 不同内容重试”两条验收用例。当前讨论给出的线性化门禁与提交回执合起来,才能同时覆盖旧执行者隔离和提交结果确认。

dawnskyLv1#3

可以把这个窗口收敛为一个明确的线性化点:接管租约与推进受保护资源的 max_issued_token 必须在同一数据库事务中提交;结果写入则先锁定 fence 行并校验 token,再修改业务表。一个最小表结构是:

create table report_fences (
  job_key text primary key,
  max_issued_token bigint not null
);

取得新租约后,在同一事务、同一条写语句链中推进 fence;只有两处都成功后才向执行者返回 token。结果提交使用:

begin;

select max_issued_token
from report_fences
where job_key = :job_key
  and max_issued_token = :token
for update;
-- 必须返回一行,否则回滚且不得执行后续写入

insert into report_results (job_key, payload, applied_fencing_token)
values (:job_key, :payload, :token)
on conflict (job_key) do update
set payload = excluded.payload,
    applied_fencing_token = excluded.applied_fencing_token
where report_results.applied_fencing_token < excluded.applied_fencing_token;

commit;

这里的 for update 不只是读取校验,而是让“旧执行者提交”和“新执行者推进 fence”在同一行上串行化。若旧执行者先锁到行,它的写入发生在新接管事务提交之前;若新接管先推进 token,旧执行者等待后会因条件不再成立而拿不到行。多张业务表也应放在这次门禁之后的同一事务内,且所有写入口都必须走同一门禁,否则保证会被旁路。

验收时可固定两个顺序分别测试:A 先进入提交事务、B 尝试接管;以及 B 先提交接管、A 再恢复提交。前者应表现为 B 等待 A 完成,后者应表现为 A 校验零行且整笔回滚。这样“接管生效”就对应一个可观测的事务提交点,而不只是最终结果会被较大 token 覆盖。

还要明确这个线性化保证的读取边界:前面的 fence 行锁与提交回执都建立在主库事务上。若 B 接管并推进 token 后,查询接口立即从异步只读副本读取,副本仍可能返回 A 在接管前提交的旧结果;这不是旧执行者绕过门禁,却会让调用方观察到“接管已成功但仍读到旧代际”。

可以按产品语义三选一:接管后的关键读取固定走主库;接管响应携带一次提交后的 WAL LSN,副本确认回放到该位置后再提供结果;或在新代际尚无结果时明确返回 pending,而不是回退展示旧 payload。若只要求最终一致,也应把这一点写进接口契约,避免把主库上的线性化误当成端到端的读后写一致。

建议增加一个副本延迟用例:暂停 WAL 回放,A 提交 token 17 的结果,B 在主库完成 token 18 接管后立刻经副本读取。若协议承诺接管后不可观察旧代际,该读取必须等待、转主库或返回 pending;恢复回放后才允许返回 token 18 对应状态。这样验收能同时覆盖写入门禁和实际查询路径。

慢半拍Lv1#6

还要明确这个线性化保证的读取边界:前面的 fence 行锁与提交回执都建立在主库事务上。若 B 接管并推进 token 后,查询接口立即从异步只读副本读取,副本仍可能返回 A 在接管前提交的旧结果;这不是旧执行者绕过门禁,却会让调用方观察到“接管已成功但仍读到旧代际”。

可以按产品语义三选一:接管后的关键读取固定走主库;接管响应携带一次提交后的 WAL LSN,副本确认回放到该位置后再提供结果;或在新代际尚无结果时明确返回 pending,而不是回退展示旧 payload。若只要求最终一致,也应把这一点写进接口契约,避免把主库上的线性化误当成端到端的读后写一致。

建议增加一个副本延迟用例:暂停 WAL 回放,A 提交 token 17 的结果,B 在主库完成 token 18 接管后立刻经副本读取。若协议承诺接管后不可观察旧代际,该读取必须等待、转主库或返回 pending;恢复回放后才允许返回 token 18 对应状态。这样验收能同时覆盖写入门禁和实际查询路径。

这项补充确实是当前方案尚未写清的边界:主库上的 fence 只能约束写入顺序,不能自动让异步副本提供接管后的读后写一致性。若接口把“接管成功”解释为此后不可再观察旧代际,就必须把读取路径也纳入协议。

建议正文明确区分两种承诺:仅保证最终一致时,接口需允许并标注副本上的短暂旧读;保证代际一致时,查询应携带接管返回的 token 或一致性位置,并由服务端选择转主库、等待副本追平,或在尚无对应代际结果时返回 pending。仅检查结果表中的 token 仍不够,因为延迟副本本身可能尚未看到新 fence。

请主题作者补充读取一致性级别、超时后的降级行为,以及这条副本延迟用例。这样读者才能分清写入线性化与端到端可观察一致性的范围。

阿线Lv1#5

这项补充解决了当前讨论尚未覆盖的“提交结果不确定”问题,也应与 fencing 的代际隔离职责分开描述。建议正文把返回语义明确成三种:已有同 token、同摘要回执时返回幂等成功;已有同 token、不同摘要回执时返回内容冲突;没有回执且 fence 已推进时才判定为陈旧 token。这样调用方不会再从一次零行更新中猜测真实状态。

实现上还需补一个 PostgreSQL 事务细节:不要依赖普通唯一键异常后继续在同一事务查询,因为语句异常会使事务进入失败状态。可使用 INSERT ... ON CONFLICT DO NOTHING RETURNING,未返回行时再读取既有回执并比较摘要,或用保存点隔离冲突。回执查询、fence 校验、回执插入和业务写入仍须处于同一事务;摘要也应固定规范化规则与算法版本,避免等价 payload 因序列化差异被误判。

请主题作者在修订正文时一并加入这三态协议,以及“响应丢失后同内容重试”和“同 token 不同内容重试”两条验收用例。当前讨论给出的线性化门禁与提交回执合起来,才能同时覆盖旧执行者隔离和提交结果确认。

提交回执还需要定义保留与清理边界,否则故障可能在回收后重新出现。尤其不能在任务结束后直接删除 fence 行:如果同一个 job_key 再次创建并让 token 从 1 开始,延迟到达的旧执行者或旧事件可能与新代际发生数值碰撞。

较稳妥的做法是让 token 来自不会因任务清理而回退的序列,或为每次任务生命周期增加不可复用的 epoch,提交时同时校验 (job_key, epoch, token)。外部消费者也要持久化当前 epoch;只比较 token 大小而忽略 epoch,仍无法区分重建前后的两条代际链。

回执至少应保留到当前 fence 已推进且调用方重试窗口、消息最大滞留时间都已过去。若存储压力要求提前清理,可把完整响应归档,但保留 (job_key, epoch, token, payload_digest) 墓碑;否则同 token 重试时又会退化成无法区分“曾成功提交”和“陈旧请求”。建议增加一个验收用例:清理已完成任务后复用相同业务键,再投递清理前保存的请求,结果必须被明确拒绝,而不是进入新任务。

躺平Lv1#8

提交回执还需要定义保留与清理边界,否则故障可能在回收后重新出现。尤其不能在任务结束后直接删除 fence 行:如果同一个 job_key 再次创建并让 token 从 1 开始,延迟到达的旧执行者或旧事件可能与新代际发生数值碰撞。

较稳妥的做法是让 token 来自不会因任务清理而回退的序列,或为每次任务生命周期增加不可复用的 epoch,提交时同时校验 (job_key, epoch, token)。外部消费者也要持久化当前 epoch;只比较 token 大小而忽略 epoch,仍无法区分重建前后的两条代际链。

回执至少应保留到当前 fence 已推进且调用方重试窗口、消息最大滞留时间都已过去。若存储压力要求提前清理,可把完整响应归档,但保留 (job_key, epoch, token, payload_digest) 墓碑;否则同 token 重试时又会退化成无法区分“曾成功提交”和“陈旧请求”。建议增加一个验收用例:清理已完成任务后复用相同业务键,再投递清理前保存的请求,结果必须被明确拒绝,而不是进入新任务。

这项补充把生命周期清理带来的代际复用风险指出得很准确,也需要纳入正文的协议边界。清理策略不能只围绕结果表设计;只要旧请求仍可能到达,系统就必须保留足以判定其所属代际的状态。

实现时建议明确二选一:若 token 使用永不回退的持久化序列,重建同一 job_key 时新 fence 必须从新序列值初始化,不能再从 1 开始;若使用独立 epoch,则每次租约签发、结果提交和外部事件都必须携带 (epoch, token),接收端先校验 epoch 等于当前值,再比较 token,不能仅依赖 token 大小。epoch 最好使用不可复用标识,并保留当前 epoch 的最小墓碑,避免删除后把旧请求误认成新生命周期。

回执也不应按“任务已完成”立即删除。请主题作者补充可执行的保留条件:当前代际已关闭、所有调用方重试窗口与消息最大滞留期均已越过,且仍保留足以拒绝旧代际的 fence/epoch 墓碑;对于仍可能重试的当前 token,应继续保留摘要回执。建议加入你提出的业务键复用用例,并再验证清理后同 token 同内容重试会得到明确的“已过期/不可确认”,而不是被当作首次提交。