Python 基础体系 · 第 105/112 篇。示例统一以 Python 3.14 为语言基线;第三方库使用与其兼容的现代稳定版本,版本敏感行为会单独说明。

Python 定时任务:时间语义、调度器、重复执行、锁和补偿

定时任务不是“每隔一段时间调用一个函数”这么简单。一个可靠的定时任务系统至少要回答六个问题:

  1. 什么时候算到期:按经过的时间,还是按日历上的某个本地时间?
  2. 由哪一个时钟判断到期:系统墙上时钟,还是不会倒退的单调时钟?
  3. 任务执行几次:允许重复、尽量一次,还是必须显式补偿?
  4. 多个线程或进程同时发现到期时怎么办:谁获得执行权?
  5. 任务执行到一半进程崩溃时怎么办:重新执行、跳过,还是恢复中间状态?
  6. 调度器重启后怎么办:依赖内存中的队列,还是从持久化记录恢复?

这些问题分别对应时间语义、调度器、重复执行、锁和补偿。它们不是互相独立的选项,而是一条故障链:

flowchart LR
    A[定义时间语义] --> B[计算下一次到期时间]
    B --> C[调度器等待]
    C --> D[任务到期]
    D --> E[竞争执行权]
    E --> F[执行外部副作用]
    F --> G[记录成功或失败]
    G --> H[重试/补偿/下一次调度]

如果第一步的时间语义不明确,后面的锁、重试和补偿都可能建立在错误的时间点上。


一、先区分三个“时间”

1. 墙上时钟:回答“现在几点”

墙上时钟,也称为日历时钟或实时时钟,用来表示现实世界中的日期和时间,例如:

from datetime import datetime
from zoneinfo import ZoneInfo

now = datetime.now(ZoneInfo("Asia/Shanghai"))
print(now)

这个值适合表达:

  • 2026 年 9 月 1 日 09:00;
  • 用户所在时区的营业时间;
  • 每天本地时间 02:30 执行;
  • 数据库中的创建时间和完成时间。

带有 tzinfodatetime 称为 aware datetime,能够表示时区;没有 tzinfo 的称为 naive datetime。Python 的 zoneinfo 提供基于 IANA 时区数据库的时区实现,系统没有时区数据时还可以使用 tzdata 包作为数据源。(docs.python.org)

2. 单调时钟:回答“经过了多久”

单调时钟只保证时间向前推进,不代表某个现实日期。time.monotonic() 的参考起点没有业务意义,只有两个读数之间的差值有效;系统时间被管理员修改或通过 NTP 校准时,单调时钟不会因此倒退。(docs.python.org)

import time

started = time.monotonic()
# 执行任务
elapsed = time.monotonic() - started
print(f"任务耗时 {elapsed:.3f} 秒")

它适合:

  • 等待 10 秒;
  • 设置锁的超时时间;
  • 判断任务是否执行超过 5 分钟;
  • 计算调度延迟;
  • 实现“从现在起经过一个间隔”。

它不适合直接写入“下次执行时间”:

# 不要把 monotonic() 的值写成业务时间戳
deadline = time.monotonic() + 3600

因为这个值不能被其他机器、数据库或运维人员解释为某个日期。

3. Unix 时间戳:在系统之间传输时刻

Unix 时间戳通常表示自 1970 年 1 月 1 日 00:00:00 UTC 以来经过的秒数。它表达的是一个绝对时刻,但不表达用户希望看到的时区和本地日历语义。Python 文档也明确指出,time.time() 返回的值可能因为系统时钟被调回而变小。(docs.python.org)

from datetime import datetime, timezone

instant = datetime.now(timezone.utc)
timestamp = instant.timestamp()

print(instant.isoformat())
print(timestamp)

可以把三者概括为:

类型 表示什么 是否可能倒退 典型用途
墙上时钟 现实世界的日期时间 可能 “每天 9 点”
单调时钟 经过的时间 不会 “等待 10 秒”
Unix 时间戳 UTC 绝对时刻 取决于来源时钟 存储、传输、比较

一个定时系统通常需要同时使用两种时间:

  • 墙上时钟计算业务上的到期时刻;
  • 单调时钟等待和测量执行时长。

二、时间语义:间隔任务和日历任务不是一回事

“每隔一天执行一次”和“每天本地时间 09:00 执行”看起来相近,但它们是两个不同的调度模型。

1. 间隔语义

间隔语义定义为:

tn+1=tn+Δt_{n+1} = t_n + \Delta

其中:

  • tnt_n 是第 nn 次任务的参考时刻;
  • Δ\Delta 是固定间隔,例如 3600 秒;
  • tn+1t_{n+1} 是下一次参考时刻。

例如,任务在 09:00:00 开始,间隔为 1 小时:

09:00:00
10:00:00
11:00:00
12:00:00

这里的“一小时”是经过时间,适合:

  • 每 5 分钟采集一次指标;
  • 每 30 秒刷新一次缓存;
  • 任务完成后等待 10 分钟再次检查。

间隔任务应该使用单调时钟计算等待时间:

import time

interval = 10.0
next_deadline = time.monotonic()

for _ in range(3):
    next_deadline += interval
    delay = max(0.0, next_deadline - time.monotonic())
    time.sleep(delay)
    print("执行一次")

