本节目标:掌握定时与延时任务的调度方式,理解「至少一次」投递为什么要求幂等,能设计出可靠的幂等键与状态机,并搭建死信队列与人工重放的闭环。
适用版本:Python 3.12+(实测 3.14.6);fakeredis 2.39.0
8.3 定时任务、幂等与死信处理
8.2 节我们做出了一个会重试、会进死信的队列。但它还留着两个没解决的问题:一是时间——怎么让任务「每天早上 9 点跑」或「3 分钟后跑」;二是重复——可见性超时带来的「至少一次」投递意味着同一个任务可能被执行两次,业务怎么保证不会重复扣款、重复发货。
这一节把这两块补齐,并给死信一个正式的出口:人工重放。这三件事看起来独立,其实指向同一个工程主题——让异步系统在不可靠的网络与进程里,表现得像可靠的一样。
8.3.1 定时任务的三种形态
「定时」其实有三种不同的时间语义,混在一起讲会乱:
| 形态 | 语义 | 例子 | 实现 |
|---|---|---|---|
| 周期性(cron) | 固定时刻反复触发 | 每天 9:00 发日报 | celery beat / arq cron / 系统 crontab |
| 周期性(interval) | 每隔固定时长触发一次 | 每 30 秒探活 | 循环 sleep 或调度器 |
| 一次性延时(delayed) | N 秒后执行一次 | 下单 15 分钟后未支付则关单 | 延时队列(ZSET) |
前两种靠调度器反复触发,第三种是一次性的、带具体执行时间点的任务。延时任务在业务里极其常见(超时关单、延迟重试、定时提醒),但大多数队列框架的「延时」能力都很弱,所以实战中往往自己用 Redis 的有序集合实现。
固定间隔的任务用 asyncio 写起来很轻,sleep 加循环就是一个最朴素的 interval 调度器:
import asyncio
async def run_every(interval: float, job, ticks: int):
"""每隔 interval 秒执行一次 job,共执行 ticks 次"""
for i in range(ticks):
await asyncio.sleep(interval) # 固定间隔
await job(i)
async def heartbeat(i):
print(f" [t={i}] 心跳任务执行")
async def main():
await run_every(0.05, heartbeat, 3)
asyncio.run(main())
[t=0] 心跳任务执行
[t=1] 心跳任务执行
[t=2] 心跳任务执行
三个 tick 各间隔 0.05 秒,总耗时约 0.15 秒。但这个朴素写法有个累积漂移的毛病:job 本身若耗时 0.02 秒,下一个 tick 实际是 0.05 + 0.02 之后才开始,越跑越偏。要求严格对齐时钟的场景(「每天 9:00」)应该算出下一个绝对触发时刻再 sleep(下一个时刻 - now),而不是简单地 sleep(interval)。
8.3.2 延时队列:用 ZSET 按时间排序
有序集合(Sorted Set)的 score 是浮点数,把它当成执行时间戳,就得到一个按「何时该执行」排序的延时队列。入队时 ZADD,取出时用 ZRANGEBYSCORE 捞所有 score 小于等于当前时刻的任务:
import json
import time
import fakeredis
redis = fakeredis.FakeRedis(decode_responses=True)
SCHED = "app:q:delayed" # score = 计划执行的时间戳(秒)
def enqueue_at(task_id: str, delay_s: float) -> float:
score = time.time() + delay_s
redis.zadd(SCHED, {json.dumps({"id": task_id}): score})
return score
def poll_due(now: float | None = None) -> list[dict]:
"""取出所有到点的任务(score <= now)"""
now = now if now is not None else time.time()
due = redis.zrangebyscore(SCHED, "-inf", now)
if due:
redis.zremrangebyscore(SCHED, "-inf", now)
return [json.loads(x) for x in due]
enqueue_at("t1", 0.05)
enqueue_at("t2", 0.15)
print("入队后待调度:", redis.zcard(SCHED))
time.sleep(0.08)
print("t=0.08 到点:", poll_due())
time.sleep(0.1)
print("t=0.18 到点:", poll_due())
print("剩余待调度:", redis.zcard(SCHED))
入队后待调度: 2
t=0.08 到点: [{'id': 't1'}]
t=0.18 到点: [{'id': 't2'}]
剩余待调度: 0
两个任务分别在 0.05s、0.15s 后到期,poll_due 在正确的时间点各捞出一个。这里有一个必须注意的原子性问题:ZRANGEBYSCORE 取出后、ZREMRANGEBYSCORE 删除前,如果 worker 崩溃,任务会被下次轮询重复取出。生产环境要用 Lua 脚本把「取出 + 删除」做成原子操作,或改成「取出 + 标记为处理中」。这正是「至少一次」语义的又一次现身——重复执行是常态,不是意外。
8.3.3 幂等键:让重复执行无害
既然任务会被跑两次,「保证只执行一次」这个目标本身就不现实。工程上的正确目标不是「不重复」,而是**「重复执行和一次执行的效果相同」**——这就是幂等。
幂等的第一把工具是幂等键(idempotency key):给每个业务操作分配一个全局唯一标识(订单号、请求 ID),执行前先抢占这个键,抢到才真正执行。Redis 的 SET NX 天生适合做这件事:
def run_once(idem_key: str, fn):
if not redis.set(f"app:idem:{idem_key}", "1", nx=True, ex=3600):
return "已处理过,跳过"
return fn()
calls = 0
def charge():
global calls
calls += 1
return f"扣款成功 #{calls}"
print("第一次:", run_once("order-1001-charge", charge))
print("重复提交:", run_once("order-1001-charge", charge))
print("实际扣款次数:", calls)
第一次: 扣款成功 #1
重复提交: 已处理过,跳过
实际扣款次数: 1
三次调用只真正扣款一次。但只用 SET NX 有一个缺点:重复请求拿到的是「跳过」,而不是原本的结果。对于「用户重复点了支付按钮」这种场景,调用方期望拿到和第一次一样的响应,而不是一句「已处理」。这时要把幂等键升级成状态机 + 结果缓存:
def process_payment(order_id: str, amount: int):
key = f"app:idem:pay:{order_id}"
if redis.set(key, json.dumps({"status": "processing"}), nx=True, ex=3600):
result = {"order": order_id, "charged": amount} # 真正扣款
redis.set(key, json.dumps({"status": "done", "result": result}), ex=3600)
return result, "首次执行"
state = json.loads(redis.get(key))
if state["status"] == "done":
return state["result"], "重复请求,返回缓存结果"
return None, "处理中,请稍后"
print(process_payment("O-1", 99))
print(process_payment("O-1", 99))
print(process_payment("O-1", 99))
({'order': 'O-1', 'charged': 99}, '首次执行')
({'order': 'O-1', 'charged': 99}, '重复请求,返回缓存结果')
({'order': 'O-1', 'charged': 99}, '重复请求,返回缓存结果')
这里藏着三个关键设计。第一,抢占和写入状态是两步,中间有极小的时间窗,所以 processing 状态是对外可见的——第二个请求看到它就知道「有人正在处理」,而不是误以为「没人在做」。第二,结果被缓存,重复请求返回完全相同的数据,对调用方来说这就是「只执行了一次」。第三,processing 状态需要超时兜底:如果第一次执行的进程在写 done 之前崩了,键会永远停在 processing。解决办法是给 processing 状态设一个较短的 TTL,或记录一个「处理开始时间」让后续请求能接管超时的任务。
8.3.4 幂等的边界:哪些操作天然幂等
不是所有操作都需要幂等键。区分「天然幂等」和「需要额外处理」能省下大量代码:
| 操作 | 是否天然幂等 | 说明 |
|---|---|---|
SET k v | 是 | 覆盖写,跑几次结果一样 |
DELETE k | 是 | 删不存在的键不报错 |
INCR k | 否 | 每跑一次加一次,必须去重 |
INSERT(无唯一约束) | 否 | 会插入多行,需唯一索引兜底 |
INSERT(有唯一约束) | 是 | 冲突即拒绝,天然去重 |
| 发送通知 / 邮件 | 否 | 会重复发,需幂等键或去重表 |
规律很清楚:「赋值型」操作天然幂等,「累加型 / 追加型」操作需要额外保护。最优雅的做法往往不是在应用层写幂等键,而是把幂等性下沉到数据库——给业务表加唯一约束(如 (order_id, action)),重复插入自然被拒绝,比应用层的 SET NX 更可靠,因为它和业务数据在同一个事务里。
8.3.5 死信队列与人工重放
重试耗尽的、参数永远不对的、依赖服务长期不可用的任务,继续留在主队列只会阻塞后面的任务。它们应该被挪进死信队列(DLQ)——一个「失败任务收容所」,等人工介入。
死信不是终点,它需要一条重放(replay)路径:运维看过失败原因、修好问题后,把死信重新投回主队列。重放的关键是重置重试计数,否则任务一进队列就又立刻进死信:
DLQ = "app:q:default:dead"
QUEUE = "app:q:default"
def list_dead(limit: int = 10) -> list[dict]:
return [json.loads(x) for x in redis.lrange(DLQ, 0, limit - 1)]
def replay(task_id: str) -> bool:
"""把指定死信重新投回主队列,并从死信移除"""
for raw in redis.lrange(DLQ, 0, -1):
d = json.loads(raw)
if d["task"]["id"] == task_id:
task = d["task"]
task["attempts"] = 0 # 重置重试计数
redis.lrem(DLQ, 1, raw)
redis.lpush(QUEUE, json.dumps(task, ensure_ascii=False))
return True
return False
for tid, reason in [("a1", "SMTP 拒收"), ("b2", "下游 503")]:
redis.lpush(DLQ, json.dumps({"task": {"id": tid, "name": "send_email", "attempts": 3},
"reason": reason}, ensure_ascii=False))
print("死信列表:")
for d in list_dead():
print(" ", d["task"]["id"], "->", d["reason"])
print("重放 a1:", replay("a1"))
print("重放后 主队列:", redis.llen(QUEUE), " 死信剩余:", redis.llen(DLQ))
死信列表:
b2 -> 下游 503
a1 -> SMTP 拒收
重放 a1: True
重放后 主队列: 1 死信剩余: 1
replay("a1") 把任务从死信挪回主队列,attempts 归零,死信从 2 条减到 1 条。生产环境的死信处理通常还有三件事:给死信设上限(超过阈值就丢弃最老的,防止 DLQ 无限膨胀)、按原因分类(SMTP 拒收这种是永久失败,重放也没用;下游 503才是可恢复的)、重放要审计(谁、何时、重放了多少条,否则重复扣款这类事故无从追溯)。
8.3.6 可观测性:让失败被看见
死信队列最大的风险不是「任务进了死信」,而是没人知道任务进了死信。异步任务的失败发生在后台,没有用户投诉来提醒你,所以必须靠指标和告警主动发现。
至少要盯四个指标:
| 指标 | 含义 | 告警条件 |
|---|---|---|
| 队列深度 | 主队列积压的任务数 | 持续增长或超阈值 |
| 处理延迟 | 任务入队到执行的耗时 | P99 超 SLA |
| 重试率 | 重试次数 / 总执行次数 | 突增说明下游异常 |
| 死信增长 | 单位时间进入 DLQ 的数量 | 任何非零都值得看 |
队列深度持续增长通常意味着消费能力不足(worker 太少或任务太慢);死信突增则意味着下游故障或代码 bug。两类问题的处理方向完全不同,所以指标要分开看。这些指标接到 3.3 节讲过的监控体系里,就能在用户察觉之前发现问题。
延伸阅读:Celery 的 beat 定时、结果后端与死信策略可参考专题 Python Celery 任务队列 。
小结
- 定时任务分 cron(固定时刻)、interval(固定间隔)、delayed(一次性延时)三类;延时任务用 Redis 的 ZSET 以时间戳为 score 实现。
- ZSET 延时队列的「取出 + 删除」不是原子操作,worker 崩溃会导致重复执行——「至少一次」是异步系统的常态。
- 幂等的目标不是「不重复执行」,而是「重复执行与一次执行效果相同」;用
SET NX抢占幂等键是最小实现。 - 需要返回一致结果时,把幂等键升级为「状态机 + 结果缓存」,并给
processing状态设超时兜底。 - 赋值型操作天然幂等,累加/追加型必须保护;最可靠的做法是把幂等下沉到数据库的唯一约束。
- 死信队列要有重放路径(重放时重置重试计数)、上限保护、按原因分类和审计;更要盯住队列深度、重试率、死信增长等指标主动告警。
到这里,8.1 到 8.3 把「缓存 + 队列 + 定时 + 幂等 + 死信」这条异步链路串了起来。下一章我们转向另一个绕不开的工程主题——认证与授权:用户是谁(认证)、他能做什么(授权),以及多租户下如何隔离数据。
阅读导航:上一节:任务队列与重试 · 下一节:认证与会话:JWT / OAuth2 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。