本文使用「署名 4.0 国际 (CC BY 4.0)」许可协议,欢迎转载、或重新修改使用,但需要注明来源。 [署名 4.0 国际 (CC BY 4.0)](https://creativecommons.org/licenses/by/4.0/deed.zh) 本文作者: 苏洋 创建时间: 2026年09月08日 统计字数: 12018字 阅读时间: 24分钟阅读 本文链接: https://soulteary.com/2026/09/08/phorge-modernization-part-15-take-over-work-queue-with-go.html ----- # Phorge 现代化改造实战(十五):用 Go 接管 Phorge 工作队列,如何避免新旧消费者互相抢任务 本文是“Phorge 现代化改造实战”系列第十五篇。上一篇为 Go 服务增加了统一的 Phorge API 入口;这一篇进一步接管有状态的后台任务,把队列存储、任务调度和仍留在 PHP 中的任务实现拆开迁移。 ## 系列导航 1. [从没有官方镜像到 Docker Compose 跑起来:建立最小可运行基线](https://soulteary.com/2026/09/06/phorge-modernization-part-1-from-no-official-image-to-docker-compose.html); 2. [改进容器化的七个细节:补齐权限、持久化、依赖和探活](https://soulteary.com/2026/09/06/phorge-modernization-part-2-seven-containerization-details.html); 3. [接入 Stargate:把 Forward Auth 的信任边界做完整](https://soulteary.com/2026/09/06/phorge-modernization-part-3-stargate-forward-auth-trust-boundary.html); 4. [拆分模块到 Gorge:无侵入改造不等于不碰文件](https://soulteary.com/2026/09/06/phorge-modernization-part-4-split-modules-to-gorge.html); 5. [替换 diff 子进程:兼容不等于逐字一致](https://soulteary.com/2026/09/07/phorge-modernization-part-5-replace-diff-subprocess.html); 6. [替换实时通知服务:为什么 HTTP 501 反而表示正常](https://soulteary.com/2026/09/07/phorge-modernization-part-6-replace-realtime-notification-service.html); 7. [迁移邮件服务:先分清哪些失败不该重试](https://soulteary.com/2026/09/07/phorge-modernization-part-7-migrate-mail-service.html); 8. [迁移搜索服务:写进索引不等于搜得到](https://soulteary.com/2026/09/07/phorge-modernization-part-8-migrate-search-service.html); 9. [迁移文件存储:写得进去也要读得回来](https://soulteary.com/2026/09/07/phorge-modernization-part-9-migrate-file-storage.html); 10. [迁移 Webhook 投递服务:先解决重复投递](https://soulteary.com/2026/09/07/phorge-modernization-part-10-migrate-webhook-delivery.html); 11. [升级 Gorge 的 HTTP 框架:接口没变,行为也不能变](https://soulteary.com/2026/09/08/phorge-modernization-part-11-upgrade-http-framework.html); 12. [兼容 Elasticsearch 5、6、7:版本配置决定整个索引结构](https://soulteary.com/2026/09/08/phorge-modernization-part-12-elasticsearch-version-compatibility.html); 13. [联调六项外部服务:容器在运行,不代表业务已经切换](https://soulteary.com/2026/09/08/phorge-modernization-part-13-integrate-six-external-services.html); 14. [为 Go 服务建立统一的 Phorge API 入口:网关可以换框架,Conduit 协议不能变](https://soulteary.com/2026/09/08/phorge-modernization-part-14-unified-conduit-api-gateway.html); 15. **用 Go 接管 Phorge 工作队列:如何避免新旧消费者互相抢任务。** ## 写在前面 前几篇迁走的高亮、搜索、发信和文件存储,大多可以概括成“Phorge 发一个 HTTP 请求,Gorge 返回结果”。工作队列还包含状态、租约、重试和消费权。两个 Go 服务解决代码实现,Phorge 新接口完成接线;接管还需要停掉旧的 taskmaster,并证明任务只由预期的消费者处理。 这次对应 [Gorge 2026.09.08-r4](https://github.com/soulteary/gorge/releases/tag/2026.09.08-r4) 与 [Phorge 2026.09.08-r4](https://github.com/soulteary/phorge/releases/tag/2026.09.08-r4)。Gorge 一侧通过 [PR #14](https://github.com/soulteary/gorge/pull/14) 加入 `gorge-taskqueue` 和 `gorge-worker`,Phorge 一侧则补上 HTTP 客户端、Conduit 执行入口、配置项、Setup Check 与 Compose 接线。 这次迁移需要继承 Phorge 队列的线协议:ID 从哪里分配、哪几列属于活跃表、租约如何抢占、yield 怎样表示,以及 Conduit 到底接受什么请求格式。 ## 先拆清楚队列、消费者和任务实现 Phorge 原本由 `PhabricatorTaskmasterDaemon` 不断执行“租约 → 运行 → 归档”。任务主要落在三张表里: - `worker_activetask`:还在排队、等待重试或正在被租用的任务; - `worker_taskdata`:任务负载,与活跃任务通过 `dataID` 关联; - `worker_archivetask`:成功、永久失败或取消后的归档记录。 r4 把这条链拆成两个二进制: ```mermaid flowchart LR P["Phorge PHP"] -->|"入队与队列操作"| Q["gorge-taskqueue :8090"] Q --> S[("MySQL 或 Redis")] W["gorge-worker :8170"] -->|"租约与回报"| Q W -->|"worker.execute"| C["gorge-conduit :8150"] C --> P ``` `gorge-taskqueue` 管队列状态,`gorge-worker` 管消费循环。两者保留 HTTP 边界,便于独立扩缩容和隔离故障: | 组件 | 负责什么 | 不负责什么 | | ----------------- | ------------------------- | ------------------- | | `gorge-taskqueue` | 入队、租约、重试、yield、取消、归档、队列统计 | 不执行具体业务任务 | | `gorge-worker` | 领取任务、选择 handler、并发执行、回报结果 | 不直接读写 MySQL 或 Redis | | Phorge PHP | 产生任务;暂未 Go 化的任务仍在 PHP 中执行 | 配置接管后不再直接驱动队列 | 因此,一个 taskqueue 前可以挂多个 worker,也可以给重任务单独开一个带 `TASK_CLASS_FILTER` 的 worker 池。若把两者塞进同一进程,worker 要么绕过 HTTP 直接碰 `Store`,要么仍然兜一圈本机网络;两种做法都削弱了独立扩缩容和故障隔离。 这也解释了为什么 render 与 diff 可以合在一个二进制里,而 taskqueue 与 worker 不应合并:前者是无状态计算,后者之间存在明确的状态所有权和部署边界。 ## taskqueue 的难点集中在旧表语义 `gorge-taskqueue` 在 `/api/queue` 下提供 `enqueue`、`lease`、`complete`、`fail`、`yield`、`cancel`、`awaken`、`stats` 与任务查询接口。HTTP 层较薄,难点集中在底下三张旧表特殊的 ID、租约和归档语义。 ### `worker_activetask.id` 不能直接用 `LastInsertId()` 最初按常规思路插入 `worker_activetask`,MySQL 会直接报: ```text Error 1364 Field 'id' doesn't have a default value ``` 原因是 `PhabricatorWorkerActiveTask` 使用 `IDS_COUNTER`。它的 `id` 由 Phorge 在应用层通过 `lisk_counter` 分配,没有使用 `AUTO_INCREMENT`。Go 侧如果另建一条自增序列,即使短期能写进去,也会和原生 `phd` 从两个来源取号,最终发生碰撞。 修复后的 `enqueue` 在同一个事务里复用 Phorge 的计数器写法: ```go func nextCounterValue( ctx context.Context, tx *sql.Tx, counterName string, ) (int64, error) { res, err := tx.ExecContext(ctx, `INSERT INTO lisk_counter (counterName, counterValue) VALUES (?, LAST_INSERT_ID(1)) ON DUPLICATE KEY UPDATE counterValue = LAST_INSERT_ID(counterValue + 1)`, counterName) if err != nil { return 0, err } return res.LastInsertId() } ``` 这里虽然仍调用了 `LastInsertId()`,拿到的却是 `LAST_INSERT_ID(...)` 主动发布的计数器结果。`worker_activetask` 没有自增值;`worker_taskdata` 仍然使用 `AUTO_INCREMENT`,所以同一次入队事务里的两张表采用不同的 ID 机制。 这条约束最容易被“统一风格”的重构破坏。统一目标是与 Phorge 共用同一条序列,代码形状可以不同。 ### 归档时要丢弃 `failureTime` 任务成功、永久失败或取消时,taskqueue 会在一个事务中把活跃行写入 `worker_archivetask`,再从 `worker_activetask` 删除。这样任务不会同时出现在两张表,也不会在两个操作之间凭空消失。 归档表与活跃表的字段并不对称。归档表没有 `failureTime`:这个字段只用于活跃任务的重试退避,归档时会被丢弃。归档行的结构可以概括成: ```text active columns - failureTime + result + duration + archivedEpoch ``` 若把 `failureTime` 也塞进归档 SQL,MySQL 会报 `Unknown column 'failureTime'`,已经执行完的任务留在活跃表,租约过期后又会被领走。 另一个容易写错的是 `result`。它使用 Phorge 定义的整数编码:`0` 表示成功,`1` 表示失败,`2` 表示取消。写成字符串后,Go 侧虽然能写入,Phorge 的管理界面却无法按原有语义读取。 ### 没有 `status` 列,状态都藏在租约字段里 租约分两阶段完成:先取从未租过的任务,再用过期任务补足名额。 ```sql -- 第一阶段:新任务优先 SELECT id FROM worker_activetask WHERE leaseOwner IS NULL AND leaseExpires IS NULL ORDER BY priority ASC, id ASC LIMIT ?; -- 第二阶段:租约已经过期的任务 SELECT id FROM worker_activetask WHERE leaseExpires < ? ORDER BY priority ASC, id ASC LIMIT ?; ``` `priority` 数字越小越紧急,同优先级按 `id` 保持 FIFO。选出候选行后,还要用带原条件的 `UPDATE` 取得租约,再只返回本次调用成功持有的行。这个二次条件非常重要:`SELECT` 只给出候选列表,不能证明所有权。 Phorge 的表里也没有单独的 `status`。yield 通过 `leaseOwner="(yield)"` 这个哨兵和未来的 `leaseExpires` 表示,`awaken` 再精确识别它并提前结束等待。字段名、空值语义和哨兵字符串都属于兼容协议,不能按 Go 结构体的内部习惯自由重命名。 ## MySQL 和 Redis 的回退能力不同 taskqueue 的 `Store` 同时有 MySQL 与 Redis 两个实现。Redis 用 hash、sorted set 和 Lua 脚本复现优先级、两阶段租约、yield 与归档移动,让并发操作保持原子性。worker 完全不知道后端是哪一个,它只通过 HTTP 领取和回报任务。 “配置缺失就回到原生 SQL”只描述了 PHP 的代码分支,两种后端的回滚能力仍有明显差异: | 后端 | 数据放在哪里 | 清空 `gorge.taskqueue.uri` 后会怎样 | | ----- | ------------------------------------ | -------------------------------------- | | MySQL | 与 Phorge 共用 `{namespace}_worker` 三张表 | 原生 SQL 能继续看见原队列,适合直接切回 | | Redis | Gorge 自己的 Redis key | Phorge 只会回到 MySQL;Redis 中未完成的任务不会自动搬过去 | 所以 MySQL 模式可以理解为“同一份数据,换一个调度入口”;Redis 模式则涉及数据面迁移。切到 Redis 前要定义排空、迁移或回放方案,不能因为代码有 fallback 就假定数据也能无缝回退。 这里还有一个配置层的小口值得补上:r4 的 `newStore()` 对 `redis` 之外的所有值都回落到 MySQL。把 `GORGE_TASKQUEUE_BACKEND=redsi` 拼错后,程序不会在启动时明确拒绝,而会尝试连接 MySQL。更稳妥的做法是只接受 `mysql` 与 `redis`,未知值直接启动失败,让错误尽量靠近配置现场。 ## worker 不碰数据库,未重写的任务经 Conduit 回到 PHP `gorge-worker` 的核心是一个租约循环: 1. 从 taskqueue 批量领取任务; 2. 按 `taskClass` 从 registry 找 handler; 3. 在租约截止时间内执行; 4. 成功走 `complete`,永久失败走 `fail(permanent=true)`,yield 走 `yield`,其余错误进入临时失败与重试; 5. 回报短暂失败时最多重试五次,仍失败则保留租约,等待正常的过期恢复。 本地 handler 目前只原生实现少量任务,例如 `FeedPublisherHTTPWorker`。配置 `GORGE_WORKER_CONDUIT_URL` 后,registry 会安装一个 fallback,把未在本地实现的任务统一交给 Phorge 新增的 `worker.execute` 方法。 于是执行路径形成一个回环:Phorge 产生任务,Go 负责排队与调度,尚未迁出的业务代码仍由 PHP 执行,结果再由 Go 回报给 taskqueue。这个过渡结构很实用,因为“替换队列运行时”和“重写所有任务业务逻辑”不必一次完成。 ### Conduit 不接受 JSON body 委派最初报的是: ```text invalid character '<' looking for beginning of value ``` 这个错误通常说明客户端准备解析 JSON,服务器却返回了 HTML 错误页。Phorge 的 `PhabricatorConduitAPIController` 不接受 `application/json` 请求体;它要求表单编码,并把参数 JSON 放进 `params` 字段。 Go 客户端最终按 Phorge 自家客户端的线格式发送: ```go params["__conduit__"] = conduitMeta paramsJSON, err := json.Marshal(params) if err != nil { return nil, err } form := url.Values{} form.Set("params", string(paramsJSON)) form.Set("output", "json") req, err := http.NewRequestWithContext( ctx, http.MethodPost, baseURL+"/api/"+method, strings.NewReader(form.Encode()), ) req.Header.Set("Content-Type", "application/x-www-form-urlencoded") req.Header.Set("Accept", "application/json") ``` 网关自己的认证仍放在 `X-Service-Token`。遇到非 JSON 响应时,客户端会返回 HTTP 状态和截断后的正文片段,不再只留下一个脱离上下文的字符解析错误。 ### `worker.execute` 只执行任务,不接管队列 Phorge r4 新增 `PhabricatorWorkerExecuteConduitAPIMethod`。它根据 `taskClass` 和 `data` 构造一个临时的 `PhabricatorWorkerActiveTask`,取得真实 worker 实例并调用 `executeTask()`,然后把结果分类返回: | PHP 执行结果 | 返回给 Go | taskqueue 动作 | | -------------------------------------------- | ------------------- | -------------- | | 正常返回 | `success` | `complete` | | `PhabricatorWorkerYieldException` | `yield` | `yield` | | `PhabricatorWorkerPermanentFailureException` | `permanent-failure` | 永久 `fail` 并归档 | | 其它异常 | `failure` | 临时 `fail`,等待重试 | 临时对象只给 PHP worker 提供 task ID、类型和负载上下文,不写数据库。租约、失败次数与归档仍由 gorge-worker 通过 taskqueue 完成,避免两边同时做账。 ## PHP 侧保留原生路径,但“调用失败再写 SQL”要分后端看 Phorge 保留了原生 SQL 路径,并在几个关键入口检查 `PhabricatorGorgeTaskQueueClient::isConfigured()`: - `PhabricatorWorker::scheduleTask()`:入队; - `PhabricatorWorker::awakenTaskIDs()`:唤醒 yield 的任务; - `PhabricatorWorkerLeaseQuery::execute()`:租约或列表查询; - `PhabricatorWorkerActiveTask`:完成、失败与 yield 回报。 以租约为例,结构很清楚: ```php if (PhabricatorGorgeTaskQueueClient::isConfigured()) { if ($this->skipLease) { return $this->executeListViaGoService(); } return $this->executeLeaseViaGoService(); } // 原生 SQL 路径继续保留。 ``` 这让未配置 Gorge 的安装维持原行为,也让 MySQL 共享表模式具备直接回退能力。 不过完成、失败与 yield 的通知函数在 HTTP 调用失败时,会尝试回到本地 SQL。这个设计在 MySQL 后端有意义,因为服务与 PHP 看到的是同一行;在 Redis 后端,PHP 手里的任务对象只是临时对象,MySQL 里没有对应活跃行,本地 SQL 无法替 Redis 完成归档。Redis 模式应把 taskqueue 调用失败视为需要恢复的基础设施故障,不能依赖 SQL 兜底。 ## 接管开关:为什么要把 `phd.taskmasters` 设为 0 当 `GORGE_TASKQUEUE_URI` 非空时,entrypoint 会写入 `gorge.taskqueue.uri/token`,并在默认开关下执行: ```bash bin/config set phd.taskmasters 0 ``` 这里需要纠正一个容易流传的说法:**旧 taskmaster 和 gorge-worker 并存,并不意味着每个任务必然执行两次。**两边如果都正确遵循同一套条件更新和租约规则,通常会竞争并瓜分任务;同一行只能由一次带条件的更新成功租用。 仍然应该关掉旧消费者,理由是所有权必须唯一: - 两套运行时会争抢任务,无法预测某个 task class 最终在哪个环境执行; - 监控、容量规划和故障定位被拆成两套; - 灰度结果会被竞争比例污染,难以证明新 worker 覆盖了预期流量; - 租约过期、执行超时或结果回报失败时,任何“至少一次”队列都可能重试并产生重复副作用,这与是否同时部署两种消费者是另一层问题。 因此,`phd.taskmasters=0` 用来完成明确的消费权交接,解决的是所有权问题。r4 也提供 `GORGE_TASKQUEUE_DISABLE_PHD_TASKMASTER=0` 作为灰度逃生开关;使用它时需要接受混合消费,不宜当作常态扩容方案。 建议按下面顺序切换: 1. 启动 taskqueue,并验证实际后端; 2. 启动 worker 与 conduit,确认 worker 能领取并回报测试任务; 3. 写入 Phorge 的 taskqueue 配置,让新入队和队列操作改走 HTTP; 4. 把 `phd.taskmasters` 设为 0,确认旧消费者退出; 5. 观察活跃队列、归档队列与业务副作用,再宣布接管完成。 其中第 3、4 步应在同一个切换窗口内连续完成,尽量缩短新旧消费者同时竞争队列的时间。先把新链路验证好,再交接消费权,既不会在准备阶段留下无人处理的任务,也不会把混合消费长期当成正常状态。 这正是系列第十三篇反复强调的区别:服务“在运行”“可达”“就绪”“已被选择”和“业务生效”是五个不同状态。 ## 启动顺序:避免循环依赖,也别高估健康探针 `{namespace}_worker` 数据库由 Phorge 的 `bin/storage upgrade` 创建。若让 Phorge 等 taskqueue 健康,而 taskqueue 又要先连到尚未创建的数据库,就会形成首次启动死锁。 r4 的 Compose 因此这样编排: - taskqueue 先等 MySQL 与授权初始化完成; - Phorge 对 taskqueue、worker 使用 `service_started`,不等它们健康; - worker 可以等 taskqueue 的 `service_healthy`; - taskqueue 和 Redis 客户端都采用惰性连接,后端暂时不可达时进程仍可启动。 这里还要再收紧一次表述:taskqueue 的 `/readyz` 对 MySQL 只执行 `PingContext()`,不会查询 `worker_activetask` 是否存在。因此: | 探针 | 能证明什么 | 不能证明什么 | | -------------------------- | -------------------- | ------------------------- | | taskqueue `/healthz` | HTTP 进程活着 | 后端可用、表已创建 | | taskqueue `/readyz` | MySQL 或 Redis 可以连接 | 三张 worker 表一定存在、完整队列操作已通过 | | worker `/readyz` | 与 `/healthz` 等价,进程活着 | 能连 taskqueue、正在消费任务 | | worker `/api/worker/stats` | 当前进程的处理计数 | 全队列状态;计数重启后会清零 | 所以首次启动后,至少要真实跑一次 `enqueue → lease → complete`,才能证明 schema、凭据、命名空间与队列协议同时正确。只看到两个绿色健康检查,还不足以宣告队列已接管。 ## Conduit 回调修了 Host 匹配,但信任边界仍需闭合 worker 把任务委派回 Phorge 时,gorge-conduit 默认访问 `http://phorge:80`。上游请求带 `Host: phorge`,若它不在 `phabricator.base-uri`、`production-uri` 或 `allowed-uris` 中,Phorge 会返回 `Site Not Found`。r4 的 entrypoint 会从 `GORGE_CONDUIT_UPSTREAM_URL` 解析主机并写入 `phabricator.allowed-uris`,解决这条内网回调链路。 这里有两个部署注意点。 第一,r4 生成的是一个只包含网关上游的数组,并通过 `bin/config set --stdin phabricator.allowed-uris` 写回。已有环境若本来维护了多个 allowed URI,应先合并旧值,避免启动时覆盖原列表。 第二点影响更大:`worker.execute` 当前设置了 `shouldRequireAuthentication() = false` 和 `shouldAllowUnguardedWrites() = true`。代码的安全前提是“调用只能经过带 `X-Service-Token` 的 gorge-conduit”,但这个前提并没有由该 Conduit method 自己验证。`allowed-uris` 只解决 Host 匹配,无法认证调用者。 因此,生产环境至少要阻止外部绕过网关,直接访问 Phorge 的这个内部入口。更稳妥的后续实现,是让 `worker.execute` 自己验证独立的服务凭据,或使用受限的 Conduit 身份,并对可执行的 `taskClass` 做白名单控制。机器间调用没有 CSRF 风险,仍然需要身份认证;这两件事不能混在一起。 ## 验证闭环:不要只数容器,要数状态转移 这次验证分成三层。 第一层是 Go 的单元与契约测试。taskqueue 通过可注入的内存 Store 验证 HTTP 契约,MySQL 和 Redis 分别测试租约、归档与错误路径;worker 则验证类过滤、fallback、结果分类与回报重试。平台分层测试也加入 `internal/taskqueue` 和 `internal/worker`,防止公共层反向依赖业务域。 第二层是仓库里的 `tests/e2e/taskqueue.sh`,跑通真实 HTTP 的入队、租约、完成等状态转移。它比单纯访问健康探针更有价值,因为只有真实操作才会碰到库名、表结构与字段约束。 第三层是在完整 Phorge + Gorge 栈上,用 `bin/worker flood` 注入真实的 `PhabricatorTestWorker`: - taskqueue 能持续领取任务; - gorge-worker 经 gorge-conduit 调用 `worker.execute`; - 停止注入后,`worker_activetask` 下降、`worker_archivetask` 上升; - 成功归档行的 `result=0`,字段名是 `failureCount`,不是 `failedCount`; - `lisk_counter` 中 `worker_activetask` 的计数持续递增,证明 Go 与 PHP 共用同一 ID 来源; - 清空 taskqueue 配置后,MySQL 模式能重新走原生 SQL 入队。 最有说服力的证据来自状态转移与所有权:谁入队、谁领取、谁执行、谁归档,每一步都能对上。 ## 其他 把 Phorge 的工作队列搬到 Go,核心在于讲清三种所有权:队列状态由谁保存,租约由谁授予,业务代码由谁执行。 r4 已经搭出了一个实用的过渡结构:Phorge 仍然产生任务,gorge-taskqueue 统一管理状态,gorge-worker 接管调度,尚未迁移的 task class 再经 Conduit 回到 PHP。它允许队列运行时先独立出来,而不要求一次性重写所有 worker。 增加两个容器还不足以代表服务化完成。完整交接还要满足:旧消费者明确退出、后端回退边界可解释、真实任务完成一次闭环、镜像发布件确实存在,以及内部执行入口没有可绕过的信任缺口。 ## 最后 对这类“在旧系统旁边接入新实现”的改造而言,代码只提供新路径;租约、表结构、启动顺序、身份验证和回滚方式共同决定这条路径能否长期可信。 --EOF