这里使用“目标时刻减当前时刻”的方式,而不是简单地在任务结束后 sleep(interval)。两者差异如下:

固定延迟:
执行开始 ──耗时 3 秒── 执行结束 ──等待 10 秒── 下一次开始
下一次间隔 = 3 + 10 = 13 秒

固定频率:
计划开始 ──每 10 秒一个目标点──
09:00:00、09:00:10、09:00:20

2. 日历语义

日历语义定义为:

tn=f(日期n,本地时间,时区规则)t_n = f(\text{日期}_n, \text{本地时间}, \text{时区规则})

例如“每天 09:00,Asia/Shanghai”:

2026-09-01 09:00 +08:00
2026-09-02 09:00 +08:00
2026-09-03 09:00 +08:00

它不是简单地给上一次时间加 timedelta(days=1)。在有夏令时的地区,“本地每天 09:00”与“每隔 24 个实际小时”可能产生不同结果。

因此,日历任务应保存:

频率:每天
本地时间:09:00
时区:America/Los_Angeles
缺失时间策略:顺延
重复时间策略:执行一次,选择第一次

而不是只保存一个 UTC 时间戳。UTC 时间戳只能表示已经计算出的某个时刻,不能完整表达“下一个本地日历事件”。


三、时区和 DST:本地时间可能不存在,也可能出现两次

DST 是 Daylight Saving Time,即夏令时。时区发生偏移变化时,本地时间会出现两个特殊区域。

1. Gap:本地时间不存在

春季拨快时,假设时钟从 01:59:59 跳到 03:00:00:

01:59:59
03:00:00

这意味着本地时间 02:30 根本不存在。

如果业务配置为:

每天 02:30 执行

那么在转换日期发生变化的那一天,必须选择策略:

  • 跳过这次;
  • 顺延到 03:00;
  • 顺延到下一个有效时间;
  • 转换为固定 UTC 时刻。

不能假设 Python 或操作系统会自动替你做出符合业务的选择。

2. Fold:本地时间出现两次

秋季拨回时,假设时钟从 02:00 回到 01:00:

01:00(夏令时)
01:30(夏令时)
01:00(标准时)
01:30(标准时)

这两个 01:30 的 UTC 时刻不同。Python 使用 datetime.fold 区分重复的本地时间:

from datetime import datetime
from zoneinfo import ZoneInfo

la = ZoneInfo("America/Los_Angeles")

first = datetime(2020, 11, 1, 1, 30, tzinfo=la, fold=0)
second = datetime(2020, 11, 1, 1, 30, tzinfo=la, fold=1)

print(first.isoformat())
print(second.isoformat())
print(first.timestamp() != second.timestamp())

fold=0 表示偏移变化前的那个本地时间,fold=1 表示偏移变化后的那个本地时间。zoneinfo 支持这一语义,并且从 UTC 转换到本地时会自动设置正确的 fold 值。(docs.python.org)

3. 不要用 replace(tzinfo=...) 做时区转换

下面的代码只是把标签贴上去:

from datetime import datetime
from zoneinfo import ZoneInfo

naive = datetime(2026, 9, 1, 9, 0)
wrong = naive.replace(tzinfo=ZoneInfo("Asia/Shanghai"))

对于一个原本不知道时区的本地时间,这表示“把它解释为上海时间”,并没有把它从某个时区转换到另一个时区。

如果已知原始时间是 UTC,应使用 astimezone()

from datetime import datetime, timezone
from zoneinfo import ZoneInfo

utc_time = datetime(2026, 9, 1, 1, 0, tzinfo=timezone.utc)
shanghai_time = utc_time.astimezone(ZoneInfo("Asia/Shanghai"))

print(shanghai_time)

astimezone() 会调整日期和时钟字段,使转换前后的对象表示同一个 UTC 时刻。(docs.python.org)

4. 计算“每天某个本地时间”的一个可验证实现

下面的实现展示一种明确策略:

  • 输入是带时区的当前时刻;
  • 目标是下一个本地日期的指定时间;
  • 如果目标时间属于 Gap,则顺延到转换后的有效时间;
  • 如果目标时间属于 Fold,则选择 fold=0
  • 通过 UTC 往返转换验证本地时间是否真实存在。
from __future__ import annotations

from datetime import datetime, date, time, timedelta, timezone
from zoneinfo import ZoneInfo


def resolve_local(
    local_date: date,
    local_time: time,
    tz: ZoneInfo,
) -> datetime:
    """
    将本地日期和本地时间解析为带时区 datetime。

    策略:
    - Fold:选择 fold=0;
    - Gap:向前移动,直到往返转换后的本地时间不再变化。
    """
    candidate = datetime.combine(
        local_date,
        local_time,
        tzinfo=tz,
        fold=0,
    )

    # 本地时间 -> UTC -> 本地时间。
    # 若结果不同,说明 candidate 落在时区转换造成的缺口中。
    round_trip = candidate.astimezone(timezone.utc).astimezone(tz)

    if round_trip.replace(fold=0) == candidate.replace(fold=0):
        return candidate

    # Gap 的处理策略:向前寻找第一个有效的本地时间。
    probe = candidate
    for _ in range(3 * 60 * 60):
        probe += timedelta(seconds=1)
        round_trip = probe.astimezone(timezone.utc).astimezone(tz)
        if round_trip.replace(fold=0) == probe.replace(fold=0):
            return round_trip

    raise ValueError("无法解析本地时间")


