From a4046c2b81d84734d9088b65a81be76423d153cf Mon Sep 17 00:00:00 2001 From: wangbo Date: Sat, 1 Aug 2026 23:24:07 +0800 Subject: [PATCH] =?UTF-8?q?perf(worker):=20=E8=A7=A3=E8=80=A6=E5=87=86?= =?UTF-8?q?=E5=85=A5=E8=B0=83=E5=BA=A6=E4=B8=8E=E8=BF=9C=E7=AB=AF=E6=89=A7?= =?UTF-8?q?=E8=A1=8C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 原因:线上 P24 同构验收确认两台香港 Worker 的全局容量为 48,但跨地域逐行准入事务每轮只能形成 8 个活跃租约,队列最老等待超过 15 分钟。\n\n影响:新增可配置的异步准入 dispatcher 角色;生产两地 API 负责准入和过期回收,Worker 仅执行 River job。当前主库同站点 API 可低延迟填满容量,主库切换后另一地 API 通过既有数据库锁安全接管;未配置环境变量时保持原 Worker 一体化行为。\n\n验证:Go 全量测试、go vet、gofmt、真实 PostgreSQL 1000 任务双进程回归、前端 lint/test/build、Kubernetes server-side dry-run、Compose、ShellCheck、cluster/manual release tests 全部通过。 --- README.md | 2 +- apps/api/internal/config/config.go | 18 +++++++ apps/api/internal/config/config_test.go | 18 +++++++ ...sync_worker_acceptance_integration_test.go | 52 ++++++++++--------- apps/api/internal/httpapi/server.go | 3 ++ apps/api/internal/runner/queue_worker.go | 15 +++++- apps/api/internal/runner/service.go | 1 + .../easyai-ai-gateway-cluster-release | 5 ++ deploy/kubernetes/production/application.yaml | 8 +++ docs/operations/production-acceptance.md | 1 + 10 files changed, 96 insertions(+), 27 deletions(-) diff --git a/README.md b/README.md index 7beb1c1..9595f99 100644 --- a/README.md +++ b/README.md @@ -168,7 +168,7 @@ AI_GATEWAY_DATABASE_URL=postgresql://easyai:easyai2025@localhost:5432/easyai_ai_ 如果现有 `easyai-pgvector` 没有把 `5432` 映射到宿主机,就需要补端口映射,或者把 AI Gateway 后端容器化后接入同一个 `easyai` Docker network。 -异步队列 worker 不使用固定业务并发。服务分别汇总启用平台模型和活跃用户组的有效 `concurrent` 策略,采用两者中更严格的集群容量,默认每 5 秒在线调整 River 执行容量。策略解析兼容历史 `platformLimits/modelLimits.max_concurrent_requests`,运行时统一转换为 `rules`。`AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT` 默认 `2048`,限制策略推导出的集群目标;`AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT` 默认 `32`,按单 Worker 的内存安全容量限制实例分配,即使其他实例失活也不会突破。Worker 的 `AI_GATEWAY_DATABASE_MAX_CONNS` 必须高于实例执行容量,为心跳、选主、租约续期和健康检查保留连接;生产环境按每实例执行容量 `24` 配置连接池上限 `32`。`AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE` 默认 `8`,把短任务准入、容量租约和唯一 River job 合并为有界原子微批次,摊薄跨地域同步提交,同时禁止重新形成整窗大事务。异步任务首次排队时会持久化已经鉴权的候选和 admission scope 快照;Worker 批量读取快照并仅刷新动态执行容量,只有显式重新选路才重新水化媒体和计算候选,避免队列越大、重复路由查询越多。`AI_GATEWAY_DATABASE_MIN_IDLE_CONNS` 控制启动预热连接数,默认 `0`,生产 K3s 配置为 `4`;`AI_GATEWAY_DATABASE_MAX_CONN_IDLE_SECONDS` 生产配置为 `300` 秒。平台模型和用户组的 PostgreSQL concurrency lease 仍是业务并发真值。可通过 `AI_GATEWAY_ASYNC_WORKER_REFRESH_INTERVAL_SECONDS` 调整刷新周期。 +异步队列 worker 不使用固定业务并发。服务分别汇总启用平台模型和活跃用户组的有效 `concurrent` 策略,采用两者中更严格的集群容量,默认每 5 秒在线调整 River 执行容量。策略解析兼容历史 `platformLimits/modelLimits.max_concurrent_requests`,运行时统一转换为 `rules`。`AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT` 默认 `2048`,限制策略推导出的集群目标;`AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT` 默认 `32`,按单 Worker 的内存安全容量限制实例分配,即使其他实例失活也不会突破。Worker 的 `AI_GATEWAY_DATABASE_MAX_CONNS` 必须高于实例执行容量,为心跳、选主、租约续期和健康检查保留连接;生产环境按每实例执行容量 `24` 配置连接池上限 `32`。`AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE` 默认 `8`,把短任务准入、容量租约和唯一 River job 合并为有界原子微批次,摊薄同步提交,同时禁止重新形成整窗大事务。异步任务首次排队时会持久化已经鉴权的候选和 admission scope 快照;启用 `AI_GATEWAY_ASYNC_ADMISSION_DISPATCHER_ENABLED=true` 的进程批量读取快照并仅刷新动态执行容量,只有显式重新选路才重新水化媒体和计算候选。生产 K3s 在两地 API 启用 dispatcher、在 Worker 禁用,使当前 PostgreSQL 主库同站点 API 负责低延迟准入,远端 Worker 只执行 River job;数据库任务锁和 scope 锁保证主库切换时另一地 API 可安全接管。未显式配置该变量时保持兼容行为,由执行 Worker 同时运行 dispatcher。`AI_GATEWAY_DATABASE_MIN_IDLE_CONNS` 控制启动预热连接数,默认 `0`,生产 K3s 配置为 `4`;`AI_GATEWAY_DATABASE_MAX_CONN_IDLE_SECONDS` 生产配置为 `300` 秒。平台模型和用户组的 PostgreSQL concurrency lease 仍是业务并发真值。可通过 `AI_GATEWAY_ASYNC_WORKER_REFRESH_INTERVAL_SECONDS` 调整刷新周期。 ## 迁移原则 diff --git a/apps/api/internal/config/config.go b/apps/api/internal/config/config.go index 7d89312..88b0a15 100644 --- a/apps/api/internal/config/config.go +++ b/apps/api/internal/config/config.go @@ -80,6 +80,8 @@ type Config struct { AsyncWorkerInstanceHardLimit int AsyncWorkerRefreshIntervalSeconds int AsyncAdmissionMicrobatchSize int + AsyncAdmissionDispatcherEnabled bool + AsyncAdmissionDispatcherConfigured bool WorkerAutoscalingEnabled bool WorkerReplicasNingbo int WorkerReplicasHongkong int @@ -188,6 +190,10 @@ func Load() Config { AsyncWorkerInstanceHardLimit: envIntValidated("AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT", 32), AsyncWorkerRefreshIntervalSeconds: envIntValidated("AI_GATEWAY_ASYNC_WORKER_REFRESH_INTERVAL_SECONDS", 5), AsyncAdmissionMicrobatchSize: envIntValidated("AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE", 8), + AsyncAdmissionDispatcherEnabled: envValue("AI_GATEWAY_ASYNC_ADMISSION_DISPATCHER_ENABLED") == "true", + AsyncAdmissionDispatcherConfigured: envValue( + "AI_GATEWAY_ASYNC_ADMISSION_DISPATCHER_ENABLED", + ) != "", WorkerAutoscalingEnabled: env("AI_GATEWAY_WORKER_AUTOSCALING_ENABLED", "false") == "true", WorkerReplicasNingbo: envOptionalIntValidated("AI_GATEWAY_WORKER_REPLICAS_NINGBO", 1), WorkerReplicasHongkong: envOptionalIntValidated("AI_GATEWAY_WORKER_REPLICAS_HONGKONG", 1), @@ -463,6 +469,18 @@ func (c Config) RunsAsyncExecutionWorker() bool { } } +// RunsAsyncAdmissionDispatcher keeps the legacy all-in-one/Worker behaviour +// unless a deployment explicitly separates admission coordination from task +// execution. Production enables this on both API sites so the instance nearest +// the current PostgreSQL primary wins the existing database scope locks, while +// remote Worker replicas only execute River jobs. +func (c Config) RunsAsyncAdmissionDispatcher() bool { + if c.AsyncAdmissionDispatcherConfigured { + return c.AsyncAdmissionDispatcherEnabled + } + return c.RunsAsyncExecutionWorker() +} + func (c Config) RunsBackgroundWorkers() bool { switch c.EffectiveProcessRole() { case "api", "capacity-controller": diff --git a/apps/api/internal/config/config_test.go b/apps/api/internal/config/config_test.go index 074ef08..48fbd58 100644 --- a/apps/api/internal/config/config_test.go +++ b/apps/api/internal/config/config_test.go @@ -159,12 +159,30 @@ func TestProcessRolePrecedenceAndCompatibility(t *testing.T) { if cfg.RunsPublicHTTP() || !cfg.RunsAsyncExecutionWorker() || !cfg.RunsBackgroundWorkers() { t.Fatalf("worker role capabilities are inconsistent: %+v", cfg) } + if !cfg.RunsAsyncAdmissionDispatcher() { + t.Fatal("legacy Worker role no longer runs the admission dispatcher") + } t.Setenv("AI_GATEWAY_PROCESS_ROLE", "api") cfg = Load() if !cfg.RunsPublicHTTP() || cfg.RunsAsyncExecutionWorker() || cfg.RunsBackgroundWorkers() { t.Fatalf("api role capabilities are inconsistent: %+v", cfg) } + if cfg.RunsAsyncAdmissionDispatcher() { + t.Fatal("API role unexpectedly runs the admission dispatcher without explicit configuration") + } + + t.Setenv("AI_GATEWAY_ASYNC_ADMISSION_DISPATCHER_ENABLED", "true") + cfg = Load() + if !cfg.RunsAsyncAdmissionDispatcher() || cfg.RunsAsyncExecutionWorker() { + t.Fatalf("API admission-only role is inconsistent: %+v", cfg) + } + t.Setenv("AI_GATEWAY_ASYNC_ADMISSION_DISPATCHER_ENABLED", "false") + t.Setenv("AI_GATEWAY_PROCESS_ROLE", "worker") + cfg = Load() + if cfg.RunsAsyncAdmissionDispatcher() || !cfg.RunsAsyncExecutionWorker() { + t.Fatalf("Worker execution-only role is inconsistent: %+v", cfg) + } t.Setenv("AI_GATEWAY_PROCESS_ROLE", "capacity-controller") cfg = Load() diff --git a/apps/api/internal/httpapi/async_worker_acceptance_integration_test.go b/apps/api/internal/httpapi/async_worker_acceptance_integration_test.go index e4227e7..bcb0edf 100644 --- a/apps/api/internal/httpapi/async_worker_acceptance_integration_test.go +++ b/apps/api/internal/httpapi/async_worker_acceptance_integration_test.go @@ -136,33 +136,37 @@ func TestAsyncWorkerThousandConcurrentSchedulingAcceptance(t *testing.T) { serverCtx, cancelServer := context.WithCancel(ctx) defer cancelServer() server := httptest.NewServer(NewServerWithStores(serverCtx, config.Config{ - AppEnv: "test", - HTTPAddr: ":0", - DatabaseURL: databaseURL, - IdentityMode: "hybrid", - JWTSecret: "test-secret", - BillingEngineMode: "observe", - CORSAllowedOrigin: "*", - ProcessRole: "api", - AsyncQueueWorkerEnabled: true, - AsyncWorkerHardLimit: 1000, - AsyncWorkerInstanceHardLimit: 1000, - AsyncWorkerRefreshIntervalSeconds: 1, + AppEnv: "test", + HTTPAddr: ":0", + DatabaseURL: databaseURL, + IdentityMode: "hybrid", + JWTSecret: "test-secret", + BillingEngineMode: "observe", + CORSAllowedOrigin: "*", + ProcessRole: "api", + AsyncQueueWorkerEnabled: true, + AsyncAdmissionDispatcherEnabled: true, + AsyncAdmissionDispatcherConfigured: true, + AsyncWorkerHardLimit: 1000, + AsyncWorkerInstanceHardLimit: 1000, + AsyncWorkerRefreshIntervalSeconds: 1, }, apiDB, apiCoordinationDB, apiDB, slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelError})))) defer server.Close() workerServer := httptest.NewServer(NewServerWithStores(serverCtx, config.Config{ - AppEnv: "test", - HTTPAddr: ":0", - DatabaseURL: databaseURL, - IdentityMode: "hybrid", - JWTSecret: "test-secret", - BillingEngineMode: "observe", - CORSAllowedOrigin: "*", - ProcessRole: "worker", - AsyncQueueWorkerEnabled: true, - AsyncWorkerHardLimit: 1000, - AsyncWorkerInstanceHardLimit: 1000, - AsyncWorkerRefreshIntervalSeconds: 1, + AppEnv: "test", + HTTPAddr: ":0", + DatabaseURL: databaseURL, + IdentityMode: "hybrid", + JWTSecret: "test-secret", + BillingEngineMode: "observe", + CORSAllowedOrigin: "*", + ProcessRole: "worker", + AsyncQueueWorkerEnabled: true, + AsyncAdmissionDispatcherEnabled: false, + AsyncAdmissionDispatcherConfigured: true, + AsyncWorkerHardLimit: 1000, + AsyncWorkerInstanceHardLimit: 1000, + AsyncWorkerRefreshIntervalSeconds: 1, }, workerDB, workerCoordinationDB, workerDB, slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelError})))) defer workerServer.Close() diff --git a/apps/api/internal/httpapi/server.go b/apps/api/internal/httpapi/server.go index a54991c..45579b8 100644 --- a/apps/api/internal/httpapi/server.go +++ b/apps/api/internal/httpapi/server.go @@ -138,6 +138,9 @@ func NewServerWithStores( logger.Info("asynchronous queue worker disabled for this process") } } + if cfg.RunsAsyncAdmissionDispatcher() { + server.runner.StartAsyncAdmissionDispatcher(ctx) + } if cfg.RunsBackgroundWorkers() { server.runner.StartBillingSettlementWorker(ctx) server.runner.StartTaskHistoryWorkers(ctx) diff --git a/apps/api/internal/runner/queue_worker.go b/apps/api/internal/runner/queue_worker.go index e1d1371..82ff0fd 100644 --- a/apps/api/internal/runner/queue_worker.go +++ b/apps/api/internal/runner/queue_worker.go @@ -274,13 +274,24 @@ func (s *Service) startRiverQueue(ctx context.Context, workerEnabled bool) error return err } go s.refreshAsyncWorkerCapacity(ctx) - go s.dispatchWaitingAsyncAdmissions(ctx) - go s.reapExpiredTaskAdmissions(ctx) go s.recoverOrphanedAsyncRiverJobs(ctx) go s.stopAsyncWorkersOnShutdown(ctx) return nil } +// StartAsyncAdmissionDispatcher starts the lightweight admission/reaping +// coordinator independently from River execution workers. This allows the +// dispatcher to run beside PostgreSQL while Worker pods scale on other nodes; +// the existing database task/scope locks make multiple API-site dispatchers +// safe during primary failover. +func (s *Service) StartAsyncAdmissionDispatcher(ctx context.Context) { + s.asyncAdmissionRunner.Do(func() { + go s.dispatchWaitingAsyncAdmissions(ctx) + go s.reapExpiredTaskAdmissions(ctx) + s.logger.Info("asynchronous admission dispatcher started") + }) +} + func (s *Service) newRiverAsyncExecutionClient(capacity int) (*river.Client[pgx.Tx], error) { workers := river.NewWorkers() if err := river.AddWorkerSafely(workers, &asyncTaskWorker{service: s}); err != nil { diff --git a/apps/api/internal/runner/service.go b/apps/api/internal/runner/service.go index 7e2a98a..4eaa998 100644 --- a/apps/api/internal/runner/service.go +++ b/apps/api/internal/runner/service.go @@ -45,6 +45,7 @@ type Service struct { asyncAdmissionWake chan struct{} admissionTaskWaiters map[string]*admissionTaskWaiter admissionListener sync.Once + asyncAdmissionRunner sync.Once asyncClientFactory func(int) (asyncExecutionClient, error) mediaResultSlots chan struct{} imageNormalizeSlots chan struct{} diff --git a/deploy/kubernetes/easyai-ai-gateway-cluster-release b/deploy/kubernetes/easyai-ai-gateway-cluster-release index 6ecfc99..dd21cbe 100755 --- a/deploy/kubernetes/easyai-ai-gateway-cluster-release +++ b/deploy/kubernetes/easyai-ai-gateway-cluster-release @@ -403,6 +403,7 @@ rollout_worker_site() { "AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT=$AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT" \ "AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT=$AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT" \ "AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE=$AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE" \ + "AI_GATEWAY_ASYNC_ADMISSION_DISPATCHER_ENABLED=false" \ "AI_GATEWAY_DATABASE_MAX_CONNS=$AI_GATEWAY_DATABASE_MAX_CONNS" \ "AI_GATEWAY_DATABASE_CRITICAL_MAX_CONNS=$AI_GATEWAY_DATABASE_CRITICAL_MAX_CONNS" \ "AI_GATEWAY_DATABASE_RIVER_MAX_CONNS=$AI_GATEWAY_DATABASE_RIVER_MAX_CONNS" \ @@ -433,6 +434,8 @@ rollout_worker_site() { -n "$NAMESPACE" --timeout=300s || return 1 [[ $("${kubectl[@]}" get deployment "easyai-worker-$site" -n "$NAMESPACE" \ -o jsonpath='{.spec.template.spec.containers[?(@.name=="worker")].env[?(@.name=="AI_GATEWAY_PROCESS_ROLE")].value}') == worker ]] || return 1 + [[ $("${kubectl[@]}" get deployment "easyai-worker-$site" -n "$NAMESPACE" \ + -o jsonpath='{.spec.template.spec.containers[?(@.name=="worker")].env[?(@.name=="AI_GATEWAY_ASYNC_ADMISSION_DISPATCHER_ENABLED")].value}') == false ]] || return 1 ready_replicas=$("${kubectl[@]}" get deployment "easyai-worker-$site" -n "$NAMESPACE" \ -o json | jq -r '.status.readyReplicas // 0') || return 1 [[ $ready_replicas == "$replicas" ]] || return 1 @@ -592,6 +595,7 @@ rollout_api_capacity_site() { "${kubectl[@]}" set env "deployment/easyai-api-$site" -n "$NAMESPACE" \ "AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT=$AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT" \ "AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT=$AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT" \ + "AI_GATEWAY_ASYNC_ADMISSION_DISPATCHER_ENABLED=true" \ "AI_GATEWAY_DATABASE_MAX_CONNS=$AI_GATEWAY_API_DATABASE_MAX_CONNS" \ "AI_GATEWAY_DATABASE_CRITICAL_MAX_CONNS=$AI_GATEWAY_DATABASE_CRITICAL_MAX_CONNS" \ "AI_GATEWAY_DATABASE_RIVER_MAX_CONNS=$AI_GATEWAY_API_DATABASE_RIVER_MAX_CONNS" \ @@ -637,6 +641,7 @@ rollout_site() { "${kubectl[@]}" set env "deployment/easyai-api-$site" -n "$NAMESPACE" \ "AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT=$AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT" \ "AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT=$AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT" \ + "AI_GATEWAY_ASYNC_ADMISSION_DISPATCHER_ENABLED=true" \ "AI_GATEWAY_DATABASE_MAX_CONNS=$AI_GATEWAY_API_DATABASE_MAX_CONNS" \ "AI_GATEWAY_DATABASE_CRITICAL_MAX_CONNS=$AI_GATEWAY_DATABASE_CRITICAL_MAX_CONNS" \ "AI_GATEWAY_DATABASE_RIVER_MAX_CONNS=$AI_GATEWAY_API_DATABASE_RIVER_MAX_CONNS" \ diff --git a/deploy/kubernetes/production/application.yaml b/deploy/kubernetes/production/application.yaml index 1561035..447d1c2 100644 --- a/deploy/kubernetes/production/application.yaml +++ b/deploy/kubernetes/production/application.yaml @@ -52,6 +52,8 @@ spec: value: api - name: AI_GATEWAY_ASYNC_QUEUE_WORKER_ENABLED value: "false" + - name: AI_GATEWAY_ASYNC_ADMISSION_DISPATCHER_ENABLED + value: "true" - name: AI_GATEWAY_DATABASE_MAX_CONNS value: "32" - name: AI_GATEWAY_DATABASE_CRITICAL_MAX_CONNS @@ -330,6 +332,8 @@ spec: value: api - name: AI_GATEWAY_ASYNC_QUEUE_WORKER_ENABLED value: "false" + - name: AI_GATEWAY_ASYNC_ADMISSION_DISPATCHER_ENABLED + value: "true" - name: AI_GATEWAY_DATABASE_MAX_CONNS value: "32" - name: AI_GATEWAY_DATABASE_CRITICAL_MAX_CONNS @@ -501,6 +505,8 @@ spec: value: worker - name: AI_GATEWAY_ASYNC_QUEUE_WORKER_ENABLED value: "true" + - name: AI_GATEWAY_ASYNC_ADMISSION_DISPATCHER_ENABLED + value: "false" - name: AI_GATEWAY_DATABASE_MAX_CONNS value: "32" - name: AI_GATEWAY_DATABASE_CRITICAL_MAX_CONNS @@ -676,6 +682,8 @@ spec: value: worker - name: AI_GATEWAY_ASYNC_QUEUE_WORKER_ENABLED value: "true" + - name: AI_GATEWAY_ASYNC_ADMISSION_DISPATCHER_ENABLED + value: "false" - name: AI_GATEWAY_DATABASE_MAX_CONNS value: "32" - name: AI_GATEWAY_DATABASE_CRITICAL_MAX_CONNS diff --git a/docs/operations/production-acceptance.md b/docs/operations/production-acceptance.md index 3c33ccb..2b3d05c 100644 --- a/docs/operations/production-acceptance.md +++ b/docs/operations/production-acceptance.md @@ -76,6 +76,7 @@ WebP、合法 4K 图片和越过 Seedance 2.0 官方输入边界的 `6144x2160` ```text AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE +AI_GATEWAY_ASYNC_ADMISSION_DISPATCHER_ENABLED AI_GATEWAY_DATABASE_MAX_CONNS AI_GATEWAY_DATABASE_MIN_IDLE_CONNS AI_GATEWAY_DATABASE_MAX_CONN_IDLE_SECONDS