消息队列配置
本文引用的文件 - config/vhttpd.example.toml - config/vhttpd.multi.example.toml - config/vhttpd.vjsx.example.toml - examples/config/hello-v2.toml - examples/config/feishu-paseo.toml - examples/config/mcp-relay-local-v2.toml - examples/config/stream-dispatch.toml - src/server_lifecycle/runtime_config.v - src/config/v1_plan_compat.v - src/config/v2_plan_compiler.v - src/http_pipeline_runtime.v - src/relay/delivery.v - src/relay/outbound_test.v - src/relay/delivery_test.v - src/http_response_relay_delivery_test.v - src/upstream/transport/worker_protocol.v - articles/11-observability.md - tests/e2e/config_acceptance_test.sh
目录
简介
本文件面向在 vhttpd 中构建“类消息队列”能力的工程团队,聚焦以下目标: - 队列后端选择:内存队列、进程内队列(Worker 队列)、基于 WebSocket 的分布式中继(Relay) - 消息路由与调度:HTTP/WS 管道匹配、亲和性路由、优先级策略 - 消费者组与确认:WebSocket 亲和与派发、完成模式(accepted/wait)、超时控制 - 持久化与清理:当前仓库未提供内置消息持久化;建议通过外部系统或应用层实现 - 可靠性保证:幂等键、重试与退避、错误分类与可观测性 - 性能参数:批量大小、异步处理、资源限制 - 监控告警:Admin API、事件日志、关键指标采集
说明:vhttpd 并非传统意义上的 MQ 服务器,但通过 Worker 队列、WebSocket 亲和派发与 Relay 能力,可以组合出高吞吐、低延迟的消息通道。
项目结构
与“消息队列”相关的配置与实现主要分布在如下位置: - 配置样例:config/.toml、examples/config/.toml - 运行时计划编译与校验:src/config/v2_plan_compiler.v、src/config/v1_plan_compat.v - 运行时参数解析与注入:src/server_lifecycle/runtime_config.v - HTTP 管道与最大请求体限制:src/http_pipeline_runtime.v - 中继(Relay)投递与完成策略:src/relay/delivery.v、outbound_test.v、delivery_test.v、http_response_relay_delivery_test.v - Worker 协议帧定义:src/upstream/transport/worker_protocol.v - 可观测性与性能调优参考:articles/11-observability.md、tests/e2e/config_acceptance_test.sh
graph TB
A["配置文件<br/>TOML"] --> B["V2 计划编译器<br/>v2_plan_compiler.v"]
A --> C["V1 兼容映射<br/>v1_plan_compat.v"]
B --> D["运行时配置<br/>runtime_config.v"]
C --> D
D --> E["HTTP 管道执行<br/>http_pipeline_runtime.v"]
D --> F["Relay 投递与完成策略<br/>relay/delivery.v"]
E --> G["应用/引擎适配器"]
F --> H["外部载体(WebSocket/HTTP)"]
图表来源 - src/config/v2_plan_compiler.v:698-745 - src/config/v1_plan_compat.v:504-544 - src/server_lifecycle/runtime_config.v:89-112 - src/http_pipeline_runtime.v:37-62 - src/relay/delivery.v:1-39
章节来源 - config/vhttpd.example.toml:1-67 - config/vhttpd.multi.example.toml:1-73 - config/vhttpd.vjsx.example.toml:1-37 - examples/config/hello-v2.toml:1-31 - examples/config/feishu-paseo.toml:1-64 - examples/config/mcp-relay-local-v2.toml:1-36 - examples/config/stream-dispatch.toml:1-30
核心组件
- Worker 进程内队列
- 用于将请求排队到 PHP/VJSX 工作进程,支持容量与等待超时控制
- 关键参数:pool_size、queue_capacity、queue_timeout_ms、max_requests、read_timeout_ms
- WebSocket 亲和与派发
- 基于 key 的亲和路由,可将同一会话的请求稳定派发到同一 lane/worker
- 支持从 header/query/app 计算 key,以及 fallback/reject 策略
- Relay 中继(分布式通道)
- 以 carrier=websocket 为载体,在节点间转发消息
- 支持 completion_mode=accepted/wait 与超时控制,返回 202/209/5xx 语义
- HTTP 管道与路由
- 通过 listeners/adapters/pipelines 声明式编排
- 支持路径匹配、状态码直返、最大请求体限制等
章节来源 - src/server_lifecycle/runtime_config.v:89-112 - src/server_lifecycle/runtime_config.v:146-171 - src/config/v1_plan_compat.v:504-544 - src/relay/delivery.v:1-39 - src/http_pipeline_runtime.v:37-62
架构总览
下图展示了“生产者 -> 管道 -> 适配器 -> 引擎/工作者 -> 中继/外部载体”的整体流程,并标注了关键配置点。
sequenceDiagram
participant P as "生产者"
participant L as "监听器(listener)"
participant R as "路由(pipeline)"
participant A as "适配器(adapter)"
participant W as "工作进程(worker)"
participant RL as "中继(relay)"
participant C as "消费者/外部系统"
P->>L : "HTTP/WS 请求"
L->>R : "按路径/方法匹配"
R->>A : "命中规则后进入适配器"
alt "本地处理"
A->>W : "入队/派发(受 queue_capacity/timeout 约束)"
W-->>A : "结果/流式帧"
A-->>P : "响应"
else "中继投递"
A->>RL : "准备投递(含 completion_mode/timeout)"
RL-->>A : "ready/accepted 或 unavailable"
A-->>P : "202 accepted 或 503 unavailable"
RL->>C : "通过载体(如 WS)投递"
C-->>RL : "完成帧(可选)"
RL-->>A : "完成回调(可选)"
A-->>P : "209 done 或 504 超时"
end
图表来源 - examples/config/hello-v2.toml:10-31 - examples/config/mcp-relay-local-v2.toml:18-36 - src/relay/delivery.v:1-39 - src/http_response_relay_delivery_test.v:1-247
详细组件分析
队列后端配置
- 内存/进程内队列(Worker 队列)
- 适用场景:单实例或同机多进程的高吞吐短任务
- 关键参数
- pool_size:并发工作进程数
- queue_capacity:队列容量
- queue_timeout_ms:入队等待超时
- max_requests:生命周期内最大请求数(配合重启避免内存泄漏)
- read_timeout_ms:读超时(长连接/AI 流式需增大)
- 参考样例
- 持久化存储
- 仓库未提供内置消息持久化;如需持久化,建议在应用层落盘或通过外部系统(数据库/对象存储)实现
- 分布式队列
- 使用 Relay + WebSocket 作为载体,实现跨实例的消息分发与回传
- 参考样例
章节来源 - config/vhttpd.example.toml:19-26 - examples/config/stream-dispatch.toml:9-15 - tests/e2e/config_acceptance_test.sh:2233-2256 - examples/config/mcp-relay-local-v2.toml:18-27 - examples/config/feishu-paseo.toml:50-64
消息路由配置
- 路由规则
- 使用 pipelines 声明 ingress/egress 与 match.paths
- 参考样例
- 优先级调度
- WebSocket 亲和支持 priority(high/low/数字),结合 should_pin_lane 决定是否钉住 lane
- 参考实现
- 死信队列
- 仓库未提供内置死信队列;可在应用层对失败消息进行重放或归档
章节来源 - examples/config/hello-v2.toml:26-31 - src/executor/inproc_vjsx_websocket_policy.v:67-73 - src/executor/inproc_vjsx_websocket_affinity_runtime.v:148-185
消费者组配置
- 负载均衡与亲和
- 通过 websocket_affinity.source/key/scope/fallback 控制亲和来源与作用域
- 参考样例
- 消费确认
- 使用 completion_mode=accepted 或 wait,配合 completion_timeout_ms 控制等待时长
- 参考测试与实现
章节来源 - examples/config/feishu-paseo.toml:56-60 - src/config/v1_plan_compat.v:504-528 - src/relay/outbound_test.v:28-43 - src/http_response_relay_delivery_test.v:209-247
消息持久化配置
- 现状
- 仓库未提供内置消息持久化机制
- 建议方案
- 应用层落盘(文件/DB)
- 通过外部消息系统(Kafka/RabbitMQ/Redis Streams)对接
- 使用 Relay 的 channel_buffer 做缓冲(非持久化)
章节来源 - examples/config/mcp-relay-local-v2.toml:26
可靠性保证配置
- 事务支持
- 仓库未提供内置事务型消息;需在应用层实现两阶段提交或补偿逻辑
- 幂等性
- 利用 frame/request_id/correlation_id/channel_id 等字段在应用层去重
- 参考
- 重复消息处理
- 在消费者侧基于唯一键去重;必要时结合幂等写入
- 错误分类与可观测性
- 通过 error_class 与 metadata 暴露错误原因,便于告警与追踪
- 参考
章节来源 - src/relay/delivery_test.v:12-34 - src/http_response_relay_delivery_test.v:26-40
性能调优参数
- Worker 池与队列
- pool_size、queue_capacity、queue_timeout_ms、max_requests、read_timeout_ms
- 参考
- 中继缓冲
- channel_buffer 控制内部通道缓冲
- 参考
- 资源限制
- 最大请求体限制(防止过大负载)
- 参考
章节来源 - config/vhttpd.example.toml:19-26 - articles/11-observability.md:622-666 - examples/config/mcp-relay-local-v2.toml:26 - src/http_pipeline_runtime.v:40-47
依赖关系分析
- 配置到运行时的链路
- TOML -> V2 计划编译器 -> 运行时配置 -> 管道/适配器/中继
- 关键耦合点
- v2_plan_compiler 对适配器语义进行校验(如 relay-delivery 必须指定 target)
- runtime_config 将 worker 队列参数注入到 AppRuntimeBuildConfig
- http_pipeline_runtime 在执行前检查请求体大小并短路返回
classDiagram
class V2计划编译器 {
+校验适配器语义
+生成运行时计划
}
class 运行时配置 {
+worker_max_requests
+worker_queue_capacity
+worker_queue_timeout_ms
}
class HTTP管道执行 {
+最大请求体检查
+状态码直返
}
class 中继投递 {
+completion_mode
+completion_timeout_ms
}
V2计划编译器 --> 运行时配置 : "生成"
运行时配置 --> HTTP管道执行 : "注入参数"
运行时配置 --> 中继投递 : "注入参数"
图表来源 - src/config/v2_plan_compiler.v:698-745 - src/server_lifecycle/runtime_config.v:89-112 - src/http_pipeline_runtime.v:37-62 - src/relay/delivery.v:1-39
章节来源 - src/config/v2_plan_compiler.v:698-745 - src/server_lifecycle/runtime_config.v:89-112 - src/http_pipeline_runtime.v:37-62 - src/relay/delivery.v:1-39
性能调优
- Worker 池与队列
- 根据 CPU 核数设置 pool_size;合理设置 queue_capacity 与 queue_timeout_ms 避免雪崩
- 设置 max_requests 定期重启,降低内存泄漏风险
- 超时与限流
- 普通请求 read_timeout_ms 较小;AI 流式请求适当放大
- 通过 Admin API 观察 queue_depth、inflight_requests 等指标
- 中继缓冲
- 调整 channel_buffer 平衡吞吐与内存占用
- 资源限制
- 启用最大请求体限制,避免异常大负载拖垮服务
章节来源 - articles/11-observability.md:622-666 - tests/e2e/config_acceptance_test.sh:3144-3163 - examples/config/mcp-relay-local-v2.toml:26 - src/http_pipeline_runtime.v:40-47
故障排查指南
- 队列满导致拒绝
- 现象:返回 503,error_class=worker_queue_full
- 定位:查看 admin/workers 与 admin/runtime 中的 queue_depth
- 解决:提升 queue_capacity、优化处理耗时、扩容 pool_size
- 中继不可用
- 现象:返回 503,error_class=relay_carrier_unavailable
- 定位:检查 relays 配置与载体健康
- 完成等待超时
- 现象:返回 504,error_class=relay_response_missing
- 定位:核对 completion_timeout_ms 与下游处理耗时
- 事件日志
- 开启 event_log,结合 e2e 脚本中的 wait_event_contains 辅助定位
章节来源 - tests/e2e/config_acceptance_test.sh:3144-3163 - src/http_response_relay_delivery_test.v:26-40 - src/http_response_relay_delivery_test.v:223-247 - config/vhttpd.example.toml:12-14
结论
- vhttpd 通过 Worker 队列、WebSocket 亲和派发与 Relay 能力,能够组合出高性能的消息通道
- 当前仓库未提供内置持久化与死信队列,建议通过应用层或外部系统补齐
- 通过合理的队列与中继参数、完善的错误分类与可观测性,可获得良好的可靠性与可运维性
附录:完整配置示例与监控告警
- 最小可用示例(HTTP 管道)
- examples/config/hello-v2.toml:1-31
- 多站点示例(PHP/VJSX)
- config/vhttpd.multi.example.toml:44-73
- 中继示例(WebSocket 载体)
- examples/config/mcp-relay-local-v2.toml:18-36
- 亲和与派发示例
- examples/config/feishu-paseo.toml:56-60
- 监控告警
- Admin API 指标:/admin/runtime、/admin/workers
- 事件日志:files.event_log
- 参考