def next_daily(
    now: datetime,
    at: time,
    tz: ZoneInfo,
) -> datetime:
    if now.tzinfo is None:
        raise ValueError("now 必须是 aware datetime")

    local_now = now.astimezone(tz)
    target_date = local_now.date()

    candidate = resolve_local(target_date, at, tz)
    if candidate <= local_now:
        candidate = resolve_local(
            target_date + timedelta(days=1),
            at,
            tz,
        )
    return candidate


if __name__ == "__main__":
    tz = ZoneInfo("America/Los_Angeles")
    now = datetime(2020, 10, 31, 12, tzinfo=tz)

    next_run = next_daily(now, time(1, 30), tz)
    print(next_run)
    print(next_run.astimezone(timezone.utc))

这段代码的关键不是“找出一个能运行的日期加法”,而是把 Gap 和 Fold 转成了显式策略。生产系统还应把策略作为配置保存,否则系统升级或迁移后无法解释历史调度结果。


四、Python 标准库中的调度器

Python 标准库没有提供 cron 服务或持久化分布式调度器,但提供了 sched.scheduler。它是一个通用事件调度器,默认使用 time.monotonic 作为时间函数、使用 time.sleep 作为延迟函数。事件按照时间和优先级执行;同一时间点的事件中,优先级数字越小越先执行。(docs.python.org)

1. 一次性任务

import sched
import time
from datetime import datetime

scheduler = sched.scheduler(
    timefunc=time.monotonic,
    delayfunc=time.sleep,
)


def notify(message: str) -> None:
    print(datetime.now().isoformat(), message)


scheduler.enter(
    delay=2,
    priority=1,
    action=notify,
    argument=("两秒后执行",),
)

scheduler.run()

执行过程是:

  1. enter() 将事件放入优先队列;
  2. run() 读取最早事件;
  3. 使用单调时钟计算还需等待多久;
  4. 调用 delayfunc()
  5. 到期后执行 action(*argument, **kwargs)
  6. 队列为空后返回。

如果任务函数抛出异常,调度器会保持队列状态并传播异常;出错事件不会在下一次 run() 中自动再次执行。若事件执行时间超过下一个事件的计划时间,调度器会落后,但不会自动丢弃事件。(docs.python.org)

这意味着:

scheduler.enter(0, 1, slow_task)
scheduler.enter(1, 1, quick_task)
scheduler.run()

如果 slow_task() 执行 10 秒,quick_task() 会在大约 10 秒后立即执行,而不是被自动取消。

2. 用 sched 实现固定频率重复任务

import sched
import time

scheduler = sched.scheduler(time.monotonic, time.sleep)

INTERVAL = 5.0
count = 0


def run_periodic(deadline: float) -> None:
    global count

    count += 1
    print(f"run={count}, lateness={time.monotonic() - deadline:.3f}s")

    if count >= 3:
        return

    # 根据原定目标时间推进,而不是根据本次完成时间推进。
    next_deadline = deadline + INTERVAL
    scheduler.enterabs(
        next_deadline,
        priority=1,
        action=run_periodic,
        argument=(next_deadline,),
    )


first_deadline = time.monotonic() + INTERVAL
scheduler.enterabs(
    first_deadline,
    priority=1,
    action=run_periodic,
    argument=(first_deadline,),
)
scheduler.run()

这里使用:

dn+1=dn+Δd_{n+1}=d_n+\Delta

而不是:

dn+1=nowfinish+Δd_{n+1}=\text{now}_{\text{finish}}+\Delta

前者是固定频率,后者是固定延迟。

3. 固定频率和固定延迟的反例

假设任务每次耗时 3 秒,配置间隔为 10 秒。

固定频率:

计划:00、10、20、30
实际:00、13、23、33

第一次任务从 00 执行到 03,因此目标 10 秒时仍然可执行,实际开始时间可能是 10 或 13,取决于调度器是否被其他任务阻塞。

固定延迟:

开始:00
执行:00~03
等待:03~13
开始:13

下一次开始:26

固定延迟的两次开始之间至少包含“任务耗时 + 间隔”;固定频率则努力追赶预先确定的时间轴。

选择哪一种取决于业务:

  • 采样、心跳通常关心计划时间,应采用固定频率;
  • “上一次处理完后休息 10 分钟”应采用固定延迟;
  • 不能无限追赶的任务需要定义错过策略。

五、重复执行:调度频率不等于执行语义

“每分钟执行一次”至少有三种不同解释。

1. 允许重叠

每个到期点都创建一个执行实例:

00:00 任务 A 开始
00:01 任务 B 开始
00:02 任务 C 开始

如果任务运行 90 秒,A、B 会重叠。它适合任务相互独立且下游能够承受并发的场景。

2. 不允许重叠,跳过忙碌时间点

如果 A 仍未完成,00:01 的 B 被跳过:

00:00 A 开始
00:01 B 跳过
00:02 C 开始或继续跳过

