07. Queue 与 Scaling:main、worker 和 Redis
7.1 两种执行模式
单进程模式中,HTTP 服务和执行器在同一进程完成。queue mode 中,main 负责接收请求、管理工作流和入队,worker 从 Redis/Bull 领取 job 并执行。
7.2 WorkflowRunner 是边界对象
packages/cli/src/workflow-runner.ts 接收 execution 请求,决定:
- 是否是 manual execution。
- 是否应该进入 queue mode。
- 是否需要等待、优先级、重试或恢复。
- 执行数据如何构造并传给 core。
- execution 开始/结束后如何保存和通知。
它不是执行算法本身,而是“平台执行编排器”。这层把本地执行和 queue execution 统一成相近的输入契约。
7.3 Worker job 的边界
worker 取出的 job 需要包含足够的信息:workflow id 或序列化 workflow、execution mode、开始节点/触发数据、retry/priority、webhook response metadata,以及恢复所需的上下文。job 太小,worker 还要同步访问 main;job 太大,会增加 Redis 和序列化开销。
packages/cli/src/scaling/job-processor.ts 和 worker-server.ts 负责把 Bull job 转成 core 执行调用,再把结果映射回 execution service、push 和 webhook relay。
7.4 并发和背压
规模化不等于无限开 worker。n8n 同时存在:
- queue worker concurrency。
- 单实例/全局 execution concurrency limit。
- 节点级外部 API 限速或重试。
- 数据库连接池和 Redis 连接。
- task runner 的并发。
packages/cli/src/concurrency 中的 ConcurrencyControlService、capacity reservation 和 queue 共同实现背压:只有获得 capacity 的 execution 才继续推进,其他 execution 在队列或等待状态停留。
7.5 多主与事件传播
多 main 实例时,激活状态、停止 execution、worker 状态和 webhook response 不能只存在本地内存。packages/cli/src/scaling/pubsub 使用发布/订阅事件传播跨进程变化,Redis lock 和 leader election 避免多个主实例重复承担全局任务。
7.6 失败与恢复
队列系统要处理:
- worker 进程崩溃:Bull 认为 job stalled,重试或转失败。
- main 崩溃:恢复 enqueued execution,避免丢失用户请求。
- execution 停止:停止信号通过 scaling/pubsub 到达实际 worker。
- webhook 已收到但 worker 未完成:response relay 需要保存关联关系。
- 重试导致外部副作用:节点和工作流设计需要幂等策略。
因此 queue mode 是可靠性协议,而不只是把函数丢给 Redis。
7.7 设计取舍
main/worker 分离提高吞吐和隔离性,但让实时反馈、取消、Webhook response 和错误回传需要跨进程协议。
Bull/Redis提供成熟的 job 生命周期和 stalled 检测,但 Redis 不是最终事实来源;execution 状态和业务数据仍要落到数据库/二进制存储。
显式并发控制保护数据库和外部 API,也让系统在高峰时保持可预测;代价是用户看到 queued/waiting,需要 UI 和运维指标解释“没跑”与“跑失败”的区别。