fix(worker): 队列消息 reclaim 不查活跃租约——排队中的长任务被「死信→重入队」循环反复冲洗(DLQ 增长 + churn)
Nobody has claimed this yet.
Assessment
- Difficulty
- 4/5
- Estimated time
- 3-5 days
- Newbie friendliness
- 70/100
- Issue type
- Bug
- Clarity
- Clearly specified
- Activity status
- Active
- Tech stack
- python, redis
- Domain
- distributed-systems
Research direction
Start at plaita/server/task_queue.py _reclaim_one() (~line 481) and trace how XCLAIM fires without consulting the lease key {ns}:execution:lease:{execution_id}, then read plaita/server/flow_worker.py _dead_letter_guard() (~line 1914) to see why 'queued, no lease' is treated as 'holder dead'. First run: reproduce queue saturation so a message waits > max_deliveries × claim period and watch times_delivered on plaita:flow:queue:v2 via XPENDING. Done means the four acceptance conditions in the issue hold: no delivery growth/DLQ additions while queued, kill -9 orphan recovery still re-enqueues, short-task throughput unchanged, and XPENDING entries stable in the saturation window.
Written by the indexing model from the issue text.
Description
优先级 P2 · 依赖:无
基线 HEAD = 44febf5
Summary
现象(实测):plaita:flow:queue:v2:dlq 3 条,全部 reason=max_deliveries=5(21:21:45 / 21:28:00 / 21:34:50 生成),payload 全为 #87 的同一条 start 任务(execution_id=99c142cce2b14539aca7a95b0e966391,tenant=default,任务入队时间 19:19:07);三条的 source_id 互为链条(1791371947969 → 1791379305016 → 1791379680771)=「死信→重入队新副本→再烧 5 次投递→再死信」的循环,循环周期 ≈6.4min。同窗口 XPENDING 全部压在 worker-mac-1(主控实测 15 条);21:40 主机 worker 重启后 claim 存量 PEL 又触发一波「交接重入队」,XLEN 15→25。
代码链路(亲读,origin/main):
- reclaim 只按 idle 判空闲,不查执行租约:
plaita/server/task_queue.py:481_reclaim_one()——xpending_range扫描(窗口 8)→idle < claim_min_idle_ms跳过 → 否则XCLAIM抢走并直接返回该任务重投(return task,deliveries < max 分支);租约/执行状态在本路径零查询。默认DEFAULT_CLAIM_MIN_IDLE_MS = 60_000(task_queue.py:18)。 - 死信守卫只覆盖「死信决策」,且把「非终态+无租约」当「持有者已死」:
plaita/server/flow_worker.py:1902-1911接线dead_letter_guard→_dead_letter_guard()(:1914):非终态 + 租约空 → 重新入队一份同体消息(delivery 归 1)后放行原消息死信——该分支在「派发≠执行、排队等待中(exec 已建但未开跑、无租约)」的常态下会被反复触发(当前队列饱和常态:keeper 超闸 + 双 worker 满载,排队数小时并不罕见)。 - 即:「排队中(未开跑)」与「持有者崩溃(孤儿)」在守卫眼里不可区分——前者本应静静等槽,实际被判定「死了」→ 死信+重投循环;后者才需要恢复路径。
- 补充观察:当 exec 真在跑(租约活)时,守卫会跳过死信(行为正确),但其队列消息仍在每个 claim 周期被
XCLAIM一次 →delivery_count持续虚增;触顶后每轮「跳过」也是 churn 噪声(消息留 pending 至终态)。
对照租约证据:同一 exec 当前 plaita:execution:lease:99c142cc… 存在且 TTL=86s(活)、status=running(start_time=21:50:04)——即队列里那条被反复死信重投的 start 任务,对应的执行此刻正在正常跑。
影响
- DLQ 常态增长 + 队列 churn:只要队列饱和(排队时长 > max_deliveries × claim 周期 ≈ 30-35min),排队消息就会被「死信→重入队」循环反复冲洗——DLQ 被无意义条目污染(真死信判别被淹没)、XLEN 虚增(15→25)、worker 重启窗口还会叠加「交接重入队」风暴。
- 恢复路径可信度受损:真孤儿(持锁 worker 崩溃)与「只是在排队」走同一分支,运维无法从 DLQ 分辨「需要处理的死信」与「排队过久的正常单」。
- 无双重执行(租约+ack 兜底有效)——本条只针对 churn 与 DLQ 噪声,不是数据正确性问题。
建议
- reclaim 前查活跃租约:
_reclaim_one在XCLAIM前查{ns}:execution:lease:{execution_id}(或先XINFO/状态键):租约活跃 → 跳过该条(不 XCLAIM、不计数),等其终态;仅对「无租约 && 非终态 && 超过更大宽限(如按 start_time 判真正孤儿)」者进入恢复路径。 - 死信判据与「排队中」解耦:给「非终态+无租约」分支加排队宽限(例如距
start_time/入队时间 < N 小时一律跳过,只留 pending),或在 exec 状态里区分queued/started两态供守卫判据使用(避免用租约缺失代理「持有者已死」)。 - 验收条件(可执行):
- 制造队列饱和(>30min 排队):该消息
times_delivered不再增长、DLQ 零新增; - 真孤儿(可 kill -9 持有 worker):仍能按恢复路径重入队并被接管(不要因修复 #1/#2 而切断恢复);
- 短任务(<60s 处理)吞吐不受影响;
XPENDING在饱和窗内条目数稳定,不出现「同一条目反复刷新 idle」。
- 制造队列饱和(>30min 排队):该消息
边界
- 不改
max_deliveries契约与 DLQ 本身;修复对象是 claim/死信判据。 - 与 reaper(60m idle + 租约门)分工:reaper 管「执行态终态化」,本条管「队列消息不因排队被误判死信」——两者判据应对齐(都以租约+状态为准)。
- Dominant language
- Python
- Stars
- 0
- Forks
- 1
- Avg merge
- 1h 44m
- Merged PRs (30d)
- 7
Getting set up
This project ships no dev container, Dockerfile or contributing guide, so setting up is up to you: start from its README, and see our first-contribution guide for the general steps.
First steps
- Read the whole issue, then the project's contributing guide.
- Comment on the issue to say you are picking it up — it saves two people doing the same work.
- Fork the repository and make your change on a branch.
- Open a pull request that references the issue number.
More from jeffkit/plaita
-
Difficulty 1/5 Under an hour Newbie friendliness 92/100
-
keeper-ignore
Difficulty 5/5 Over a week Newbie friendliness 25/100
-
Difficulty 5/5 Over a week Newbie friendliness 38/100
-
Difficulty 5/5 Over a week Newbie friendliness 35/100
-
Difficulty 4/5 3-5 days Newbie friendliness 55/100
Similar issues
-
area/install reliability
Difficulty 2/5 1-3 hours Newbie friendliness 75/100
Maintainers usually reply within 1 day
-
Difficulty 2/5 1-3 hours Newbie friendliness 83/100
FluidNumerics/fluid-walk-blocker#191 ·
Maintainers usually reply within 1 day
-
Difficulty 2/5 1-3 hours Newbie friendliness 62/100
TransformerLensOrg/TransformerLens#1868 ·
Maintainers usually reply within 1 day
-
Difficulty 2/5 1-3 hours Newbie friendliness 76/100
Maintainers usually reply within 1 day
-
Difficulty 1/5 Under an hour Newbie friendliness 85/100
climate-analytics-lab/jax-gcm#1057 ·
Maintainers usually reply within 1 day