这适合只关心“当前最新状态”的刷新任务,但不适合账务、结算或消息处理,因为跳过可能意味着数据永久遗漏。

3. 不允许重叠,排队等待

每一个时间点都形成一个待执行实例:

00:00 A 执行
00:01 B 排队
00:02 C 排队

任务恢复后会依次处理。它不会丢失计划实例,但可能产生越来越大的延迟。

因此,重复任务至少需要定义:

@dataclass(frozen=True)
class RepeatPolicy:
    interval_seconds: int
    overlap: str       # "allow"、"skip"、"queue"
    misfire: str       # "run_once"、"catch_up"、"skip"

其中:

  • overlap 控制任务尚未完成时是否允许新实例进入;
  • misfire 控制调度器停机或阻塞后,如何处理已经错过的时间点。

4. 错过执行点:misfire

假设任务每分钟执行一次,调度器停机 10 分钟后恢复:

计划点:10:00、10:01、...、10:10
恢复时间:10:10:30

常见补偿语义有:

只执行一次

恢复后执行一次,代表“尽快完成当前周期”

适合刷新缓存。

追赶所有实例

依次执行 10:00 到 10:10 的 11 个实例

适合必须逐周期处理的业务,但要控制积压。

跳过过期实例

只执行下一个未来计划点

适合过期后已无业务价值的报表或临时通知。

错误的实现通常是:

while True:
    do_work()
    time.sleep(60)

它同时具有三个问题:

  1. 任务耗时会累积到下一次间隔;
  2. 进程退出后没有历史计划;
  3. 任务失败时没有失败记录和补偿依据。

六、调度器的生命周期和故障路径

一个内存调度器通常经历以下状态:

stateDiagram-v2
    [*] --> Scheduled: 创建任务
    Scheduled --> Due: 到达计划时间
    Due --> Claimed: 获得执行权
    Due --> Skipped: 被策略跳过
    Claimed --> Running: 开始执行
    Running --> Succeeded: 成功提交结果
    Running --> Failed: 抛出异常
    Failed --> Retryable: 可重试
    Failed --> Dead: 超过重试次数
    Retryable --> Scheduled: 重新安排
    Running --> Unknown: 进程崩溃
    Unknown --> Scheduled: 超时恢复

最危险的状态是 Unknown:系统无法确定任务到底执行到了哪一步。

例如:

1. 调用支付接口
2. 支付服务已经扣款
3. 本地进程在写入“成功”之前崩溃
4. 重启后再次调用支付接口

如果支付接口不幂等,就可能重复扣款。

因此,任务系统不能只记录:

last_run = "成功"

而应为每个计划实例记录独立状态:

job_id
scheduled_for
status
attempt
lease_until
started_at
finished_at
result
error

scheduled_for 是计划实例的身份组成部分。例如同一个任务的:

daily-report / 2026-09-01
daily-report / 2026-09-02

应被视为两个不同实例。


七、锁:防止并发进入,但不能保证业务只执行一次

1. 线程锁只解决进程内竞争

import threading

lock = threading.Lock()
balance = 0


def add_money(amount: int) -> None:
    global balance

    with lock:
        balance += amount

with lock 等价于获取锁、执行代码、无论是否异常都释放锁。Python 3.14 的 threading.Lock 是一个锁类,锁的基本状态只有“已锁定”和“未锁定”;等待中的线程谁先获得锁没有规范保证。(docs.python.org)

锁保护的是共享内存中的临界区:

线程 A:读取 balance
线程 A:加 1
线程 A:写回 balance

如果没有锁,线程 B 可能在 A 写回之前读取旧值。

但线程锁不能防止:

  • 两个进程各自持有一把不同的锁;
  • 两台机器同时执行;
  • 进程崩溃后锁状态丢失;
  • 外部 HTTP 请求已成功但本地记录未提交。

2. GIL 不是业务锁

CPython 默认构建中的 GIL 限制同一时刻执行 Python 字节码的线程数,但它不构成业务级事务锁,也不能把多个语句自动变成原子操作。Python 3.14 还支持自由线程构建,是否启用 GIL 不能作为程序正确性的依据。(docs.python.org)

错误思路:

# “有 GIL,所以不会冲突”
if key not in cache:
    cache[key] = build_value()

正确思路是显式同步:

with lock:
    if key not in cache:
        cache[key] = build_value()

3. 进程间锁需要外部协调者

多个 Worker 之间通常需要:

  • 数据库行锁;
  • 数据库唯一约束;
  • Redis 等外部存储的租约;
  • 文件锁;
  • 专门的任务队列协调机制。

核心目标不是“永远只有一个进程看到任务”,而是:

同一计划实例在同一时刻最多一个有效执行租约\text{同一计划实例在同一时刻最多一个有效执行租约}

这里的“有效”很重要。永久锁会在持有者崩溃后阻塞任务;因此生产系统通常使用带过期时间的租约:

未领取
  └─ Worker A 领取,lease_until = 10:05
       ├─ A 正常续租
       ├─ A 成功完成并标记 succeeded
       └─ A 崩溃,10:05 后允许其他 Worker 重新领取

