定义任务¶
三种写法(能力等价)¶
from taskmq import Task
class EmailTask(Task):
queue = "email"
retry_policy = Retry(max_attempts=5, backoff="exp") # 注意策略叫 retry_policy
def __init__(self, app): # 每进程一次:连接池、客户端都放这里
super().__init__(app)
self.client = SmtpClient()
def before_start(self, ctx):
self.client.acquire()
def run(self, to: str, subject: str) -> str:
self.request.update_meta(stage="sending") # 进度上报(写进 job.meta)
return self.client.send(to, subject)
def on_failure(self, ctx, exc):
alert(f"{ctx.id} failed: {exc!r}")
def after_return(self, ctx, state, result=None, exc=None):
self.client.release()
app.register(EmailTask)
约定:
self是每进程一个的实例,只放进程级资源;请求级状态放self.request(就是ctx);- 钩子顺序:
before_start → run → on_success / on_retry / on_failure → after_return; - 钩子异常不改投递语义(只记
hook_error);例外是before_start,它抛异常即任务失败。
任务级选项¶
写在装饰器参数里,或写成 Task 子类的类属性(装饰器参数优先):
| 选项 | 默认 | 说明 |
|---|---|---|
name |
模块.函数名 | 任务名(跨进程/跨版本必须稳定) |
queue |
Config.default_queue |
投到哪个队列(worker 必须订阅它) |
priority |
Config.default_priority |
-9..9,越大越优先 |
retry / retry_policy |
无 | Retry(...) 策略,见下 |
timeout |
无 | 软超时:入队时算 deadline,超时抛 TaskTimeout(协作式) |
hard_timeout |
无 | 硬超时:只有 processes 池能强杀子进程 |
rate_limit |
无 | "100/m" / "10/s" / "1000/h",worker 内令牌桶,超限 defer |
concurrency_key |
无 | 模板(如 "user:{user}"),同 key 在集群内串行 |
expires |
无 | 秒;超过这个时间还没开始执行就作废(EXPIRED) |
ack |
on_success |
on_receipt(收到即 ack)/ on_success / on_completion |
max_deliveries |
5(或 Retry.max_attempts) |
毒丸保护:投递次数上限,超过进 DLQ |
提交选项¶
send_email.delay("a@b.com", "hi") # 只传任务参数
send_email.apply_async( # 带投递选项
("a@b.com", "hi"),
queue="email",
priority=Priority.HIGH,
delay=30, # 相对延迟(秒)
eta=time.time() + 30, # 绝对时间
expires=3600, # 秒;入队时算 expires_at
key="welcome:a@b.com", # 幂等键
timeout=120, # 覆盖任务的软超时
)
delay()只接任务参数(有静态类型检查,IDE 能补全);apply_async()才接队列/优先级/可见性等投递选项;- 两者都返回
TaskHandle:.id、.state、.info、.successful()、.wait(timeout)、.get(timeout)(失败抛RemoteError,带远端 traceback 文本)、.forget()。
重试策略¶
from taskmq import Retry
Retry(
max_attempts=5, # 总尝试次数
backoff="exp", # "fixed" | "linear" | "exp"
base=1.0, # 基准延迟(秒)
factor=2.0, # exp/linear 的倍数
max_delay=600.0, # 延迟上限
jitter=True, # ±50% 抖动,避免惊群
retry_on=(Exception,), # 只有这些异常才重试
)
- 重试由 worker 侧执行:失败后按退避把消息放回
delayed,deliveries继续累加; - 想不重试直接进 DLQ:抛
Reject("原因");想显式重试:raise ctx.retry(reason=..., delay=...); - 重试次数用尽 → job 变
FAILED,消息进 DLQ。
限流与按 key 串行¶
`app.task(queue="email", rate_limit="100/m")
def send_email(to: str) -> str: ...
`app.task(queue="report", concurrency_key="user:{user}")
def build_report(user: str) -> str: ...
两者都在还没真正执行时用 transport.defer() 放回队列,不消耗 deliveries,
所以不会被毒丸保护误送 DLQ。concurrency_key 用 transport 的命名租约做跨 worker 互斥,租约到期自动释放
(需要 transport 支持 supports_leases;不支持的后端会在启动时报错)。
进度上报与日志¶
`app.task(queue="etl")
def extract(path: str) -> int:
ctx = current_task() # 或 bind=True 的 self.request
ctx.update_meta(stage="reading", rows=0) # 写进 job.meta,status/info 可见
ctx.log.info("reading", path=path) # 结构化日志(带 job/task/attempt 字段)
...
handle.info 能看到 meta;taskmq status / 工作流状态页也读同一份数据。