taskmq 设计文档(v0.1 草案)¶
定位:一个「类 Celery」的 Python 任务队列框架。继承 Celery 的心智模型(task / queue / worker / beat), 但重新设计投递语义、配置模型和开发体验,干掉那些让人无语的隐式魔法。
状态:v0.2。§20 的决策已确认(仅第 7 条工作流优先级为 🟡 暂定)。Phase 0 实现进行中。
0. 一句话定位¶
零外部服务就能跑起来(不需要 Redis / RabbitMQ / 外部 DB)、投递语义可预测、配置显式、调试不用猜的 Python 分布式任务队列。
from taskmq import App, Config, Retry
app = App(Config(transport="sqlite:///./taskmq.db", concurrency=8))
@app.task(queue="email", retry=Retry(max_attempts=5, backoff="exp"))
def send_email(to: str, subject: str) -> str:
...
send_email.delay("a@b.com", "hi") # 生产者侧,不需要 broker
不需要先装 Redis、不需要先起 RabbitMQ、不需要先写 200 行配置。
1. 为什么要重写:Celery 的痛点 → 对策¶
| # | Celery 的坑 | 实际表现 | taskmq 对策 |
|---|---|---|---|
| 1 | 强依赖外部 broker | 本地开发和 CI 都要先起 Redis/RabbitMQ | 内置 SQLite transport(WAL + 原子 claim),生产可直接用;内存 transport 用于测试;Redis/AMQP 是可选增强而非前提 |
| 2 | 默认「早 ack」 | worker 被杀 / OOM / 部署重启 → 任务静默丢失 | 默认 at-least-once(成功后才 ack)+ 可见性租约 + 最大投递次数 + DLQ |
| 3 | prefetch 是乘法 | prefetch_multiplier * concurrency,长任务被一个 worker 全吞掉,其他 worker 饿死 |
prefetch 是绝对值,默认等于并发数,可 per-queue;语义写进文档 |
| 4 | 配置是懒加载代理 | app.conf 读的时候才知道值,拼错 key 不报错 |
单一 frozen dataclass,构造即校验,未知字段直接抛错;环境变量前缀显式映射 |
| 5 | autodiscover 隐式 import | 循环导入、配置未就绪就执行模块级代码、任务注册顺序依赖 import 顺序 | 显式 include=[...] / app.discover("pkg");导入期零副作用是硬约束 |
| 6 | 两套重试机制 | self.retry() 用异常控制流 + autoretry_for,语义重叠、容易写出死循环 |
单一 Retry 声明式策略 + ctx.retry() 手动逃生口,二者语义互补且明确 |
| 7 | canvas 过于复杂 | chain/group/chord 互相嵌套,chord 靠轮询 result backend,出问题没法调 | Phase 2 用原生 DAG 工作流替代,调度器直接感知依赖,不走轮询 |
| 8 | 默认允许 pickle | 安全风险和跨版本反序列化灾难 | 默认 msgspec(serializer="json" 可退回纯标准库),类型白名单注册,pickle 需要显式开启并打警告 |
| 9 | 结果语义混乱 | ignore_result / backend 组合爆炸,Redis 里结果永远堆积 |
结果后端是显式配置的一个对象,强制 TTL;不配置就只保留轻量状态,不存返回值 |
| 10 | beat 多实例重复触发 / 时区坑 | 部署两个副本就双倍执行;DST 切换时任务漂移 | SQLite/Redis lease 选主;zoneinfo 显式时区;misfire 策略显式声明 |
| 11 | signals 无类型、难测 | 字符串事件名 + 弱签名,IDE 无补全,测试要全局注册再清理 | 类型化 hooks + middleware,可直接实例化测试 |
| 12 | 测试难 | task_always_eager 是全局开关,monkeypatch 满天飞 |
MemoryTransport + eager=True + Worker.run_until_idle(),确定性、无全局状态 |
| 13 | 版本频繁破坏性变更 | 升级一次改一堆配置项 | 消息协议带 v,语义化版本 + 明确兼容窗口,配置字段新增不删除 |
2. 设计原则¶
- P1 零外部服务可运行:不需要 Redis / RabbitMQ / 外部数据库。core 的唯一第三方依赖是
msgspec(默认序列化器,决策 §20-5);显式配置serializer="json"时退化为纯标准库路径。 - P2 显式优于隐式:路由、并发、重试、序列化、超时全部显式声明;默认值写在文档里且可预测。
- P3 投递语义可证明:每条消息处于什么状态、谁持有、何时可见,都是可查询的数据库行,不是内存里的黑盒。
- P4 失败必须可见:失败进 DLQ、进事件流、进 CLI,不吞异常、不静默丢任务。
- P5 单进程开发体验:
taskmq dev一条命令起全套,本地零配置。 - P6 不做分布式事务:不承诺 exactly-once;用「at-least-once + 幂等键」把问题还给业务,但把工具给足。
- P7 小核心 + 明确插拔点:Transport / ResultBackend / Pool / Scheduler / Serializer 五个接口,其余不抽象。
- P8 类型友好:
py.typed,任务句柄泛型化,Result[T]可被 mypy/pyright 推断。
3. 非目标(这一版明确不做)¶
- 不默认提供 Celery 的 drop-in 兼容层:
from celery import ...只在显式安装可选包taskmq-celery(Phase 2,决策 §20-8)后可用,默认安装不遮蔽真实的celery。 - 不支持 Python < 3.10(决策 §20-2)。代码不得使用 3.11+ 专有语法/API(
asyncio.timeout、TaskGroup、ExceptionGroup、tomllib、typing.Self、StrEnum、datetime.UTC等),确有需要时收敛到taskmq/_compat.py。 - 不做 supervisor / 自动扩缩容 / K8s operator(交给 systemd、K8s、Nomad)。
- 不内置 Web UI(Phase 3 再议)。
- 不默认启用 pickle,不追求「任意对象都能进队列」。
4. 架构总览¶
flowchart LR
subgraph PROD["生产者进程"]
TASK["@app.task<br/>task.delay / task.submit"]
RES["Result 句柄 / 状态查询"]
end
subgraph WORKER["worker 进程"]
FETCH["Fetcher<br/>拉取 · 预取 · 背压"]
ROUTE["Router<br/>按队列分派"]
POOL["Pool<br/>solo · threads · processes · asyncio"]
ACK["Ack / Nack / Requeue / Retry"]
FETCH --> ROUTE --> POOL
POOL -.->|"TaskContext"| ACK
ACK -.-> FETCH
end
subgraph STORE["共享存储层"]
TRANS[("Transport<br/>Memory · SQLite · Redis · AMQP")]
BACK[("ResultBackend<br/>Memory · SQLite · Redis")]
EVENTS["EventSink<br/>stdout JSON · OTel · 自定义"]
end
BEAT["Scheduler beat<br/>lease 选主 · 多副本安全"]
TASK -->|"envelope JSON"| TRANS
TRANS -->|"reserve 投递"| FETCH
ACK -->|"ack / nack / dead_letter"| TRANS
POOL -->|"写状态 / 结果"| BACK
RES -->|"get / state"| BACK
BACK -.->|"state / result"| RES
POOL -->|"事件"| EVENTS
ACK -->|"事件"| EVENTS
BEAT -->|"到点入队"| TRANS
WORKER -.->|"heartbeat"| TRANS
代码分层(依赖只能向下,不能反向):
taskmq/
__init__.py # 公开 API:App, task, Config, Retry, Result, current_task
config.py # frozen dataclass 配置 + 校验 + 环境变量映射
app.py # App:注册表、transport/backend/pool 装配
task.py # @task 装饰器、TaskHandle、TaskContext、Retry
protocol.py # envelope 定义、版本、编解码、大小限制
transport/
base.py # Transport 抽象 + 不变量文档
memory.py
sqlite.py
redis.py # Phase 1
amqp.py # Phase 2
result/
base.py memory.py sqlite.py redis.py
worker/
runner.py # 主循环、生命周期、优雅退出
pool.py # solo / threads / processes / asyncio
fetcher.py # 拉取、预取、背压、per-queue 公平
hooks.py # 类型化钩子 + 中间件
scheduler/
beat.py # cron/interval 调度
lease.py # 选主
cli.py # taskmq worker / beat / dev / status / dlq / call / purge
testing.py # fixtures、run_until_idle、时间冻结、故障注入
5. 核心概念¶
| 概念 | 定义 |
|---|---|
| Task | 被 @app.task 装饰的可调用对象 + 它的静态策略(队列、重试、超时、序列化) |
| Job | 一次逻辑任务(一个 ULID),可能包含多次投递/尝试 |
| Delivery | Job 的一次投递尝试,有独立的 attempt 计数和租约 |
| Envelope | 在 transport 中流转的 JSON 消息体 |
| Queue | 逻辑队列(名字 + 优先级 + 容量 + TTL),不是 broker 的物理队列 |
| Worker | 一个进程,持有 N 个并发槽位,订阅若干队列,定期心跳 |
| Pool | Worker 内的执行模型:solo / threads / processes / asyncio |
| Transport | 消息的可靠投递层:enqueue / reserve / ack / nack / requeue / dead_letter |
| ResultBackend | 返回值与最终状态的持久化层(可选、显式) |
| DLQ | 死信队列:超过最大投递次数、不可重试、或过期失败的任务落在这里,可查询可重放 |
| TaskContext | 任务运行时的显式上下文对象(attempt、deadline、日志、ctx.retry()、ctx.publish()) |
6. 目标开发体验(先写用例,再写实现)¶
6.1 生产者¶
from taskmq import App, Config
app = App(Config(transport="sqlite:///./taskmq.db"))
@app.task(queue="email", timeout=30, retry="default")
def send_email(to: str, subject: str) -> str:
return f"sent:{to}"
h = send_email.delay("a@b.com", "hi") # 语法糖
h2 = send_email.submit( # 全功能入口
"a@b.com", "hi",
queue="email",
eta=None, # 延迟执行(秒或 datetime)
expires=3600, # 超过 TTL 未执行则丢弃并标记 expired
timeout=30,
priority=5,
key="order-42:email", # 幂等键:同 key 在窗口内只执行一次
)
print(h.id, h.state) # 立即返回,不阻塞
print(h.get(timeout=10)) # 需要结果时才阻塞
6.2 消费者(任务实现)¶
from taskmq import task, current_task, Retry
@app.task(queue="email", retry=Retry(max_attempts=5, backoff="exp", jitter=True))
def send_email(to: str, subject: str) -> str:
ctx = current_task() # contextvar,显式取用,测试里可直接构造
ctx.log.info("sending", to=to, attempt=ctx.attempt)
try:
return smtp_send(to, subject)
except TemporaryError as e:
raise ctx.retry(reason=str(e)) # 交由策略计算退避时间
6.3 CLI¶
taskmq worker -Q email,default -c 8 --pool threads
taskmq worker -Q heavy -c 4 --pool processes
taskmq beat --schedule taskmq_schedule.py
taskmq dev # worker + beat 同一进程,本地开发
taskmq status # worker 列表、队列深度、速率
taskmq dlq list --queue email
taskmq dlq replay --all --queue email
taskmq call myapp.tasks.send_email --args '["a@b.com","hi"]' # 同步触发一次,方便调试
taskmq purge --queue email
7. 公开 API 设计¶
7.1 App¶
app = App(
Config(
transport="sqlite:///./taskmq.db", # 或 "memory://" / "redis://..." / 自定义实例
result=None, # 不配则不存返回值,只保留状态
default_queue="default",
concurrency=4,
pool="threads",
prefetch=4, # 绝对值,不是倍数
serializer="json",
timezone="Asia/Shanghai",
),
include=["myapp.tasks"], # 显式导入任务模块
)
规则:
- 构造即校验:字段类型、URL scheme、pool 与任务类型的兼容性全部在 App() 时检查,未知配置项直接抛 ConfigError。
- 无全局单例:可以有多个 App(测试里很需要),current_app() 只在 task 执行上下文中有效。
- 无模块级副作用:import myapp.tasks 只做函数定义,不建立连接、不读环境变量。
7.2 任务定义参数¶
两种等价写法:@app.task(函数式糖)与 继承 Task 基类(类式,可 mixin、可覆写生命周期钩子);函数式任务用 bind=True 拿到任务实例 self,self.request 即类型化的 TaskContext。
钩子顺序、异常策略、实例生命周期与 Celery 对照见 tasks.md。
| 参数 | 默认 | 说明 |
|---|---|---|
name |
module.func |
跨进程唯一标识,必须稳定(改名 = 协议不兼容) |
queue |
default_queue |
静态路由;也可用 router 函数动态决定 |
priority |
0 |
静态默认优先级(数值越大越优先);submit(priority=…) 可覆盖;作用域与防饥饿见 priority.md |
retry |
None |
Retry 对象或注册表里的名字 |
timeout |
None |
软超时(协作式 deadline),见 §10.7 |
ack |
"on_success" |
"on_receipt" / "on_success" / "on_completion" |
expires |
None |
默认 TTL,超时未执行则丢弃 |
max_deliveries |
5 |
超过则进 DLQ,防止毒丸消息无限重投 |
rate_limit |
None |
如 "100/m",token bucket |
concurrency_key |
None |
同 key 的任务在集群内串行(如按 user_id 串行) |
serializer |
全局 | 允许 per-task 覆盖 |
store_result |
True |
若配了 result backend |
7.3 状态机¶
stateDiagram-v2
[*] --> QUEUED: submit
QUEUED --> RUNNING: reserve 投递
RUNNING --> SUCCEEDED: 执行成功
RUNNING --> RETRYING: 可重试异常
RETRYING --> QUEUED: 退避到期,重新可见
RUNNING --> FAILED: 不可重试 / 超过 max_deliveries
RUNNING --> REVOKED: revoke / 硬超时 kill
QUEUED --> EXPIRED: 超过 expires 未执行
SUCCEEDED --> [*]
FAILED --> [*]: 保留 traceback 并写入 DLQ
EXPIRED --> [*]
REVOKED --> [*]
- 状态变更先写 transport/backend 再产生副作用,保证崩溃后状态一致。
PENDING只用于「生产者还没提交成功」或「查询不存在的 id」,不参与正常工作流(消除 Celery 里 PENDING 的歧义)。
7.4 Result 句柄¶
h = send_email.delay("a@b.com", "hi")
h.id # ULID,按时间可排序,便于按时间范围排查
h.state # 单次读取,不阻塞
h.get(timeout=10) # 阻塞取结果;失败抛 RemoteError(含远端 traceback)
await h.aget(timeout=10) # asyncio 场景
h.wait(timeout=60, poll=0.2) # 只等状态不取结果
h.revoke(terminate=False) # 撤销(未执行则直接丢弃)
h.info # {"attempt": 2, "worker": "w-1", "runtime": 0.31, ...}
7.5 TaskContext¶
ctx.id # job id
ctx.attempt # 当前是第几次尝试(从 1 开始)
ctx.deliveries # 累计投递次数
ctx.queue
ctx.deadline # 绝对时间;协作式超时用
ctx.remaining() # 剩余秒数
ctx.log # 结构化 logger,自动带 job_id / task / attempt
ctx.publish(t, *a, **kw) # 在任务里派发子任务(不自己建连接)
ctx.retry(reason="") # 按策略重试,返回要抛的异常
ctx.heartbeat() # 长任务手动续租,避免被判定为孤儿
ctx.check_cancelled() # 协作式取消检查点
8. 消息协议(Envelope)¶
{
"v": 1,
"id": "01HQ8Z9K2M4T6V8X0Z2B4D6F8H",
"task": "myapp.tasks.send_email",
"args": ["a@b.com", "hi"],
"kwargs": {},
"queue": "email",
"priority": 0,
"eta": null,
"expires_at": null,
"deadline": null,
"attempt": 1,
"max_attempts": 5,
"ack": "on_success",
"key": "order-42:email",
"trace": { "traceparent": "00-...-...-01", "origin": "web-1" },
"enqueued_at": 1735000000.123,
"headers": {}
}
约束:
- v 是协议版本,反序列化时按版本分支;不认识的大版本直接拒绝并告警,不猜测。
- 消息体默认上限 256 KiB(可配),超限在生产者侧就抛错,不等到 broker 炸。
- args/kwargs 必须是 JSON 可编码;自定义类型需要 app.register_codec(Type, encode, decode)。
- 时间统一用 UTC 时间戳,展示层再转时区。
9. Transport 抽象¶
9.1 接口¶
class Transport(Protocol):
def enqueue(self, env: Envelope, *, queue: str | None = None, delay: float = 0.0,
priority: int | None = None) -> str: ... # 返回 job id(env.id)
def reserve(self, queues: list[str], *, worker_id: str, lease: float, limit: int) -> list[Delivery]: ...
def ack(self, d: Delivery) -> None: ...
def nack(self, d: Delivery, *, requeue: bool = True, delay: float = 0.0) -> None: ...
def dead_letter(self, d: Delivery, reason: str) -> None: ...
def extend_lease(self, d: Delivery, seconds: float) -> None: ...
# 插队原语(决策 §20-9,见 priority.md §3.3)
def peek_max_priority(self, queues: list[str]) -> int | None: ...
def yield_reservation(self, d: Delivery, *, delay: float = 0.0, max_yields: int = 100) -> bool: ...
def next_visible_at(self, queues: list[str]) -> float | None: ...
def set_state(self, job_id: str, state: str, **meta) -> JobRecord: ...
def get_state(self, job_id: str) -> JobRecord | None: ...
def queue_stats(self, queues: list[str] | None = None) -> list[QueueStat]: ...
def reap_expired_leases(self, now: float | None = None) -> int: ... # 孤儿回收
def reap_expired_jobs(self, now: float | None = None) -> int: ... # expires
def close(self) -> None: ...
Phase 0 已实现
base.py+memory.py(sqlite.py是下一步)。
9.2 必须成立的不变量¶
- 一条消息在任意时刻至多被一个 worker 持有(租约保证)。
- 持有者崩溃且租约到期后,消息必须重新可见(除非已 ack)。
ack之后消息不可再被任何 worker 看到。reserve返回的消息,其attempt/deliveries计数已自增。- 队列顺序:同优先级下 FIFO;有
priority时高优先级先出(不保证严格公平,除非显式开启)。 - 所有状态转换都是幂等的(重复 ack 不报错)。
reap_*可由任意进程并发调用而不产生重复投递。- 任何错误都不允许「假装成功」——宁可重投。
9.3 内置 SQLite Transport(重点:这是「零依赖」的关键)¶
PRAGMA journal_mode=WAL;
PRAGMA busy_timeout=5000;
PRAGMA synchronous=NORMAL;
CREATE TABLE messages (
id INTEGER PRIMARY KEY AUTOINCREMENT,
job_id TEXT NOT NULL,
queue TEXT NOT NULL,
task TEXT NOT NULL,
envelope BLOB NOT NULL,
priority INTEGER NOT NULL DEFAULT 0,
state TEXT NOT NULL DEFAULT 'queued', -- queued|reserved|acked|dead
visible_at REAL NOT NULL, -- eta/退避后的可见时间
expires_at REAL,
claimed_by TEXT,
claimed_at REAL,
lease_until REAL,
deliveries INTEGER NOT NULL DEFAULT 0,
last_error TEXT
);
CREATE INDEX idx_claim ON messages(state, queue, visible_at, priority DESC, id);
CREATE INDEX idx_job ON messages(job_id);
CREATE TABLE jobs ( -- 轻量状态机 + 可选结果
job_id TEXT PRIMARY KEY,
task TEXT NOT NULL,
state TEXT NOT NULL,
result BLOB,
error TEXT,
created_at REAL NOT NULL,
updated_at REAL NOT NULL,
expires_at REAL
);
CREATE TABLE workers ( -- 心跳,用于 status 与孤儿判定
id TEXT PRIMARY KEY, queues TEXT, pool TEXT, concurrency INTEGER,
started_at REAL, heartbeat_at REAL, meta BLOB
);
CREATE TABLE leases ( -- beat 选主 / 全局互斥
name TEXT PRIMARY KEY, owner TEXT, expires_at REAL
);
CREATE TABLE idempotency ( -- 幂等键
key TEXT PRIMARY KEY, job_id TEXT, created_at REAL, expires_at REAL
);
原子 claim(多 worker 并发安全):
BEGIN IMMEDIATE;
UPDATE messages
SET state='reserved', claimed_by=?, claimed_at=?, lease_until=?, deliveries=deliveries+1
WHERE id = (SELECT id FROM messages
WHERE state='queued'
AND queue IN (...)
AND visible_at <= ?
AND (expires_at IS NULL OR expires_at > ?)
ORDER BY priority DESC, id
LIMIT 1)
RETURNING *;
COMMIT;
实现要点:
- enqueue 用独立短事务,避免长事务阻塞 writer。
- 空闲时轮询退避:0.05s → 0.5s(可配上限),有积压时立刻回落到最小值。
- 吞吐预期:SSD + WAL 下 enqueue ≈ 3k–10k msg/s,claim ≈ 1k–3k msg/s(单文件单 writer)。
超过这个量级请上 Redis transport。文档必须写明这条边界,不做「什么都能扛」的承诺。
实测(本机 M 系列 SSD、单连接、小载荷、每消息 1 次 reserve + 1 次 ack):
enqueue ≈ 11.3k msg/s,claim+ack ≈ 4.6k msg/s;端到端吞吐由 worker 的 poll 间隔与任务本身决定
(验收里的 1000 任务 7.3s 主要是 run_until_idle 的 50ms 轮询等待,不是 SQL 瓶颈)。
- 多进程共享同一个 db 文件是支持的一等场景(worker 部署在同一台机器 / 共享盘)。
跨主机(决策 §20-6):一期只覆盖同机多进程与共享盘(NFS/SMB,需可靠的 POSIX 文件锁);
跨主机高吞吐请用 redis://(Phase 1)或 postgres://(Phase 2)。
对象存储(oss:// / s3://)不作为一期 transport:即便对象存储提供条件写类原语,每次 reserve 仍需
额外的租约对象 + 轮询/长轮询,延迟与请求成本都不适合放在热路径;定位为 Phase 2+ 的低频/归档队列候选。
9.4 其他 Transport¶
| Transport | 阶段 | 用途 |
|---|---|---|
memory:// |
Phase 0 | 单测、eager 模式、本地试跑 |
sqlite:// |
Phase 0 | 零依赖默认、单机生产 |
redis:// |
Phase 1 ✅ | 高吞吐;零依赖 RESP 客户端 + Lua 原子操作(见下) |
postgres:// |
Phase 2 | 已经有 PG 的团队:SKIP LOCKED + LISTEN/NOTIFY |
postgresql:// |
Phase 2 ✅ | 多机 + 强一致:FOR UPDATE SKIP LOCKED 原子 claim、ON CONFLICT CAS 租约、幂等键唯一约束;需要 taskmq-py[postgres](psycopg) |
amqp://(RabbitMQ) |
Phase 2 ✅ | 投递走 broker(x-max-priority + per-message TTL/DLX 延迟 + DLX 死信),状态走侧车 ?state=sqlite:///…(AMQP 无 KV/DAG 索引);跨队列不做全局严格优先(已声明降级);需要 taskmq-py[amqp] |
<自研>:// |
插件 | 第三方后端经注册表接入,核心零改动(design/plugins.md) |
Redis transport 实现要点(taskmq/transport/redis.py + taskmq/redis_client.py):
- 零第三方依赖:自带极简 RESP2 客户端(
socket+ 一把锁),EVAL走 Lua; - 键前缀隔离:URL 支持
redis://host:port/db?prefix=app1,所有操作只碰自己的前缀, 不会 FLUSH 别人的库(测试用随机前缀互不干扰); - 优先级:ready ZSET 的
score = -priority * 2**40 + seq。取件规则 = 先定「所有队头里的最高 优先级档位」,档内再按served/weight最小者取(平手按队列名)——即「跨档严格优先 + 档内 平级加权轮询」,与 memory/sqlite 一致(redis-cluster.md); - 可见性统一:延迟/退避/让位/重投都进 delayed ZSET(
score = visible_at),reserve 先 promote 到期的; 过期在 promote/claim 时判定,因此reap_expired_jobs()在 Redis 上是 no-op(返回 0); - 原子性:claim、ack/nack/dead/defer/extend/yield、孤儿回收、命名租约默认用 Lua 一次往返; 没有 Lua 也能跑(见下);
- 无 Lua 回退:启动时探测
EVAL。很多托管 Redis 会禁用/重命名 EVAL(安全合规), 这时自动切到WATCH/MULTI/EXEC乐观事务——语义与 Lua 路径一致,代价只是争抢时多几个往返 (乐观重试)。开关:?lua=off强制回退、?lua=on要求必须有 Lua(没有就启动即报错)、不写则自动探测; 当前模式在transport.lua_enabled上可查。测试两种模式各跑一遍全部用例,另有自动探测与人工对拍用例。 - Cluster(
?cluster=1):键按逻辑队列打 hash tag({prefix}{q}:ready),单队列的 promote/claim 走显式 KEYS 的 Lua(一次往返、仍是原子的);跨队列取件跨 slot,没有单次原子可言, 降级为「每队列各取一次 + Python 侧按 band/权重选」,并在cluster_limitations里声明 (一致性套件据此跳过global_priority)。客户端是自带的RedisClusterClient: CRC16 算 slot +CLUSTER SLOTS拓扑 +MOVED/ASK跟随,同槽校验在发命令前就拦下 (Redis 7 的 Lua 会放行跨槽访问,静默写错节点)。见 redis-cluster.md; - 租约按队列:
{prefix}leases:{queue}ZSET(score = lease_until),孤儿回收逐队列扫描。
10. Worker 运行时¶
10.1 生命周期¶
stateDiagram-v2
[*] --> INIT
INIT --> WARMUP: 导入任务 · 校验注册表 · 连接 transport
WARMUP --> RUNNING: 注册心跳,开始拉取
RUNNING --> DRAINING: SIGTERM / SIGINT —— 停止拉取,跑完在途任务
DRAINING --> STOPPED: 在途任务跑完
DRAINING --> STOPPING: 超过 shutdown_timeout
STOPPING --> STOPPED: nack 未开始的任务
STOPPED --> [*]
SIGTERM/SIGINT触发优雅退出;shutdown_timeout(默认 30s)后可被 SIGKILL。- 未开始执行的消息立即 nack 回队,不留在内存里丢。
- 启动时校验:本 worker 注册表里缺失的 task 名 → 启动即失败(提前发现部署不一致),而不是运行时才报错。
10.2 执行池¶
| Pool | 适用 | 说明 |
|---|---|---|
solo |
调试 | 同步顺序执行,异常直接冒到终端 |
threads |
IO 密集(默认,决策 §20-3) | 轻量、共享内存;无法强杀,超时只能协作式 |
processes |
CPU 密集 | ✅ 真隔离、可强杀、可硬超时;任务参数必须可序列化。硬超时任务每次投递起一个独立子进程(超时 terminate/kill),其余任务走 ProcessPoolExecutor 复用子进程;子进程需要 app_spec(--app module:attr / TASKMQ_APP)重建 App |
asyncio |
任务本身是 async def |
✅ 单事件循环 + 线程槽位:async 任务体在 loop 上 await,同步任务体在池线程里跑;状态写入/ack/钩子都不占 loop |
- 混用:
async def任务在 threads/processes 池里会被明确拒绝(启动即ConfigError),不做隐式asyncio.run()。 - 预 fork:processes 池的普通任务复用
ProcessPoolExecutor的子进程;声明hard_timeout的任务为了可强杀, 每次投递起独立子进程(macOS spawn 约 0.2–0.4s)——所以 processes 池适合 CPU 密集的长任务,海量小任务用 threads。 hard_timeout在非 processes 池只告警不报错(threads 无法强杀,§10.6)。
10.3 拉取、预取与背压¶
prefetch是绝对值(默认 = concurrency),是「本 worker 最多同时持有多少条未 ack 消息」。- reserve–start 耦合:只在真有空闲执行槽位时才 reserve,池内队列有界。否则「预取占着槽位但活还没开始」会挡住紧急任务(决策 §20-9)。
- 让位(yield):更高优先级消息可见时,未开始的预留必须交回队列(
yields独立计数,不计入deliveries);正在执行的永不打断。详见 priority.md §3。 - 每个队列独立配额,避免一个大队列把 worker 的所有槽位占满(Celery 的经典队头阻塞)。
- 传输层不可用 → 指数退避重连,worker 保持存活并上报状态,不 panic 退出。
- 内存水位保护:当
len(inflight) >= prefetch时停止 reserve。
10.4 孤儿回收与心跳¶
- ✅ 已实现:worker 每
min(heartbeat_interval, lease/3)(下限 50ms)续租在途任务的租约。 没有这一步,跑得比lease还久的任务会被自家 reap 判成孤儿并重投——at-least-once 下就是重复执行 (实测:lease=0.3s 的 1s 任务被重投 14 次后超时;tests/test_acceptance.py::test_long_running_task_lease_is_renewed守住)。 - ✅ 已实现:memory / sqlite(
workers表 + heartbeat 索引)/ redis({prefix}worker:{id}hash +workersZSET)三家 transport 都支持register_worker / heartbeat_worker / deregister_worker / list_workers; worker 在首次 poll 时登记,之后每heartbeat_interval刷新,优雅退出主动注销;taskmq status输出WORKERS段(queues / pool / concurrency / 心跳年龄,超过max(3×interval, 30s)标[stale])。 worker 侧心跳失败只记 debug 日志,绝不打死 worker。 - 租约默认
lease = timeout * 2 + 30s;任何进程都可以调用reap_expired_leases()(幂等,用UPDATE ... WHERE state='reserved' AND lease_until < ?实现),把死掉 worker 的消息重新可见。 - 超过
max_deliveries的消息进 DLQ,并记录last_error和claim history。
10.5 速率限制与并发键¶
rate_limit="100/m":worker 内 token bucket(桶初始满,允许突发 = limit);超限时用transport.defer()把消息放回队列——defer 不消耗deliveries,所以被限流反复推迟的任务 不会被毒丸保护误送 DLQ。跨 worker 的分布式限流仍属 Phase 2。concurrency_key="user:{user_id}"(模板在 submit 时用任务 kwargs 渲染,渲染失败即ConfigError): 同 key 任务在集群内串行。实现是 transport 的命名租约(acquire_lease/release_lease/renew_lease, SQLite 落在leases表、memory 在进程内),拿不到就 defer 稍后再来。 租约 TTL =max(lease, task.timeout + 30s),worker 猝死也能自动释放。 注意:租约 owner 是 worker_id,所以同一 worker 内的自查必须与获取原子(Worker用一把锁 保护_held_keys),否则同 worker 的多个线程会互相「自认持有」而并发跑起来。这是 Celery 完全没有的能力。
10.6 超时(三层,语义写清楚)¶
| 类型 | 机制 | 可靠性 |
|---|---|---|
软超时 timeout |
注入 ctx.deadline,任务需要自行检查或通过 ctx.remaining() 控制 |
协作式,依赖任务配合 |
硬超时 hard_timeout |
✅ 已实现:processes 池超时直接 terminate/kill 子进程,任务按策略重试或进 DLQ;threads/solo/asyncio 池不支持(启动时警告)。被 kill 的任务不会执行 on_failure/on_retry/after_return(进程已经没了) |
强可靠(仅进程池) |
取消 revoke(terminate=True) |
processes 池 kill;否则等任务自己结束 | 同上 |
不承诺「线程池也能硬超时」——Celery 的 soft/hard time limit 正是这里最容易让人误判的地方。
10.7 优雅退出的时序¶
- 收到信号 → 停止
reserve。 - 对已 reserve 但未开始的消息
nack(requeue=True)。 - 等待在途任务至多
shutdown_timeout;期间每heartbeat续租。 - 完成任务正常 ack;未完成的保持租约到期后自动重投。
- 注销 worker 记录,退出。
11. 执行语义¶
11.1 Ack 策略¶
| 策略 | 语义 | 适用 |
|---|---|---|
on_receipt |
拿到就 ack,崩溃即丢 | 可丢的埋点/指标 |
on_success(默认) |
成功才 ack,失败/崩溃会重投 | 绝大多数任务 |
on_completion |
无论成败都 ack(失败进 DLQ,不重投) | 不希望重试的写操作 |
11.2 重试¶
Retry(
max_attempts=5,
backoff="exp", # "fixed" | "linear" | "exp"
base=1.0, factor=2.0, max_delay=600,
jitter=True, # 全抖动,防惊群
retry_on=(TemporaryError, TimeoutError), # 白名单;其余异常不重试
)
delay = min(max_delay, base * factor ** (attempt - 1)) * uniform(0.5, 1.5)- 重试通过重新入队 +
visible_at实现,不在 worker 内部 sleep(sleep 会占住并发槽位,Celery 的常见事故)。 - 非白名单异常默认不重试,直接 FAILED → DLQ。避免 Celery 里「不小心把所有 bug 重试 3 次」。
- 重试次数用尽 → DLQ,带完整的失败历史。
11.3 幂等¶
task.submit(..., key="order-42:email"):transport 在idempotency表内做原子插入。 插入成功 → 正常入队;冲突且未过期 → 返回已存在 job 的句柄(不重复执行)。- TTL 默认 24h,可配。key 只在窗口内去重,文档明确说明不是 exactly-once。
11.4 过期¶
expires到期仍未开始执行 → 标记EXPIRED,不执行、进事件流;仍在执行的任务不被中断(除非配了硬超时)。
11.5 失败分类¶
| 分类 | 行为 |
|---|---|
| 可重试(白名单) | RETRYING → 退避重投 |
| 不可重试(其他异常) | FAILED → DLQ,保留 traceback |
| 致命(协议/编码/注册表错误) | 直接 DLQ + 告警,不重试不重投 |
| 过载(队列满) | 生产者侧阻塞或抛 QueueFull(可配) |
11.6 序列化¶
- 默认
msgspec(决策 §20-5,速度优先、类型更严格)。serializer="json"退回标准库json,纯标准库可运行。 - 注册自定义类型:
app.register_codec(Money, encode=..., decode=...),未知类型在编码时报错。 - 反序列化只允许白名单类型,禁止
__reduce__之类的路径。 pickle需显式开启且启动时打警告(Phase 2 再评估是否移除)。
12. Result 后端¶
- 默认不存返回值:只要状态。需要
h.get()才配 backend,且必须给result_ttl(无 TTL 会被拒绝)。 - 接口:
set_result / get_result / set_state / get_state / forget(job_id)。 - 失败信息包含远端 traceback 字符串,
h.get()抛RemoteError(保留__cause__语义之外的原始 traceback 文本)。 - 不要拿 backend 当队列:明确写进文档。工作流依赖走调度器(Phase 2),不靠轮询 backend。
- TTL 清理由任意进程的
reap顺带完成,不需要额外守护进程。 - 实现现状:结果与状态同源落在 transport 里(SQLite 的
jobs表 / Redis 的{prefix}job:{id}hash), 没有拆出独立的ResultBackend接口 —— 好处是不会出现「状态说成功、结果丢了」的两套一致性问题; 真要换后端(比如 S3 存大结果)再按 §12 的接口拆。
13. 调度器(beat)¶
from taskmq.schedule import cron, every
app.schedule(
cron("send_daily_report", "0 9 * * *", tz="Asia/Shanghai"),
every("cleanup", seconds=300, args=[]),
)
- ✅ 独立进程
taskmq beat [--once] [--poll N],也可taskmq dev与 worker 合并(仅开发)。 - ✅ 选主:所有 beat 副本抢 transport 的
__beat__命名租约(TTL 30s,每 tick 续租; 优雅退出主动释放,follower 立刻接手);抢不到者进入待命、绝不触发。多副本不会重复触发 (Celery beat 的经典坑)。transport 不支持命名租约时降级单副本并告警。 - ✅ misfire 策略显式声明:
skip(默认,错过多次就跳过)/run_once(补一次,不补一堆)。 - ✅ 时区用
zoneinfo;cron 按本地时区计算,DST 边界(美东 2026-03-08 拨快)单测覆盖。 - ✅ cron 解析用标准库自己实现:5 字段 +
@daily等别名,支持*/a/a-b/a,b/*/n/a-b/n; weekday0(或7)= 周日;day-of-month 与 day-of-week 同时限定时按 Vixie 规则取或。 - ✅ 调度定义放代码里(
app.schedule(cron(...), every(...))),或用--schedule module:attr指向 App / Schedule 列表(生产推荐代码化,可 review、可测试)。 - ✅ 状态:每条调度的「上次观测时间」落在 JSON 文件(
--state,默认taskmq.beat.json), 只有 leader 写;首次部署只记基准不补跑,避免刚上线刷一堆任务。 - ⏳ 待办:多副本跨主机需要共享 transport(SQLite 共享盘 / Redis);beat 状态未来可以搬进 transport 表。
14. 可观测性¶
14.1 事件(结构化,默认输出 JSON 到 stdout)¶
| 事件 | 关键字段 |
|---|---|
task.submitted |
job_id, task, queue, key |
task.started |
job_id, worker, attempt, queue |
task.succeeded |
job_id, runtime, attempt |
task.failed |
job_id, error_type, error, attempt, will_retry |
task.retrying |
job_id, next_visible_at, attempt |
task.dropped |
job_id, reason(expired/revoked/max_deliveries) |
queue.depth |
queue, pending, inflight, dead(周期采样) |
worker.heartbeat |
worker, inflight, queues, concurrency |
transport.error |
op, error, retry_in |
- 所有事件都带
job_id和trace,可与 OpenTelemetry 的 traceparent 串起来。 - ✅ 已实现(
taskmq/otel.py,可选依赖taskmq-py[otel],core 不依赖):task.started→start_span("task <name>"),属性用标准 messaging 语义约定 (messaging.system=taskmq/destination.name/message.id/message.type)+taskmq.attempt/priority/worker;task.succeeded|failed|retrying→ 结束 span 并设状态(retrying 记 ERROR +taskmq.will_retry);task.deferred等其它事件 → 挂成当前 span 的 event。 用法:Config(events="otel")或app.add_sink(OtelEventSink(tracer))(可注入 tracer,便于测试/自定义 exporter)。
14.2 CLI 状态¶
$ taskmq status
WORKERS
w-1 ip 10.0.0.3 queues email,default pool=threads inflight=3/8 last_hb 2s ago
QUEUES
email pending=120 inflight=3 dead=4 rate=52/s
default pending=0 inflight=0 dead=0 rate=0/s
BEAT
leader w-1, next: send_daily_report in 3h12m
14.3 日志¶
- 结构化 JSON,字段固定:
ts, level, event, job_id, task, attempt, worker, queue, msg。 ctx.log自动注入上下文;本地开发可切--log-format=pretty。
15. 配置模型¶
@dataclass(frozen=True, slots=True)
class Config:
transport: str = "memory://"
result: str | None = None
result_ttl: int | None = None
default_queue: str = "default"
queues: dict[str, QueueConfig] = field(default_factory=dict)
concurrency: int = 4
pool: Literal["solo", "threads", "processes", "asyncio"] = "threads"
prefetch: int | None = None # None = concurrency
lease: float = 60.0
heartbeat_interval: float = 10.0
shutdown_timeout: float = 30.0
serializer: Literal["json", "msgspec"] = "msgspec"
max_message_bytes: int = 256 * 1024
timezone: str = "UTC"
log_format: Literal["json", "pretty"] = "json"
poll_interval: float = 0.05
max_poll_interval: float = 0.5
eager: bool = False
- 环境变量:
TASKMQ_TRANSPORT、TASKMQ_CONCURRENCY…(TASKMQ_<UPPER_FIELD>),显式列出,不做通配魔法。 - 优先级:显式参数 > 环境变量 > 默认值。没有配置文件、没有
app.conf.update()。 Config.validate()在 App 构造时执行:URL 合法性、pool 与任务类型兼容、result 必须有 TTL、concurrency>=1等。
16. 测试与本地开发¶
from taskmq.testing import worker_for
def test_eager_logic():
app = App(Config(transport="memory://", eager=True))
@app.task
def add(a: int, b: int) -> int:
return a + b
h = add.delay(1, 2)
assert h.state == "SUCCEEDED" and h.get() == 3
def test_retry_then_success():
app = App(Config(transport="memory://"))
with worker_for(app) as w: # 线程池 worker,跑在当前进程
h = flaky.delay()
w.run_until_idle(timeout=5) # 确定性:跑空就返回
assert h.state == "SUCCEEDED"
assert h.info["attempt"] == 2
def test_worker_crash_redelivery():
# 故障注入:在第 2 次 attempt 时模拟 worker 崩溃
with worker_for(app, fail_at={"flaky": 2}) as w:
...
# 断言消息被重新投递,且进 DLQ 前只跑 max_deliveries 次
eager=True:delay()同步执行,用于纯逻辑单测。worker_for(app):真实跑完整链路(transport → pool → ack),但全在进程内,可断点调试。- 时间冻结:
freeze_time支持,测 ETA/退避/过期不用 sleep。 - 故障注入:崩溃、transport 报错、ack 失败,都要有一等公民的测试辅助。
- 目标:框架自身测试覆盖率 ≥ 90%,且所有并发/崩溃场景用确定性测试而非「睡 2 秒碰运气」。
- 工具链:
make check= ruff + mypy + pyright + pytest,不绑 uv(默认用.venv的解释器); pip 路径pip install -e ".[dev]",uv 路径(可选)uv sync→uv run pytest;.python-version固定 3.10(最低支持版本,真机验证)。
17. 失败模式矩阵(我要的行为)¶
| 场景 | 期望行为 |
|---|---|
| worker 被 SIGKILL | 租约到期后消息重新可见;deliveries +1;超限进 DLQ |
| worker 优雅退出 | 未开始的 nack 回队;在途的跑完再退 |
| transport 短暂不可用 | worker 保持存活、指数退避重连、持续上报状态 |
| transport 长时间不可用 | 生产侧可配「阻塞」或「抛 QueueFull」;绝不静默丢弃 |
| 任务超时(进程池) | kill 进程 → 可重试则重试 → 否则 DLQ |
| 任务超时(线程池) | 只能协作式;启动时明确警告此组合的限制 |
| 反序列化失败 | 消息直接进 DLQ,附带原始字节,不 crash worker |
| 任务名不存在 | worker 启动时校验失败(fail fast) |
| 队列积压 | queue.depth 事件 + CLI 可见;可选告警钩子 |
| 结果后端不可用 | 任务本身照样执行(业务结果 > 结果记录),后端错误只记事件 |
| 重复投递 | 幂等键去重;没有 key 的任务要求业务自己幂等(文档写明) |
18. 与 Celery 对照速查¶
| 能力 | Celery | taskmq |
|---|---|---|
| 本地跑起来 | 需要 broker | 需要 0 个外部服务 |
| 默认投递语义 | 收到即 ack(可能丢) | 成功才 ack(可能重) |
| prefetch | 倍数(反直觉) | 绝对值 |
| 配置 | 懒加载全局 conf | frozen dataclass,构造即校验 |
| 任务发现 | autodiscover 魔法 | 显式 include |
| 重试 | 两套机制 | 一套声明式 + 一个手动入口 |
| 结果后端 | 语义组合爆炸 | 显式对象 + 强制 TTL |
| 多 beat | 会重复触发 | lease 选主 |
| 死信队列 | 需自己搭 | 内置 + CLI 重放 |
| 按 key 串行 | 无 | concurrency_key |
| 结构化事件 | 需 flower/自建 | 内置 JSON 事件 + OTel 适配 |
| 测试 | 全局 eager 开关 | memory transport + worker_for |
| 类型提示 | 弱 | py.typed + 泛型 Result |
19. 里程碑¶
Phase 0 — 可用的 MVP(先跑通闭环)✅ 完成¶
- [x]
Config/App/@task/TaskHandle公开 API - [x] 类式任务 + 生命周期钩子 +
bind=True:Task基类、app.register()、base=、bind=True(self/self.request)、五个钩子(设计 tasks.md,T1–T11 全部 ✅) - [x]
MemoryTransport(原子 claim、租约、孤儿回收、DLQ、幂等键) - [x]
SqliteTransport:WAL +BEGIN IMMEDIATE原子 claim、优先级定档 + 档内按 weight 轮询、yields/yieldable让位、DLQ、幂等键、结果/meta 持久化 - [x]
threads+solo池 - [x] 默认 at-least-once、
Retry策略、DLQ - [x] 优先级插队:全局严格优先 + 未开始预留让位(
yields独立计数)(决策 §20-9);status --by-priority已可用 - [x] 协议 v1:默认 msgspec 编解码,
serializer="json"可切 - [x] CLI:
taskmq worker [-Q] [-c] [--once]/status [--by-priority]/dlq list|replay/call(入口taskmq,--app module:attr或TASKMQ_APP) - [x]
taskmq.testing:eager+worker_for+run_until_idle - [x] 质量门:Python 3.10 上
ruff+mypy+pyright+ 55 个测试全绿(协议 / 两个 transport 不变量 / G1·G2·G3 插队 / 任务类与钩子 / CLI / 验收) - [x] 验收(全部达成):无 Redis 环境下 1000 任务 × 2 个 App/worker 跑完且无重复执行 ✔;
worker 进程被
os._exit(9)真杀掉后租约到期消息重新可见并被重新执行(deliveries==2)✔; 超过重试上限进 DLQ 且可重放(重放后attempt归 1 重新开始)✔;紧急任务在 claim 时插队 ✔。 见tests/test_acceptance.py。
Phase 1 — 生产可用¶
- [x]
processes/asyncio池 + 硬超时:ProcessPool(硬超时任务独立子进程可强杀,其余复用进程池) /AsyncPool(单事件循环 + 线程槽位);async 任务与非 asyncio 池的混用在启动期ConfigError - [x]
RedisTransport:零依赖 RESP + Lua 原子 claim/状态转换/孤儿回收/命名租约,键前缀隔离; 结果与 meta 存在{prefix}job:{id}(独立RedisResultBackend未拆,理由见 §12)。已对 Redis 7.0.5 实测 - [x]
workers心跳表 +taskmq status的 worker 列表(见 §10.4) - [x] PostgreSQL transport(
taskmq/transport/postgres.py):FOR UPDATE SKIP LOCKED原子 claim (多机并发零重复、零阻塞)、ON CONFLICT … WHERE单条 CAS 命名租约、幂等键唯一约束、 worker 注册表与list_jobs(DAG 可用);一致性套件 16/16 通过、零降级声明;表前缀?prefix=隔离 - [x] AMQP transport(RabbitMQ)(
taskmq/transport/amqp.py):投递走 broker(x-max-priority排序、 per-message TTL + DLX 延迟回投、DLX 死信、未确认重投),状态走侧车?state=<transport url>(AMQP 没有 KV:job 状态/DAG 索引/命名租约/worker 表都由侧车提供,缺state=直接报错); 无锁帧里如实声明降级:跨队列不做全局严格优先、priority_stats为空、reap_expired_jobs是 no-op; 一致性套件 15/16(global_priority按声明跳过),含 DAG 端到端 - [x] OTel 适配器
taskmq/otel.py+events="otel"(见 §14.2) - [x]
beat+ lease 选主 + misfire:cron()/every()(标准库 cron 解析 + zoneinfo)、taskmq beat/dev、__beat__租约选主(优雅退出释放)、skip/run_once、JSON 状态文件 - [x] 幂等键(Phase 0)、
concurrency_key、rate_limit:命名租约 + token bucket +defer(不消耗deliveries) - [x] 结构化事件:
Config.events(stdout/null)+EventSink协议 +CollectingSink;事件含task.submitted/started/succeeded/failed/retrying/deferred+worker.started/stopped;OTel 适配器见 §14.2 - [x] 覆盖率 ≥ 90%(当前 93%,
make coverage)+ 崩溃场景确定性测试(os._exit(9)租约恢复、 ack 幂等、迟到 ackLeaseLost、Redis 两种模式各跑一遍) - 验收:单机 4 worker × 8 并发压测达标;双 beat 副本无重复触发。
Phase 2 — 进阶¶
- [x] 原生 DAG 工作流(替代 chain/group/chord):
taskmq/workflow.py+ worker 侧依赖推进(事件 + 维护期补偿)taskmq workflow list/status/resume+supports_job_listing能力位;16 例测试,见 design/workflows.md(v1.0)
- [ ]
postgres://transport(SKIP LOCKED)、amqp:// - [ ] 分布式限流、优先级公平调度
- [ ] 任务级中间件(审计、多租户、权限)
- [ ] Celery 兼容层(§20-8):可选包
taskmq-celery,提供celery命名空间 shim(Celery、shared_task、delay/apply_async、Retry、chain/group/chord→ DAG 映射),存量项目渐进迁移
Phase 3 — 生态¶
- [ ] Web UI(只读看板 + DLQ 操作)
- [ ] 插件入口(
taskmq.transportentry point) - [ ] Celery 任务迁移工具(AST 级重写辅助)
20. 决策清单(✅ 已定 / 🟡 暂定 / ❓ 待拍板)¶
| # | 决策点 | 结论 |
|---|---|---|
| 1 | 包名 | ✅ import 名 taskmq;PyPI 发行名 taskmq-py(taskmq 与已有的 task-mq 被 PyPI 判为相似名称,不允许新建),命令行仍是 taskmq |
| 2 | 最低 Python 版本 | ✅ 3.9+ 运行时下限;开发与类型检查按 3.10(.python-version)。3.10 专属能力按需降级:dataclass(slots=True) 走 taskmq/_compat.py 的 _SLOTS,ParamSpec/Concatenate 在 3.9 从 typing_extensions 取;CI 在真 3.9 与 3.10 上各跑一遍全量测试。不使用 3.11+ 的 asyncio.timeout / TaskGroup / ExceptionGroup / tomllib |
| 3 | 默认执行池 | ✅ threads;CPU 密集显式切 processes;async def 任务强制 asyncio 池 |
| 4 | 默认投递语义 | ✅ at-least-once(成功才 ack)+ 可见性租约 + max_deliveries + DLQ;不承诺 exactly-once,重复用 key 幂等 |
| 5 | 默认序列化 | ✅ msgspec(core 唯一第三方依赖);serializer="json" 退回纯标准库 |
| 6 | SQLite 跨主机 | ✅ 一期只支持同机多进程 + 共享盘(NFS/SMB);跨主机走 Redis(Phase 1)/ PG(Phase 2);对象存储(OSS/S3)不做一期 transport,列为 Phase 2+ 候选 |
| 7 | 工作流优先级 | ✅ Phase 2 落地原生 DAG(声明式、依赖推进、不轮询):design/workflows.md |
| 8 | Celery 兼容层 | ✅ 做,但不是默认 drop-in:可选包 taskmq-celery 提供 celery 命名空间 shim,Phase 2 交付,便于存量项目快速切换 |
| 9 | 优先级语义 | ✅ 方案 D:全局严格优先 + 未开始预留让位(yields 独立计数)+ 平级队列轮询。G1 空闲即最高 / G2 未开始让位 / G3 不打断运行中;数值域 -9..9;aging 与兜底 Phase 2。详见 priority.md §10(P1–P19) |
相对 v0.1 的变化:Python 下限 3.11 → 3.10;序列化默认 stdlib json → msgspec(P1 相应改为「零外部服务」);Celery 兼容层从「明确不做」→ Phase 2 可选 shim。
21. 下一步¶
- ✅ 优先级设计已定稿:priority.md v1.0(P1–P19 全部 ✅ = 方案 D + 让位)。
- ✅ §20 决策已确认(仅第 7 条工作流优先级为 🟡 暂定)。
- ✅ Phase 0 完成:
protocol/transport.base/transport.memory/transport.sqlite/App/Task(类式 +bind=True+ 钩子)/Worker(含 G1–G3 插队)/testing/cli全部落地,Python 3.10 上make check(ruff + mypy + pyright + 55 测试)全绿。 - 🚧 Phase 1 进行中:
concurrency_key/rate_limit、结构化事件、processes/asyncio池 + 硬超时、beat+ 选主 + misfire、RedisTransport、workers心跳、OTel 适配、覆盖率 93% 已完成 —— Phase 1 收尾。✅ 插件机制(自注册后端)已落地:plugins.md v1.0。 ✅ Phase 2 第一刀 原生 DAG 工作流已落地(workflows.md v1.0)。 ✅ Phase 2 第二/三/四刀:amqp://(状态侧车)、Redis 平级加权轮询、Redis Cluster (?cluster=1,redis-cluster.md)已落地。 下一步 Phase 2:Celery 兼容 shim(taskmq-celery,canvas → DAG 映射)、独立result=后端(插件点已预留)。 - 🚧 分册拆分:已落地
priority.md;protocol.md/transport.md/worker.md/scheduler.md/testing.md待拆。
22. 扩展机制(插件 / 自注册后端)¶
第三方后端(例如阿里云 RocketMQ、公司内部 MQ)不改 taskmq 源码即可接入,三条路径:
- 打包 + entry point:插件包声明
[project.entry-points."taskmq.plugins"] rocketmq = "taskmq_rocketmq", 用户只写Config(transport="rocketmq://..."); - 私有环境显式加载:
app.load_plugins(["mycompany.mq"])/TASKMQ_PLUGINS=.../ CLI--plugins; - 直接给实例:
Config(transport=<Transport 实例>)(现有能力,框架不干预生命周期)。
规则:内建 scheme 优先(插件要覆盖必须显式 override=True);未命中时对 entry points 做懒发现
(标准部署不 import 任何第三方包,启动开销不变);插件必须如实声明 supports_leases /
supports_workers,缺能力启动即 ConfigError,不留到运行时;
taskmq.testing.transport_conformance() 让第三方后端自证语义与内建一致(内建三家也跑同一套件,
避免"给别人的契约"与"自己跑的契约"脱节)。插件可依赖的公开契约见 §7 稳定性边界。
已落地(v1.0,D1–D7 全部按建议确认):taskmq/plugins.py 注册表 + entry point 懒发现、
App.load_plugins/App.plugins、--plugins/TASKMQ_PLUGINS、codec/sink(+ pool 预留)扩展点、
ChildTask.plugins 子进程透传、15 个场景的 taskmq.testing.transport_conformance()、
Transport.limitations + status 展示,以及示例插件 examples/plugin_rocketmq。
设计稿与实现纪要:design/plugins.md。
文档版本 v0.2(2025)· 所有设计取舍以「可预测、可调试、失败可见」为最高优先级。 分册:priority.md(v1.0)· tasks.md(草案 T1–T8)· plugins.md(v1.0)· workflows.md(v1.0);待拆:
protocol.md/transport.md/worker.md/scheduler.md/testing.md。