租约过期并不证明 A 没有执行成功。A 可能只是网络隔离,或者完成后尚未来得及提交状态。因此,租约只能限制并发窗口,不能单独保证 exactly-once。


八、持久化领取:用数据库唯一约束消除重复创建

下面用 SQLite 演示一个最小的持久化任务表。SQLite 适合本地单机或教学示例,不应直接推导出它能替代高并发分布式队列。

CREATE TABLE IF NOT EXISTS job_run (
    job_id       TEXT NOT NULL,
    scheduled_for TEXT NOT NULL,
    status       TEXT NOT NULL,
    attempt      INTEGER NOT NULL DEFAULT 0,
    lease_until  TEXT,
    started_at   TEXT,
    finished_at  TEXT,
    error        TEXT,
    PRIMARY KEY (job_id, scheduled_for)
);

主键 (job_id, scheduled_for) 表示:

同一个任务 + 同一个计划时刻 = 一个计划实例

两个 Worker 同时尝试创建同一实例时,数据库只允许一个插入成功。

from datetime import datetime, timezone, timedelta
import sqlite3


def utc_now() -> datetime:
    return datetime.now(timezone.utc)


def ensure_run(
    conn: sqlite3.Connection,
    job_id: str,
    scheduled_for: datetime,
) -> bool:
    """
    确保计划实例存在。
    返回 True 表示本次插入成功,False 表示实例原本已经存在。
    """
    cursor = conn.execute(
        """
        INSERT OR IGNORE INTO job_run
            (job_id, scheduled_for, status)
        VALUES (?, ?, 'pending')
        """,
        (job_id, scheduled_for.isoformat()),
    )
    conn.commit()
    return cursor.rowcount == 1

INSERT OR IGNORE 只解决“不要创建两个相同实例”,还没有解决“谁可以执行”。领取动作仍然需要原子条件:

def claim_run(
    conn: sqlite3.Connection,
    job_id: str,
    scheduled_for: datetime,
    lease_seconds: int = 300,
) -> bool:
    now = utc_now()
    lease_until = now + timedelta(seconds=lease_seconds)

    cursor = conn.execute(
        """
        UPDATE job_run
           SET status = 'running',
               attempt = attempt + 1,
               started_at = ?,
               lease_until = ?
         WHERE job_id = ?
           AND scheduled_for = ?
           AND (
               status = 'pending'
               OR (
                   status = 'running'
                   AND lease_until IS NOT NULL
                   AND lease_until < ?
               )
           )
        """,
        (
            now.isoformat(),
            lease_until.isoformat(),
            job_id,
            scheduled_for.isoformat(),
            now.isoformat(),
        ),
    )
    conn.commit()
    return cursor.rowcount == 1

这个条件的含义是:

  • pending 实例可以被领取;
  • running 实例只有在租约过期后才能被重新领取;
  • succeededdead 永远不能被再次领取;
  • rowcount == 1 表示当前 Worker 获得了执行权;
  • rowcount == 0 表示其他 Worker 已经领取,或任务已完成。

生产数据库中应使用真正的事务隔离和适合该数据库的行级锁语义。SQLite 的写并发能力、锁粒度和部署方式都有边界,示例中的更新语句不能直接当作所有数据库的性能方案。


九、重复执行不可避免:必须设计幂等性

幂等性不是“函数只调用一次”,而是:

f(f(x))=f(x)f(f(x)) = f(x)

对任务系统来说,更实用的定义是:

对同一个业务操作重复提交,最终业务状态与提交一次相同。

例如,将订单状态设置为 paid 通常比“订单金额加 100”更容易做成幂等:

UPDATE orders
   SET status = 'paid'
 WHERE order_id = ?
   AND status <> 'paid';

而下面的操作天然容易重复产生副作用:

UPDATE account
   SET balance = balance + 100
 WHERE account_id = ?;

要让它可重试,需要引入业务操作 ID:

CREATE TABLE payment_operation (
    operation_id TEXT PRIMARY KEY,
    order_id     TEXT NOT NULL,
    created_at    TEXT NOT NULL
);

处理流程可以是:

1. 生成稳定的 operation_id
2. 尝试插入 payment_operation
3. 如果唯一键冲突,说明该操作已经处理过
4. 如果插入成功,在同一事务中更新业务表
5. 提交事务

伪代码:

def apply_payment(conn, operation_id: str, order_id: str) -> bool:
    with conn:
        cursor = conn.execute(
            """
            INSERT OR IGNORE INTO payment_operation
                (operation_id, order_id, created_at)
            VALUES (?, ?, datetime('now'))
            """,
            (operation_id, order_id),
        )

        if cursor.rowcount == 0:
            # 之前已经成功或正在由同一业务操作处理
            return False

        conn.execute(
            """
            UPDATE orders
               SET status = 'paid'
             WHERE order_id = ?
            """,
            (order_id,),
        )
        return True

但是,数据库幂等键只能保护数据库中的副作用。如果任务还调用外部支付服务,则必须使用外部服务提供的幂等键,或者采用 Outbox、状态查询、人工对账等补偿机制。


十、失败、重试和补偿不是同一个概念

1. 重试

重试是对一个仍然可能成功的执行再次发起相同操作:

第一次:网络超时
第二次:重新调用

