
WeKnora 文档入库流水线 Worker Pool 治理六池拓扑、并发预算与容量调优实战【免费下载链接】WeKnoraOpen-source LLM knowledge platform: turn raw documents into a queryable RAG, an autonomous reasoning agent, and a self-maintaining Wiki.项目地址: https://gitcode.com/GitHub_Trending/we/WeKnoraWorker 并发是调度预算不是模型配额、DocReader 容量、向量库或数据库连接上限的替代品。本文以 WeKnora 文档入库流水线为背景完整讲解其按阶段保底 弹性共享的 Worker Pool 治理模型六大独立池的拓扑与队列归属、全部可配置项与环境变量、容量分层思路、基于任务到达率的容量估算公式以及运行时队列仪表盘的使用方法帮助你正确地为生产环境做并发调优。为什么需要 Worker Pool 治理WeKnora 的文档入库是一个典型的解析 → 后处理 → 富化扇出流水线一篇文档解析完成后会派生摘要生成、问题生成按 chunk 批量扇出、每 chunk 的图谱实体抽取、多张图片的多模态处理等大量子任务。如果所有任务共用一个 worker 池长耗时的 DocReader 调用会堵住轻量任务批量导入会饿死交互式聊天上传Wiki 生成洪峰也可能拖慢用户面解析。为此 WeKnora 采用保证型guaranteed分阶段 Worker Pool 弹性elastic共享池的设计每个阶段拥有独立、硬隔离的并发预算共享池在出现积压的一侧借用空闲容量。需要强调的是这些并发数值只是每服务实例最多同时容纳多少个任务 handler的调度预算不能替代以下资源治理模型配额治理provider 并发、RPM、TPMDocReader 的处理能力向量库、对象存储、Postgres 的自身限制本机 CPU / 内存资源。拓扑总览六大 Worker Pool 与队列归属默认拓扑由下表给出来源internal/types/task.go 中的常量与默认值Pool默认并发每实例队列职责Core8default文档解析与手工重解析的保底保障Post-process2postprocess解析收尾与富化扇出调度Enrichment12summary、multimodal、graph、question模型密集型富化的保底保障Maintenance4sync、low数据源同步、批处理、移动/删除/清理Shared6Core 与 Enrichment 的队列有积压一侧借用的弹性容量Wiki8wiki独立治理的 Wiki 生成上游文档入库相关总并发保持默认每实例 328 2 12 4 6由WorkerPoolConcurrency.UpstreamTotal()计算internal/types/task.go。Wiki 独立于该总数它拥有自己的 asynq server 与并发预算WEKNORA_WIKI_ASYNQ_CONCURRENCY默认 8与解析流水线互为隔离Wiki 洪峰不会拖慢用户面解析上传高峰期解析任务也不会饿死 Wiki。队列注册表单一事实来源所有队列的名称、归属池、权重与承载任务类型集中定义在queueDefinitionsinternal/types/task.go供 worker server 构造与运行时巡检共同消费保证运维界面展示的调度权重不会与实际 server 配置漂移Pool队列池内权重共享权重承载任务类型Coredefault13document:process、manual:processCorechat_attachment33temporary_document:process会话聊天附件解析Post-processpostprocess1—knowledge:post_processEnrichmentsummary22summary:generation、datatable:summary、knowledge:auto_tagEnrichmentmultimodal11image:multimodalOCR VLM CaptionEnrichmentgraph11chunk:extract按 chunk 图谱抽取Enrichmentquestion11question:generation按 chunk 批次扇出Enrichmentmemory11memory:extract防抖后的长期记忆后台抽取Maintenancesync2—datasource:syncMaintenancelow1—faq:import、kb:clone、index:delete、kb:delete、批量删除/重解析、knowledge:moveWikiwiki1—wiki:ingest、wiki:finalize两个值得注意的细节聊天附件优先chat_attachment在 Core 池内的权重3高于default1确保大规模知识库批量导入不会让交互式聊天上传排队在其后。维护队列保留物理旧名lowQueueMaintenance常量故意保留旧 Redis 物理名low使滚动升级期间旧版本入队的任务仍可被消费internal/types/task.go。为何 Post-process 与 Maintenance 不共享弹性容量Post-process 拥有独立物理队列轻量扇出任务不会被长时间 DocReader 调用堵在后面保证解析收尾的低延迟。Maintenance 被刻意排除在共享池之外它的长任务KB 克隆、删除、移动等可能长期钉住本应服务于用户面流水线的弹性容量。从代码看QueueWeightsForSharedPool()只聚合SharedWeight 0的队列即 Core 与 Enrichment 的全部队列internal/types/task.go。原子出队保证asynq 的 dequeue 是原子的因此专用 server 与共享 server 可以安全地订阅同一批 Core/Enrichment 队列——同一个任务仍然只会被一个 worker 执行不会重复处理。共享池的构造逻辑见 internal/router/task.goNewSharedAsynqServer以QueueWeightsForSharedPool()作为队列配置订阅全部共享队列。配置项系统设置与环境变量所有 worker 并发设置在System settings页面中管理且修改后必须重启服务进程方可生效。每个配置项的注册信息类型、环境变量、默认值、RequiresRestart标记与中文说明定义在 internal/application/service/system_setting.go配置键环境变量默认值说明asynq.core_concurrencyWEKNORA_ASYNQ_CORE_CONCURRENCY8文档解析、手工重解析等核心任务的每实例保底并发可额外使用共享弹性池asynq.postprocess_concurrencyWEKNORA_ASYNQ_POSTPROCESS_CONCURRENCY2解析收尾的轻量编排与富化扇出专用并发避免被长文档解析阻塞asynq.enrichment_concurrencyWEKNORA_ASYNQ_ENRICHMENT_CONCURRENCY12摘要、图片、图谱、问题生成的每实例保底并发可额外使用共享弹性池asynq.maintenance_concurrencyWEKNORA_ASYNQ_MAINTENANCE_CONCURRENCY4数据源同步、批处理、移动/删除/清理的每实例保底并发与用户面流水线硬隔离asynq.shared_concurrencyWEKNORA_ASYNQ_SHARED_CONCURRENCY6核心解析与内容富化共用的每实例弹性并发空闲容量由有积压的一侧借用asynq.wiki_concurrencyWEKNORA_WIKI_ASYNQ_CONCURRENCY8Wiki 生成专用池的 worker 并发数与文档解析池相互隔离旧聚合配置已退役旧的聚合配置asynq.concurrency/WEKNORA_ASYNQ_CONCURRENCY已经退役。设置了它的部署必须迁移到上述显式池配置。系统设置列表接口在读取时会显式删除asynq.concurrency这一旧键见 internal/application/service/system_setting.go使持久化的旧行既不会被当作生效的运行时控制项也会从 System settings 页面隐藏。配置解析与合法性校验并发值的解析集中在ResolveWorkerPoolConcurrencyinternal/types/task.go它同时负责服务端构造与运行时 API 的取值避免两处逻辑漂移。解析采用键/环境变量/默认值三层的 3-tier 读取且每个值都有最小值保护任何池的取值若小于 1一律回退到默认值positive辅助函数。设置校验层面同样在 system_setting.go 强制concurrency must be at least 1非法值无法写入。以 Core 池为例server 启动时打印的日志形如asynq core-pool server starting with concurrency8 total_upstream32 redis_op_timeout500ms容量分层并发预算、模型配额与基础设施限制各自独立调节调优时必须把以下三个层次分开考虑每层独立调节Worker 并发控制每个服务实例容纳的任务 handler 数量本文主题。模型配额治理控制跨副本与配额组的 provider 并发、RPM、TPM。后台任务文档入库/富化对单个模型的默认并发上限由model.max_concurrency/WEKNORA_MODEL_MAX_CONCURRENCY控制默认 32按模型 ID 全副本共享单个模型还可在模型管理中覆盖自己的max_concurrency。它只影响后台任务不影响交互式对话且修改后立即生效、无需重启internal/application/service/system_setting.go。DocReader、向量库、对象存储、Postgres、本地 CPU/RAM各自保有独立资源限制。关键判断如果模型限流器的等待在增长而 worker 队列仍然繁忙那么增加 worker 只会制造更多等待者。此时应该提高 provider 配额或者降低 worker 准入数而不是继续加 worker。容量估算与调优策略估算公式对每个池用以下公式估算峰值所需 worker 数required workers ceil(peak task arrival rate * mean task runtime / 0.70)其中的任务到达率应为扇出之后的数值一篇已解析文档可以派生 1 个摘要任务、多个问题批次、每个 chunk 1 个图谱任务、多张图片的多个多模态任务。队列长度不是容量信号——队列积压只能说明到达超过消费不能直接反推该加多少 worker。何时扩容 / 何时收缩扩容仅当某个池的下游依赖还有余量headroom、而该池积压任务的最老等待年龄oldest pending age在持续增长时才增加该池的并发。重分配当 Core 长期低利用率而富化队列summary/multimodal/graph/question在增长时把预算重新分配给 Enrichment 池。收缩当 DocReader 的 CPU、内存或延迟成为瓶颈时应降低 Core 的准入并发而不是让它继续排队轰炸解析服务。运行时仪表盘指标运行时仪表盘Runtime Dashboard提供如下观测数据用于支撑上述决策每个池的每实例已配置并发configured concurrency per instance在线服务实例数live server instance count由活跃 asynq 心跳聚合得到的集群容量cluster capacity活跃 worker 数与利用率utilization Active / ClusterCapacity每个队列的积压、重试、死信数量与最老待处理任务年龄模型并发与限流器等待统计。从实现看RuntimeQueuesResponse把每池的配置并发、实例数、集群容量、活跃数与利用率组装在一起并附带每个队列的QueueStatpending/active/scheduled/retry/archived/completed、今日 processed/failed、是否暂停、最老待处理延迟毫秒数、近似内存占用见 internal/handler/system.go。聚合逻辑通过比对 server 心跳携带的队列权重来识别每个心跳属于哪个逻辑池再累加实例数、集群容量与活跃 workerinternal/handler/system.go因此配置的每实例并发不会与实际的集群容量混淆。数据接口为GET /system/admin/runtime/queues仅系统管理员可见。注意Lite 模式无 Redis/asynq下该接口返回availablefalse前端据此渲染当前部署不可用状态而不是空表。源码级验证拓扑约束由测试锁定并发拓扑与配置解析有专门的测试保障internal/types/task_queue_test.go可作为理解与回归的锚点TestQueueDefinitionsAreUniqueAndConsumable队列定义唯一且可消费TestChatAttachmentQueueIsIsolatedAndPrioritized聊天附件队列隔离且优先权重 3TestQueueMaintenanceKeepsLegacyPhysicalName维护队列保留物理旧名low保证滚动升级兼容TestEveryAsynqTaskTypeHasADeclaredQueue每个任务类型都有声明的队列防止新任务类型漏挂队列导致无法被消费TestDefaultWorkerPoolConcurrencyIsExplicitBudget默认并发是显式预算合计 32TestResolveWorkerPoolConcurrencyFallsBackPerPool逐池回退逻辑正确。总结WeKnora 的 Worker Pool 治理可以概括为三条原则按阶段硬隔离Core、Post-process、Enrichment、Maintenance、Wiki 各有独立 asynq server 与并发预算互不挤占Shared 池只在 Core/Enrichment 有积压时弹性借用空闲容量。并发只是调度预算Worker 并发不替代模型配额与各基础设施DocReader、向量库、数据库的容量治理调优时必须三层独立观察。用数据调优基于扇出后的任务到达率与平均运行时间做容量估算结合运行时仪表盘的积压年龄、利用率与限流器等待决定扩容、重分配还是收缩。配置变更均在 System settings 页面完成或直接设置WEKNORA_*环境变量修改后需重启服务进程并且旧聚合配置asynq.concurrency已不再生效务必迁移到六个显式池配置。更完整的异步任务系统架构可继续阅读 website-docs/02-architecture/05-async-tasks.md。【免费下载链接】WeKnoraOpen-source LLM knowledge platform: turn raw documents into a queryable RAG, an autonomous reasoning agent, and a self-maintaining Wiki.项目地址: https://gitcode.com/GitHub_Trending/we/WeKnora创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考