适合:

  • 临时网络错误;
  • 服务暂时不可用;
  • 数据库连接断开;
  • 限流后等待重试。

不适合:

  • 参数校验失败;
  • 权限不足;
  • 业务规则明确拒绝;
  • 目标对象不存在且不会自动创建。

2. 补偿

补偿是对已经发生或可能发生的业务结果采取修复动作:

扣款成功,但发货记录创建失败
→ 查询扣款状态
→ 创建发货记录,或执行退款

重试的对象通常是“原操作”,补偿的对象是“状态不一致”。

3. 指数退避

重试等待时间可以定义为:

dk=min(dmax,d02k)+jd_k = \min(d_{\max}, d_0 \cdot 2^k) + j

其中:

  • kk 是已经失败的次数;
  • d0d_0 是初始延迟;
  • dmaxd_{\max} 是最大延迟;
  • jj 是随机抖动,用于避免多个 Worker 同时重试。

例如:

import random


def retry_delay(attempt: int) -> float:
    base = min(300.0, 2.0 ** attempt)
    jitter = random.uniform(0, 0.5)
    return base + jitter

如果 100 个任务同时在 60 秒后重试,它们会再次形成流量尖峰;抖动的作用是把重试分散到一个时间窗口内。

4. 重试上限和死信状态

任务不能无限重试。一个明确的状态机可以是:

pending
  → running
  → succeeded

running
  → retryable_failed
  → dead

retryable_failed
  → running

dead
  → 人工检查
  → 手工重放

dead 不表示问题解决,而表示自动机制停止继续扩大副作用。死信记录至少应包含:

job_id
scheduled_for
attempt
last_error
last_started_at
last_finished_at
payload 或 payload 引用

十一、一个可运行的本地调度示例

下面的示例结合:

  • sched.scheduler
  • 单调时钟等待;
  • SQLite 持久化;
  • 计划实例唯一键;
  • 领取租约;
  • 幂等业务写入;
  • 失败记录。

它不提供跨机器调度,但展示了关键生命周期。

from __future__ import annotations

import sched
import sqlite3
import time
from datetime import datetime, timezone, timedelta
from pathlib import Path


DB_PATH = Path("jobs.db")
JOB_ID = "demo"
LEASE_SECONDS = 30


def utc_now() -> datetime:
    return datetime.now(timezone.utc)


def connect() -> sqlite3.Connection:
    conn = sqlite3.connect(DB_PATH)
    conn.row_factory = sqlite3.Row
    conn.execute("PRAGMA busy_timeout = 5000")
    conn.execute(
        """
        CREATE TABLE IF NOT EXISTS job_run (
            job_id TEXT NOT NULL,
            scheduled_for TEXT NOT NULL,
            status TEXT NOT NULL,
            attempt INTEGER NOT NULL DEFAULT 0,
            lease_until TEXT,
            started_at TEXT,
            finished_at TEXT,
            error TEXT,
            PRIMARY KEY (job_id, scheduled_for)
        )
        """
    )
    conn.execute(
        """
        CREATE TABLE IF NOT EXISTS result (
            operation_id TEXT PRIMARY KEY,
            value INTEGER NOT NULL,
            created_at TEXT NOT NULL
        )
        """
    )
    conn.commit()
    return conn


def create_run(conn: sqlite3.Connection, scheduled_for: datetime) -> None:
    conn.execute(
        """
        INSERT OR IGNORE INTO job_run
            (job_id, scheduled_for, status)
        VALUES (?, ?, 'pending')
        """,
        (JOB_ID, scheduled_for.isoformat()),
    )
    conn.commit()


def claim_run(
    conn: sqlite3.Connection,
    scheduled_for: datetime,
) -> bool:
    now = utc_now()
    lease_until = now + timedelta(seconds=LEASE_SECONDS)

    cursor = conn.execute(
        """
        UPDATE job_run
           SET status = 'running',
               attempt = attempt + 1,
               started_at = ?,
               lease_until = ?
         WHERE job_id = ?
           AND scheduled_for = ?
           AND (
               status = 'pending'
               OR (
                   status = 'running'
                   AND lease_until < ?
               )
           )
        """,
        (
            now.isoformat(),
            lease_until.isoformat(),
            JOB_ID,
            scheduled_for.isoformat(),
            now.isoformat(),
        ),
    )
    conn.commit()
    return cursor.rowcount == 1


def mark_success(
    conn: sqlite3.Connection,
    scheduled_for: datetime,
) -> None:
    conn.execute(
        """
        UPDATE job_run
           SET status = 'succeeded',
               finished_at = ?,
               lease_until = NULL,
               error = NULL
         WHERE job_id = ?
           AND scheduled_for = ?
           AND status = 'running'
        """,
        (utc_now().isoformat(), JOB_ID, scheduled_for.isoformat()),
    )
    conn.commit()


def mark_failure(
    conn: sqlite3.Connection,
    scheduled_for: datetime,
    error: str,
) -> None:
    conn.execute(
        """
        UPDATE job_run
           SET status = 'retryable_failed',
               finished_at = ?,
               lease_until = NULL,
               error = ?
         WHERE job_id = ?
           AND scheduled_for = ?
           AND status = 'running'
        """,
        (
            utc_now().isoformat(),
            error[:1000],
            JOB_ID,
            scheduled_for.isoformat(),
        ),
    )
    conn.commit()


def business_action(
    conn: sqlite3.Connection,
    scheduled_for: datetime,
) -> None:
    """
    用计划实例作为幂等键。
    即使这个函数被重复调用,也只会插入一条结果。
    """
    operation_id = f"{JOB_ID}:{scheduled_for.isoformat()}"

    conn.execute(
        """
        INSERT OR IGNORE INTO result
            (operation_id, value, created_at)
        VALUES (?, ?, ?)
        """,
        (operation_id, 1, utc_now().isoformat()),
    )
    conn.commit()


def run_once(scheduled_for: datetime) -> None:
    conn = connect()

    try:
        create_run(conn, scheduled_for)

        if not claim_run(conn, scheduled_for):
            print("本次实例未获得执行权")
            return

        print("开始执行:", scheduled_for.isoformat())
        business_action(conn, scheduled_for)
        mark_success(conn, scheduled_for)
        print("执行成功")

    except Exception as exc:
        mark_failure(conn, scheduled_for, repr(exc))
        print("执行失败:", repr(exc))

    finally:
        conn.close()


def main() -> None:
    scheduler = sched.scheduler(
        timefunc=time.monotonic,
        delayfunc=time.sleep,
    )

    now_wall = utc_now()
    now_mono = time.monotonic()

    # 业务计划时间使用 UTC aware datetime。
    # 调度器等待使用 monotonic 的相对时间。
    first_run = now_wall + timedelta(seconds=2)
    delay = (first_run - now_wall).total_seconds()

    scheduler.enter(
        delay=max(0.0, delay),
        priority=1,
        action=run_once,
        argument=(first_run,),
    )

    scheduler.run()

    conn = connect()
    rows = conn.execute(
        """
        SELECT job_id, scheduled_for, status, attempt
          FROM job_run
         ORDER BY scheduled_for
        """
    ).fetchall()

    for row in rows:
        print(dict(row))

    conn.close()


if __name__ == "__main__":
    main()

运行:

python scheduler_demo.py

可能输出:

开始执行: 2026-09-01T...
执行成功
{'job_id': 'demo', 'scheduled_for': '2026-09-01T...', 'status': 'succeeded', 'attempt': 1}

这个示例中的关键保证分别来自不同机制:

机制 解决的问题
time.monotonic() 等待期间系统时间调整不会让等待倒退
scheduled_for 给计划实例提供稳定身份
主键约束 防止同一实例被重复创建
claim_run() 限制多个 Worker 同时执行
lease_until 允许崩溃后的实例恢复
INSERT OR IGNORE 让业务写入具有幂等效果
statuserror 为重试和诊断保留依据

它仍然不能保证外部系统 exactly-once。例如:

business_action()
数据库提交成功
mark_success() 尚未执行
进程崩溃

重启后,任务可能被认为未成功而重新领取。此时只能依靠业务幂等键或状态查询避免重复副作用。


十二、异步任务和 threading.Timer 的边界

如果任务主体是异步 I/O,不应在事件循环中直接调用阻塞式 time.sleep()

# 错误:会阻塞整个事件循环
async def job():
    time.sleep(10)

应使用:

import asyncio


async def job():
    await asyncio.sleep(10)

asyncio 的任务是事件循环中的协作式并发单元。asyncio.create_task() 可以并发调度协程,TaskGroup 提供结构化的任务生命周期管理。(docs.python.org)

一个简单的固定延迟异步循环是:

import asyncio


async def worker() -> None:
    for index in range(3):
        print("执行", index)
        await asyncio.sleep(2)


asyncio.run(worker())

如果要实现固定频率,应计算下一次目标时间,而不是让每次执行后的 sleep() 累积任务耗时:

import asyncio


async def periodic(interval: float, count: int) -> None:
    loop = asyncio.get_running_loop()
    deadline = loop.time()

    for index in range(count):
        deadline += interval
        await asyncio.sleep(max(0.0, deadline - loop.time()))
        print("执行", index)


asyncio.run(periodic(2.0, 3))

loop.time() 是事件循环使用的单调时间概念,适合计算相对延迟;业务日期仍应由 aware datetime 表达。

threading.Timer 适合简单的一次性线程延迟:

from threading import Timer

timer = Timer(5.0, print, args=("执行",))
timer.start()

它不是持久化调度器,也不自动解决:

  • 进程重启后的恢复;
  • 多进程重复执行;
  • 时区和 DST;
  • 失败重试;
  • 任务幂等;
  • 任务执行超时后的租约回收。

十三、调度器、任务队列和 Worker 的职责边界

在更大的系统中,通常会把“决定何时产生任务”和“执行任务”拆开:

flowchart LR
    S[Scheduler] -->|发布任务消息| B[Broker]
    B --> W1[Worker 1]
    B --> W2[Worker 2]
    W1 --> DB[(Database)]
    W2 --> DB
    W1 --> E[External Service]
    W2 --> E

调度器负责:

某个计划实例现在到期
→ 发送一条任务消息

Broker 负责:

暂存消息
→ 交给 Worker
→ 在一定条件下重新投递

Worker 负责:

获取消息
→ 执行业务函数
→ 成功确认 ACK,或失败重试

Celery 等任务队列通常还会涉及:

  • Broker;
  • Worker;
  • ACK;
  • 重试;
  • 定时任务;
  • 幂等性。

这些组件可以提高吞吐和故障恢复能力,但不会消除重复执行问题。消息可能在 Worker 已经产生外部副作用后、发送 ACK 前丢失连接,于是 Broker 再次投递;这正是为什么任务函数必须设计成幂等或可检测重复。

换句话说:

调度器解决“何时产生任务”
任务队列解决“如何传递任务”
Worker 解决“在哪里执行任务”
数据库和业务协议解决“重复执行后的正确性”

把所有问题都归因于某个调度框架,通常会导致错误的可靠性预期。


十四、时间戳、序列化和数据库字段

建议区分以下字段:

scheduled_for:计划实例的业务时刻
created_at:记录被创建的时刻
started_at:某次尝试开始的时刻
finished_at:某次尝试结束的时刻
lease_until:执行租约过期时刻

这些字段都应明确时区。跨系统存储时,可以统一使用 UTC:

from datetime import datetime, timezone

value = datetime.now(timezone.utc).isoformat()
print(value)

读取后不要把字符串直接当作本地时间:

from datetime import datetime

parsed = datetime.fromisoformat(value)
local = parsed.astimezone()

如果序列化的是本地日历规则,还需要额外存储:

{
  "local_time": "09:00:00",
  "timezone": "Asia/Shanghai",
  "fold_policy": "first",
  "gap_policy": "shift_forward"
}

仅保存:

{
  "timestamp": 1788224400
}

只能恢复某一个时刻,无法恢复原始调度规则。


十五、诊断定时任务:先判断是哪一种延迟

监控中至少要区分三个时间:

调度延迟=tstarttscheduled\text{调度延迟} = t_{\text{start}} - t_{\text{scheduled}}

执行时长=tfinishtstart\text{执行时长} = t_{\text{finish}} - t_{\text{start}}

端到端延迟=tfinishtscheduled\text{端到端延迟} = t_{\text{finish}} - t_{\text{scheduled}}

例如:

scheduled_for = 10:00:00
started_at    = 10:00:08
finished_at   = 10:00:10

调度延迟 = 8 秒
执行时长 = 2 秒
端到端延迟 = 10 秒

这三者反映不同问题:

  • 调度延迟高:调度线程阻塞、Worker 不足、锁竞争或机器暂停;
  • 执行时长高:业务本身变慢;
  • 端到端延迟高:前两者之一,或消息传递和重试造成积压。

常见故障表现与原因:

表现 可能原因
每次执行逐渐变晚 使用了固定延迟,却误以为是固定频率
DST 当天少执行一次 未定义 Gap 策略
DST 当天执行两次 未定义 Fold 策略
多实例同时执行 只使用了进程内锁
重启后任务消失 计划只存储在内存
任务偶尔重复 执行成功与状态提交之间发生崩溃
重试后数据翻倍 业务操作不可幂等
调整系统时间后任务提前或延后 time.time() 实现了相对等待
任务停机后一次执行几十次 默认采用了无限 catch-up

诊断时应先打印并关联:

job_id
scheduled_for
attempt
worker_id
started_at
finished_at
lease_until
status
error

不要只记录“任务开始”和“任务结束”,否则无法判断任务是晚了、慢了、重试了,还是被多个 Worker 同时领取。


十六、如何选择实现方式

单进程、短任务、非关键业务

可以使用:

  • sched.scheduler
  • asyncio
  • threading.Timer
  • 一个前台进程中的循环。

但要接受进程重启会丢失内存状态。

单机、需要恢复

增加:

  • SQLite 或其他本地数据库;
  • 计划实例表;
  • 唯一键;
  • 任务状态;
  • 租约和超时恢复。

多进程或多机器

需要:

  • 外部持久化存储;
  • 分布式领取或任务队列;
  • 幂等键;
  • 重试和死信;
  • 监控与人工重放能力。

任务具有明显外部副作用

必须优先设计:

稳定的业务操作 ID
幂等接口或幂等表
状态查询
失败补偿
对账机制

锁只能缩小并发窗口,不能替代幂等;重试只能提高成功机会,不能替代补偿;调度器只能决定时间,不能保证业务结果。

一个可以作为设计起点的完整约束是:

正确性=明确的时间语义+可恢复的计划实例+受控的执行权+幂等的业务操作+可观测的失败状态+明确的补偿路径\text{正确性} = \text{明确的时间语义} + \text{可恢复的计划实例} + \text{受控的执行权} + \text{幂等的业务操作} + \text{可观测的失败状态} + \text{明确的补偿路径}

当任务只需要“过一段时间执行一次”时,单调时钟和简单调度器通常已经足够。当任务需要“在某个时区的某个本地时间执行”,就必须引入 aware datetime、IANA 时区和 DST 策略。当任务还涉及数据库、消息队列或外部服务时,问题的核心便从“如何定时”转变为“如何在不确定的执行次数下保持业务状态正确”。


系列导航与关联阅读

官方资料

本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。