From c89c56ca65b22ba46ef8698790426a7d5b9d772a Mon Sep 17 00:00:00 2001 From: wangbo Date: Sat, 1 Aug 2026 20:26:47 +0800 Subject: [PATCH] =?UTF-8?q?perf(worker):=20=E4=BB=A5=E6=9C=89=E7=95=8C?= =?UTF-8?q?=E5=BE=AE=E6=89=B9=E6=AC=A1=E6=8F=90=E5=8D=87=E5=87=86=E5=85=A5?= =?UTF-8?q?=E5=90=9E=E5=90=90?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 原因:线上 P24 同构验收中,单任务同步复制提交与全局容量锁串行化,使两个 Worker 的 48 个执行槽只能维持约 10–14 个运行任务,最老等待超过 15 分钟。 影响:新增可配置的 1–32 条准入微批次,默认 8;租约、任务准入与唯一 River job 在一个有界事务内原子提交,并按确定顺序预锁任务和容量范围,避免整窗 48 条大事务和多 Dispatcher 死锁。并容忍 rebind 已被其他 Dispatcher 完成的幂等竞态。 验证:Go 全量测试、go vet、真实 PostgreSQL 跨 Store 集成测试、ShellCheck、迁移安全检查、OpenAPI、前端 lint/test/build、Compose/Kubernetes 渲染及人工发布脚本均通过。 --- README.md | 2 +- apps/api/internal/config/config.go | 5 + apps/api/internal/config/config_test.go | 25 +++++ apps/api/internal/runner/admission.go | 64 ++++++++----- apps/api/internal/store/admission_queue.go | 93 +++++++++++++++++++ .../store/admission_queue_integration_test.go | 93 +++++++++++++++++++ .../easyai-ai-gateway-cluster-release | 9 +- ...ai-ai-gateway-cluster-release.conf.example | 1 + .../local-acceptance/local-config.yaml | 1 + deploy/kubernetes/production/application.yaml | 4 + docs/operations/production-acceptance.md | 1 + scripts/cluster/run-production-acceptance.sh | 14 ++- tests/release/cluster-release-helper-test.sh | 1 + 13 files changed, 282 insertions(+), 31 deletions(-) diff --git a/README.md b/README.md index 9ff9596..bbccf26 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_DATABASE_MIN_IDLE_CONNS` 控制启动预热连接数,默认 `0`,生产 K3s 配置为 `4`;`AI_GATEWAY_DATABASE_MAX_CONN_IDLE_SECONDS` 可缩短突发连接的空闲回收时间,生产配置为 `30` 秒,避免把连接池高水位长期常驻。平台模型和用户组的 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 合并为有界原子微批次,摊薄跨地域同步提交,同时禁止重新形成整窗大事务。`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 7b3ace6..7d89312 100644 --- a/apps/api/internal/config/config.go +++ b/apps/api/internal/config/config.go @@ -79,6 +79,7 @@ type Config struct { AsyncWorkerHardLimit int AsyncWorkerInstanceHardLimit int AsyncWorkerRefreshIntervalSeconds int + AsyncAdmissionMicrobatchSize int WorkerAutoscalingEnabled bool WorkerReplicasNingbo int WorkerReplicasHongkong int @@ -186,6 +187,7 @@ 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), 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), @@ -291,6 +293,9 @@ func (c Config) Validate() error { if c.AsyncWorkerRefreshIntervalSeconds < 1 { return errors.New("AI_GATEWAY_ASYNC_WORKER_REFRESH_INTERVAL_SECONDS must be positive") } + if c.AsyncAdmissionMicrobatchSize < 1 || c.AsyncAdmissionMicrobatchSize > 32 { + return errors.New("AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE must be between 1 and 32") + } if c.WorkerReplicasNingbo < 0 || c.WorkerReplicasHongkong < 0 || c.WorkerMinReplicasNingbo < 0 || c.WorkerMinReplicasHongkong < 0 || c.WorkerMaxReplicasNingbo < c.WorkerMinReplicasNingbo || diff --git a/apps/api/internal/config/config_test.go b/apps/api/internal/config/config_test.go index b10434a..074ef08 100644 --- a/apps/api/internal/config/config_test.go +++ b/apps/api/internal/config/config_test.go @@ -31,6 +31,7 @@ func TestValidateIdentityFileSecretStoreRequiresDirectory(t *testing.T) { AsyncWorkerHardLimit: 2048, AsyncWorkerInstanceHardLimit: 32, AsyncWorkerRefreshIntervalSeconds: 5, + AsyncAdmissionMicrobatchSize: 8, } if err := cfg.Validate(); err == nil || !strings.Contains(err.Error(), "IDENTITY_SECRET_DIR") { t.Fatalf("Validate() error = %v, want missing identity secret directory", err) @@ -48,6 +49,7 @@ func TestValidateIdentityKubernetesSecretStore(t *testing.T) { AsyncWorkerHardLimit: 2048, AsyncWorkerInstanceHardLimit: 32, AsyncWorkerRefreshIntervalSeconds: 5, + AsyncAdmissionMicrobatchSize: 8, } if err := cfg.Validate(); err == nil || !strings.Contains(err.Error(), "namespace") { t.Fatalf("Validate() error = %v, want missing namespace", err) @@ -66,6 +68,7 @@ func TestValidateIdentitySecurityEventTiming(t *testing.T) { AsyncWorkerHardLimit: 2048, AsyncWorkerInstanceHardLimit: 32, AsyncWorkerRefreshIntervalSeconds: 5, + AsyncAdmissionMicrobatchSize: 8, } if err := cfg.Validate(); err == nil || !strings.Contains(err.Error(), "heartbeat") { t.Fatalf("Validate() error = %v, want invalid stale threshold", err) @@ -77,6 +80,7 @@ func TestValidateAsyncWorkerSettings(t *testing.T) { AsyncWorkerHardLimit: 10001, AsyncWorkerInstanceHardLimit: 32, AsyncWorkerRefreshIntervalSeconds: 5, + AsyncAdmissionMicrobatchSize: 8, } if err := cfg.Validate(); err == nil || !strings.Contains(err.Error(), "HARD_LIMIT") { t.Fatalf("Validate() error = %v, want invalid hard limit", err) @@ -102,6 +106,27 @@ func TestValidateAsyncWorkerSettings(t *testing.T) { if err := loaded.Validate(); err == nil || !strings.Contains(err.Error(), "INSTANCE_HARD_LIMIT") { t.Fatalf("Validate() error = %v, want invalid non-integer instance hard limit", err) } + t.Setenv("AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT", "32") + t.Setenv("AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE", "33") + loaded = Load() + if err := loaded.Validate(); err == nil || !strings.Contains(err.Error(), "MICROBATCH_SIZE") { + t.Fatalf("Validate() error = %v, want invalid admission microbatch size", err) + } +} + +func TestLoadAsyncAdmissionMicrobatchSize(t *testing.T) { + cfg := Load() + if cfg.AsyncAdmissionMicrobatchSize != 8 { + t.Fatalf("admission microbatch size=%d, want 8", cfg.AsyncAdmissionMicrobatchSize) + } + t.Setenv("AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE", "4") + cfg = Load() + if err := cfg.Validate(); err != nil { + t.Fatalf("valid admission microbatch size was rejected: %v", err) + } + if cfg.AsyncAdmissionMicrobatchSize != 4 { + t.Fatalf("admission microbatch size=%d, want 4", cfg.AsyncAdmissionMicrobatchSize) + } } func TestLoadAsyncQueueWorkerEnabled(t *testing.T) { diff --git a/apps/api/internal/runner/admission.go b/apps/api/internal/runner/admission.go index 0c46349..51c6dfb 100644 --- a/apps/api/internal/runner/admission.go +++ b/apps/api/internal/runner/admission.go @@ -354,6 +354,13 @@ func (s *Service) prepareTaskAdmissionInput( if currentErr == nil && current.Status == "waiting" && (current.PlatformID != input.PlatformID || current.PlatformModelID != input.PlatformModelID || current.UserGroupID != input.UserGroupID) { if _, rebindErr := s.store.RebindWaitingTaskAdmission(ctx, input); rebindErr != nil { + // Another dispatcher may admit or remove the same FIFO row between + // the read above and the conditional rebind. Let the atomic admission + // transaction observe the committed state instead of aborting an + // entire dispatcher batch on that benign race. + if errors.Is(rebindErr, pgx.ErrNoRows) { + return input, nil + } return store.TaskAdmissionInput{}, rebindErr } s.observeTaskAdmission("candidate_migrated") @@ -715,33 +722,40 @@ func (s *Service) dispatchWaitingAsyncTasks(ctx context.Context, tasks []store.G inputs = append(inputs, input) tasksByID[task.ID] = task } - outcomes, err := s.store.TryTaskAdmissionBatchWithAdmittedHook( - ctx, - inputs, - func(tx pgx.Tx, input store.TaskAdmissionInput) error { - task, ok := tasksByID[input.TaskID] - if !ok { - return fmt.Errorf("async admission batch task %s was not prepared", input.TaskID) - } - return s.enqueueAsyncTaskTx(ctx, tx, task.ID, asyncTaskInsertOpts(task)) - }, - ) - if err != nil { - return false, err - } saturated := false - for _, outcome := range outcomes { - s.observeTaskAdmissionAttempt(outcome.Result, outcome.Err) - if outcome.Err != nil { - if errors.Is(outcome.Err, store.ErrTaskExecutionFinished) || - errors.Is(outcome.Err, store.ErrQueueTimeout) { - continue - } - return false, outcome.Err + microbatchSize := s.cfg.AsyncAdmissionMicrobatchSize + if microbatchSize <= 0 { + microbatchSize = 1 + } + for start := 0; start < len(inputs) && !saturated; start += microbatchSize { + end := min(start+microbatchSize, len(inputs)) + outcomes, err := s.store.TryTaskAdmissionAtomicBatchWithAdmittedHook( + ctx, + inputs[start:end], + func(tx pgx.Tx, input store.TaskAdmissionInput) error { + task, ok := tasksByID[input.TaskID] + if !ok { + return fmt.Errorf("async admission batch task %s was not prepared", input.TaskID) + } + return s.enqueueAsyncTaskTx(ctx, tx, task.ID, asyncTaskInsertOpts(task)) + }, + ) + if err != nil { + return false, err } - if !outcome.Result.Admitted { - saturated = true - break + for _, outcome := range outcomes { + s.observeTaskAdmissionAttempt(outcome.Result, outcome.Err) + if outcome.Err != nil { + if errors.Is(outcome.Err, store.ErrTaskExecutionFinished) || + errors.Is(outcome.Err, store.ErrQueueTimeout) { + continue + } + return false, outcome.Err + } + if !outcome.Result.Admitted { + saturated = true + break + } } } return saturated, nil diff --git a/apps/api/internal/store/admission_queue.go b/apps/api/internal/store/admission_queue.go index 6d3528d..d4b7bcf 100644 --- a/apps/api/internal/store/admission_queue.go +++ b/apps/api/internal/store/admission_queue.go @@ -502,6 +502,99 @@ func (s *Store) TryTaskAdmissionBatchWithAdmittedHook( return outcomes, nil } +// TryTaskAdmissionAtomicBatchWithAdmittedHook admits a deliberately small +// group of tasks in one transaction. The caller bounds the group size so the +// global capacity lock spans one synchronous commit without returning to the +// old failure mode where a full capacity window locked dozens of task rows. +// Admission leases and unique River jobs either commit together for the whole +// microbatch or roll back together. +func (s *Store) TryTaskAdmissionAtomicBatchWithAdmittedHook( + ctx context.Context, + inputs []TaskAdmissionInput, + onAdmitted func(pgx.Tx, TaskAdmissionInput) error, +) ([]TaskAdmissionBatchOutcome, error) { + if len(inputs) == 0 { + return nil, nil + } + for _, input := range inputs { + if err := validateTaskAdmissionInput(input); err != nil { + return nil, err + } + } + lockKeys := make([]string, 0, len(inputs)*3) + for _, input := range inputs { + lockKeys = append(lockKeys, admissionOperationLockKeys(input)...) + } + return retryAdmissionOperation(ctx, lockKeys, func() ([]TaskAdmissionBatchOutcome, error) { + return s.tryTaskAdmissionAtomicBatchOnce(ctx, inputs, onAdmitted, lockKeys) + }) +} + +func (s *Store) tryTaskAdmissionAtomicBatchOnce( + ctx context.Context, + inputs []TaskAdmissionInput, + onAdmitted func(pgx.Tx, TaskAdmissionInput) error, + lockKeys []string, +) ([]TaskAdmissionBatchOutcome, error) { + tx, err := s.pool.Begin(ctx) + if err != nil { + return nil, err + } + defer rollbackTransaction(tx) + // Lock every task and scope in deterministic order. Without this pre-lock, + // two dispatchers could each hold a different task lock while waiting for + // the shared worker-capacity lock. + for _, key := range normalizedAdmissionLockKeys(lockKeys) { + if err := tryAdmissionTransactionLock(ctx, tx, key); err != nil { + return nil, err + } + } + + outcomes := make([]TaskAdmissionBatchOutcome, 0, len(inputs)) + notify := false + for _, input := range inputs { + hook := func(tx pgx.Tx) error { + if onAdmitted == nil { + return nil + } + return onAdmitted(tx, input) + } + outcome, admissionErr := s.tryTaskAdmissionTx(ctx, tx, input, hook) + if errors.Is(admissionErr, ErrTaskExecutionFinished) { + outcomes = append(outcomes, TaskAdmissionBatchOutcome{ + TaskID: input.TaskID, + Err: admissionErr, + }) + continue + } + if admissionErr != nil { + outcomes = append(outcomes, TaskAdmissionBatchOutcome{ + TaskID: input.TaskID, + Err: admissionErr, + }) + return outcomes, admissionErr + } + outcomes = append(outcomes, TaskAdmissionBatchOutcome{ + TaskID: input.TaskID, + Result: outcome.Result, + Err: outcome.PostCommitErr, + }) + if outcome.NotifyTaskID != "" { + notify = true + } + if !outcome.Result.Admitted { + break + } + } + if err := tx.Commit(ctx); err != nil { + return outcomes, err + } + if notify { + s.notifyTaskAdmissionBestEffort(ctx, "*") + } + return outcomes, nil +} + func validateNewAdmissionCapacity(states []admissionScopeState) error { mustWait := false for _, state := range states { diff --git a/apps/api/internal/store/admission_queue_integration_test.go b/apps/api/internal/store/admission_queue_integration_test.go index c0b5b6b..b0cbca1 100644 --- a/apps/api/internal/store/admission_queue_integration_test.go +++ b/apps/api/internal/store/admission_queue_integration_test.go @@ -649,6 +649,99 @@ SELECT } } + atomicScope := batchScope + atomicScope.ScopeKey = "atomic-microbatch-" + suffix + atomicScope.ConcurrentLimit = 4 + atomicTasks := []GatewayTask{createTask(true), createTask(true), createTask(true), createTask(true)} + atomicInputs := make([]TaskAdmissionInput, 0, len(atomicTasks)) + for index, task := range atomicTasks { + input := inputFor(task, 300+index, "") + input.Scopes = []AdmissionScope{atomicScope} + if _, err := first.QueueTaskAdmissionWithHook(ctx, input, nil); err != nil { + t.Fatalf("queue atomic microbatch task %d: %v", index, err) + } + atomicInputs = append(atomicInputs, input) + } + atomicOutcomes, err := first.TryTaskAdmissionAtomicBatchWithAdmittedHook( + ctx, + atomicInputs, + func(tx pgx.Tx, input TaskAdmissionInput) error { + _, hookErr := tx.Exec(ctx, ` +UPDATE gateway_tasks +SET river_job_id = 987654600 +WHERE id = $1::uuid`, input.TaskID) + return hookErr + }, + ) + if err != nil || len(atomicOutcomes) != len(atomicTasks) { + t.Fatalf("atomic microbatch outcomes=%+v err=%v", atomicOutcomes, err) + } + var microbatchAdmissions, microbatchLeases, microbatchRiverJobs int + atomicTaskIDs := []string{atomicTasks[0].ID, atomicTasks[1].ID, atomicTasks[2].ID, atomicTasks[3].ID} + if err := first.pool.QueryRow(ctx, ` +SELECT + (SELECT count(*) FROM gateway_task_admissions WHERE task_id = ANY($1::uuid[]) AND status = 'admitted'), + (SELECT count(*) FROM gateway_concurrency_leases WHERE task_id = ANY($1::uuid[]) AND released_at IS NULL), + (SELECT count(*) FROM gateway_tasks WHERE id = ANY($1::uuid[]) AND river_job_id IS NOT NULL)`, + atomicTaskIDs, + ).Scan(µbatchAdmissions, µbatchLeases, µbatchRiverJobs); err != nil { + t.Fatalf("read atomic microbatch: %v", err) + } + if microbatchAdmissions != 4 || microbatchLeases != 4 || microbatchRiverJobs != 4 { + t.Fatalf("atomic microbatch admissions=%d leases=%d jobs=%d, want 4/4/4", + microbatchAdmissions, microbatchLeases, microbatchRiverJobs) + } + for _, task := range atomicTasks { + if err := first.DeleteTaskAdmission(ctx, task.ID); err != nil { + t.Fatalf("release atomic microbatch task %s: %v", task.ID, err) + } + } + + atomicRollbackScope := batchScope + atomicRollbackScope.ScopeKey = "atomic-rollback-" + suffix + atomicRollbackScope.ConcurrentLimit = 2 + atomicRollbackTasks := []GatewayTask{createTask(true), createTask(true)} + atomicRollbackInputs := make([]TaskAdmissionInput, 0, len(atomicRollbackTasks)) + for index, task := range atomicRollbackTasks { + input := inputFor(task, 400+index, "") + input.Scopes = []AdmissionScope{atomicRollbackScope} + if _, err := first.QueueTaskAdmissionWithHook(ctx, input, nil); err != nil { + t.Fatalf("queue atomic rollback task %d: %v", index, err) + } + atomicRollbackInputs = append(atomicRollbackInputs, input) + } + _, err = first.TryTaskAdmissionAtomicBatchWithAdmittedHook( + ctx, + atomicRollbackInputs, + func(_ pgx.Tx, input TaskAdmissionInput) error { + if input.TaskID == atomicRollbackTasks[1].ID { + return batchHookFailure + } + return nil + }, + ) + if !errors.Is(err, batchHookFailure) { + t.Fatalf("atomic rollback error=%v, want synthetic failure", err) + } + var atomicRollbackWaiting, atomicRollbackLeases int + if err := first.pool.QueryRow(ctx, ` +SELECT + (SELECT count(*) FROM gateway_task_admissions WHERE task_id = ANY($1::uuid[]) AND status = 'waiting'), + (SELECT count(*) FROM gateway_concurrency_leases WHERE task_id = ANY($1::uuid[]) AND released_at IS NULL)`, + []string{atomicRollbackTasks[0].ID, atomicRollbackTasks[1].ID}, + ).Scan(&atomicRollbackWaiting, &atomicRollbackLeases); err != nil { + t.Fatalf("read atomic rollback: %v", err) + } + if atomicRollbackWaiting != 2 || atomicRollbackLeases != 0 { + t.Fatalf("atomic rollback waiting=%d leases=%d, want 2/0", + atomicRollbackWaiting, atomicRollbackLeases) + } + for _, task := range atomicRollbackTasks { + if err := first.DeleteTaskAdmission(ctx, task.ID); err != nil { + t.Fatalf("release atomic rollback task %s: %v", task.ID, err) + } + } + result, err = second.TryTaskAdmission(ctx, inputFor(queuedAtomicTask, 100, "")) if err != nil || !result.Admitted || len(result.Leases) != 1 { t.Fatalf("worker-time queued admission result=%+v err=%v", result, err) diff --git a/deploy/kubernetes/easyai-ai-gateway-cluster-release b/deploy/kubernetes/easyai-ai-gateway-cluster-release index 6607002..6ecfc99 100755 --- a/deploy/kubernetes/easyai-ai-gateway-cluster-release +++ b/deploy/kubernetes/easyai-ai-gateway-cluster-release @@ -20,6 +20,7 @@ source "$config_file" : "${AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT:=24}" : "${AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT:=48}" : "${AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT:=$AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT}" +: "${AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE:=8}" : "${AI_GATEWAY_WORKER_MEMORY_REQUEST_MIB:=1536}" : "${AI_GATEWAY_WORKER_CPU_REQUEST_MILLICORES:=500}" : "${AI_GATEWAY_DATABASE_MAX_CONNS:=32}" @@ -71,6 +72,7 @@ for capacity_value in \ "$AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT" \ "$AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT" \ "$AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT" \ + "$AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE" \ "$AI_GATEWAY_WORKER_MEMORY_REQUEST_MIB" \ "$AI_GATEWAY_WORKER_CPU_REQUEST_MILLICORES" \ "$AI_GATEWAY_DATABASE_MAX_CONNS" \ @@ -113,6 +115,7 @@ done AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT <= 256 && AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT <= 512 && AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT <= 512 && + AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE <= 32 && AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT == AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT && AI_GATEWAY_WORKER_MEMORY_REQUEST_MIB >= 256 && AI_GATEWAY_WORKER_MEMORY_REQUEST_MIB <= 1900 && @@ -170,6 +173,7 @@ capacity_config_hash() { "$AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT" \ "$AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT" \ "$AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT" \ + "$AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE" \ "$AI_GATEWAY_WORKER_MEMORY_REQUEST_MIB" \ "$AI_GATEWAY_WORKER_CPU_REQUEST_MILLICORES" \ "$AI_GATEWAY_DATABASE_MAX_CONNS" \ @@ -211,6 +215,7 @@ record_capacity_config() { --arg autoscaling "$AI_GATEWAY_WORKER_AUTOSCALING_ENABLED" \ --arg instanceLimit "$AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT" \ --arg globalLimit "$AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT" \ + --arg admissionMicrobatch "$AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE" \ --arg workerPool "$AI_GATEWAY_DATABASE_MAX_CONNS" \ --arg materialization "$AI_GATEWAY_MEDIA_MATERIALIZATION_CONCURRENCY" \ --arg mediaRequest "$AI_GATEWAY_MEDIA_REQUEST_CONCURRENCY" \ @@ -233,6 +238,7 @@ record_capacity_config() { AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT:$instanceLimit, AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT:$globalLimit, AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT:$globalLimit, + AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE:$admissionMicrobatch, AI_GATEWAY_WORKER_DATABASE_MAX_CONNS:$workerPool, AI_GATEWAY_MEDIA_MATERIALIZATION_CONCURRENCY:$materialization, AI_GATEWAY_MEDIA_REQUEST_CONCURRENCY:$mediaRequest, @@ -396,6 +402,7 @@ rollout_worker_site() { "AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT=$AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT" \ "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_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" \ @@ -617,7 +624,7 @@ apply_capacity_config() { verify_site ningbo wait_for_url 'public readiness' "$PUBLIC_BASE_URL/api/v1/readyz" '"ok":true' record_capacity_config - echo "production_capacity=PASS config_hash=$(capacity_config_hash) ningbo_replicas=$AI_GATEWAY_WORKER_REPLICAS_NINGBO hongkong_replicas=$AI_GATEWAY_WORKER_REPLICAS_HONGKONG instance_limit=$AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT global_limit=$AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT worker_database_pool=$AI_GATEWAY_DATABASE_MAX_CONNS api_database_pool=$AI_GATEWAY_API_DATABASE_MAX_CONNS database_min_idle=$AI_GATEWAY_DATABASE_MIN_IDLE_CONNS database_max_idle_seconds=$AI_GATEWAY_DATABASE_MAX_CONN_IDLE_SECONDS media_concurrency=$AI_GATEWAY_MEDIA_MATERIALIZATION_CONCURRENCY worker_memory_request_mib=$AI_GATEWAY_WORKER_MEMORY_REQUEST_MIB worker_cpu_request_millicores=$AI_GATEWAY_WORKER_CPU_REQUEST_MILLICORES autoscaling=$AI_GATEWAY_WORKER_AUTOSCALING_ENABLED" + echo "production_capacity=PASS config_hash=$(capacity_config_hash) ningbo_replicas=$AI_GATEWAY_WORKER_REPLICAS_NINGBO hongkong_replicas=$AI_GATEWAY_WORKER_REPLICAS_HONGKONG instance_limit=$AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT global_limit=$AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT admission_microbatch=$AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE worker_database_pool=$AI_GATEWAY_DATABASE_MAX_CONNS api_database_pool=$AI_GATEWAY_API_DATABASE_MAX_CONNS database_min_idle=$AI_GATEWAY_DATABASE_MIN_IDLE_CONNS database_max_idle_seconds=$AI_GATEWAY_DATABASE_MAX_CONN_IDLE_SECONDS media_concurrency=$AI_GATEWAY_MEDIA_MATERIALIZATION_CONCURRENCY worker_memory_request_mib=$AI_GATEWAY_WORKER_MEMORY_REQUEST_MIB worker_cpu_request_millicores=$AI_GATEWAY_WORKER_CPU_REQUEST_MILLICORES autoscaling=$AI_GATEWAY_WORKER_AUTOSCALING_ENABLED" } rollout_site() { diff --git a/deploy/kubernetes/easyai-ai-gateway-cluster-release.conf.example b/deploy/kubernetes/easyai-ai-gateway-cluster-release.conf.example index ac40d5f..c44f663 100644 --- a/deploy/kubernetes/easyai-ai-gateway-cluster-release.conf.example +++ b/deploy/kubernetes/easyai-ai-gateway-cluster-release.conf.example @@ -9,6 +9,7 @@ DESIRED_STATE_DIR=/usr/local/share/easyai-ai-gateway-release/production AI_GATEWAY_WORKER_REPLICAS_NINGBO=0 AI_GATEWAY_WORKER_REPLICAS_HONGKONG=2 AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT=24 +AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE=8 AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT=48 AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT=48 AI_GATEWAY_WORKER_MEMORY_REQUEST_MIB=1536 diff --git a/deploy/kubernetes/local-acceptance/local-config.yaml b/deploy/kubernetes/local-acceptance/local-config.yaml index 5fe165d..567bd46 100644 --- a/deploy/kubernetes/local-acceptance/local-config.yaml +++ b/deploy/kubernetes/local-acceptance/local-config.yaml @@ -31,6 +31,7 @@ data: AI_GATEWAY_MEDIA_OSS_DIRECT_ENABLED: "false" AI_GATEWAY_ASYNC_QUEUE_WORKER_ENABLED: "true" AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT: "24" + AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE: "8" AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT: "24" AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT: "48" AI_GATEWAY_ASYNC_WORKER_REFRESH_INTERVAL_SECONDS: "5" diff --git a/deploy/kubernetes/production/application.yaml b/deploy/kubernetes/production/application.yaml index ed97c65..1561035 100644 --- a/deploy/kubernetes/production/application.yaml +++ b/deploy/kubernetes/production/application.yaml @@ -513,6 +513,8 @@ spec: value: "30" - name: AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT value: "24" + - name: AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE + value: "8" - name: AI_GATEWAY_MEDIA_MATERIALIZATION_CONCURRENCY value: "24" - name: AI_GATEWAY_MEDIA_REQUEST_CONCURRENCY @@ -686,6 +688,8 @@ spec: value: "30" - name: AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT value: "24" + - name: AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE + value: "8" - name: AI_GATEWAY_MEDIA_MATERIALIZATION_CONCURRENCY value: "24" - name: AI_GATEWAY_MEDIA_REQUEST_CONCURRENCY diff --git a/docs/operations/production-acceptance.md b/docs/operations/production-acceptance.md index bd79eba..0202169 100644 --- a/docs/operations/production-acceptance.md +++ b/docs/operations/production-acceptance.md @@ -75,6 +75,7 @@ WebP、合法 4K 图片和越过 Seedance 2.0 官方输入边界的 `6144x2160` ```text AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT +AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE AI_GATEWAY_DATABASE_MAX_CONNS AI_GATEWAY_DATABASE_MIN_IDLE_CONNS AI_GATEWAY_DATABASE_MAX_CONN_IDLE_SECONDS diff --git a/scripts/cluster/run-production-acceptance.sh b/scripts/cluster/run-production-acceptance.sh index a49a39f..4dcde52 100755 --- a/scripts/cluster/run-production-acceptance.sh +++ b/scripts/cluster/run-production-acceptance.sh @@ -89,6 +89,7 @@ require_commands curl git go jq node openssl sed shasum : "${AI_GATEWAY_ACCEPTANCE_WORKER_RIVER_MAX_CONNS:=8}" : "${AI_GATEWAY_ACCEPTANCE_API_RIVER_MAX_CONNS:=4}" : "${AI_GATEWAY_ACCEPTANCE_DATABASE_MAX_CONN_IDLE_SECONDS:=300}" +: "${AI_GATEWAY_ACCEPTANCE_ASYNC_ADMISSION_MICROBATCH_SIZE:=8}" : "${AI_GATEWAY_ACCEPTANCE_API_MEDIA_REQUEST_CONCURRENCY:=128}" : "${AI_GATEWAY_ACCEPTANCE_IDENTITY_SHARDS:=32}" : "${AI_GATEWAY_ACCEPTANCE_WORKER_MEMORY_REQUEST_MIB:=1536}" @@ -128,6 +129,11 @@ fi echo 'AI_GATEWAY_ACCEPTANCE_DATABASE_MAX_CONN_IDLE_SECONDS must be between 1 and 3600' >&2 exit 1 } +[[ $AI_GATEWAY_ACCEPTANCE_ASYNC_ADMISSION_MICROBATCH_SIZE =~ ^[1-9][0-9]*$ && + $AI_GATEWAY_ACCEPTANCE_ASYNC_ADMISSION_MICROBATCH_SIZE -le 32 ]] || { + echo 'AI_GATEWAY_ACCEPTANCE_ASYNC_ADMISSION_MICROBATCH_SIZE must be between 1 and 32' >&2 + exit 1 +} [[ $AI_GATEWAY_ACCEPTANCE_WORKER_RIVER_MAX_CONNS =~ ^[1-9][0-9]*$ && $AI_GATEWAY_ACCEPTANCE_API_RIVER_MAX_CONNS =~ ^[1-9][0-9]*$ && $AI_GATEWAY_ACCEPTANCE_WORKER_RIVER_MAX_CONNS -lt 28 && @@ -1547,14 +1553,14 @@ apply_capacity_profile() { global=$((slots * baseline_replicas)) local command_text capacity_output printf -v command_text \ - 'AI_GATEWAY_WORKER_REPLICAS_NINGBO=%q AI_GATEWAY_WORKER_REPLICAS_HONGKONG=%q AI_GATEWAY_WORKER_MIN_REPLICAS_NINGBO=%q AI_GATEWAY_WORKER_MIN_REPLICAS_HONGKONG=%q AI_GATEWAY_WORKER_MAX_REPLICAS_NINGBO=%q AI_GATEWAY_WORKER_MAX_REPLICAS_HONGKONG=%q AI_GATEWAY_WORKER_AUTOSCALING_ENABLED=false AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT=%q AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT=%q AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT=%q AI_GATEWAY_DATABASE_MAX_CONNS=%q AI_GATEWAY_API_DATABASE_MAX_CONNS=%q AI_GATEWAY_DATABASE_CRITICAL_MAX_CONNS=4 AI_GATEWAY_DATABASE_RIVER_MAX_CONNS=%q AI_GATEWAY_API_DATABASE_RIVER_MAX_CONNS=%q AI_GATEWAY_DATABASE_MIN_IDLE_CONNS=%q AI_GATEWAY_DATABASE_MAX_CONN_IDLE_SECONDS=%q AI_GATEWAY_MEDIA_MATERIALIZATION_CONCURRENCY=%q AI_GATEWAY_MEDIA_REQUEST_CONCURRENCY=%q AI_GATEWAY_WORKER_TARGET_OUTSTANDING_PER_REPLICA=%q AI_GATEWAY_WORKER_MEMORY_REQUEST_MIB=%q AI_GATEWAY_WORKER_CPU_REQUEST_MILLICORES=%q %q capacity' \ + 'AI_GATEWAY_WORKER_REPLICAS_NINGBO=%q AI_GATEWAY_WORKER_REPLICAS_HONGKONG=%q AI_GATEWAY_WORKER_MIN_REPLICAS_NINGBO=%q AI_GATEWAY_WORKER_MIN_REPLICAS_HONGKONG=%q AI_GATEWAY_WORKER_MAX_REPLICAS_NINGBO=%q AI_GATEWAY_WORKER_MAX_REPLICAS_HONGKONG=%q AI_GATEWAY_WORKER_AUTOSCALING_ENABLED=false AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT=%q AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT=%q AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT=%q AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE=%q AI_GATEWAY_DATABASE_MAX_CONNS=%q AI_GATEWAY_API_DATABASE_MAX_CONNS=%q AI_GATEWAY_DATABASE_CRITICAL_MAX_CONNS=4 AI_GATEWAY_DATABASE_RIVER_MAX_CONNS=%q AI_GATEWAY_API_DATABASE_RIVER_MAX_CONNS=%q AI_GATEWAY_DATABASE_MIN_IDLE_CONNS=%q AI_GATEWAY_DATABASE_MAX_CONN_IDLE_SECONDS=%q AI_GATEWAY_MEDIA_MATERIALIZATION_CONCURRENCY=%q AI_GATEWAY_MEDIA_REQUEST_CONCURRENCY=%q AI_GATEWAY_WORKER_TARGET_OUTSTANDING_PER_REPLICA=%q AI_GATEWAY_WORKER_MEMORY_REQUEST_MIB=%q AI_GATEWAY_WORKER_CPU_REQUEST_MILLICORES=%q %q capacity' \ "$AI_GATEWAY_ACCEPTANCE_BASE_REPLICAS_NINGBO" \ "$AI_GATEWAY_ACCEPTANCE_BASE_REPLICAS_HONGKONG" \ "$AI_GATEWAY_ACCEPTANCE_BASE_REPLICAS_NINGBO" \ "$AI_GATEWAY_ACCEPTANCE_BASE_REPLICAS_HONGKONG" \ "$AI_GATEWAY_ACCEPTANCE_BASE_REPLICAS_NINGBO" \ "$AI_GATEWAY_ACCEPTANCE_BASE_REPLICAS_HONGKONG" \ - "$slots" "$global" "$global" "$pool" "$api_pool" \ + "$slots" "$global" "$global" "$AI_GATEWAY_ACCEPTANCE_ASYNC_ADMISSION_MICROBATCH_SIZE" "$pool" "$api_pool" \ "$AI_GATEWAY_ACCEPTANCE_WORKER_RIVER_MAX_CONNS" "$AI_GATEWAY_ACCEPTANCE_API_RIVER_MAX_CONNS" \ "$min_idle" "$max_idle_seconds" \ "$media" "$AI_GATEWAY_ACCEPTANCE_API_MEDIA_REQUEST_CONCURRENCY" "$((slots * 2))" \ @@ -1721,14 +1727,14 @@ apply_autoscaling_profile() { global=$((slots * (max_replicas_ningbo + max_replicas_hongkong))) local command_text capacity_output printf -v command_text \ - 'AI_GATEWAY_WORKER_REPLICAS_NINGBO=%q AI_GATEWAY_WORKER_REPLICAS_HONGKONG=%q AI_GATEWAY_WORKER_MIN_REPLICAS_NINGBO=%q AI_GATEWAY_WORKER_MIN_REPLICAS_HONGKONG=%q AI_GATEWAY_WORKER_MAX_REPLICAS_NINGBO=%q AI_GATEWAY_WORKER_MAX_REPLICAS_HONGKONG=%q AI_GATEWAY_WORKER_AUTOSCALING_ENABLED=true AI_GATEWAY_WORKER_SCALE_UP_WINDOW_SECONDS=20 AI_GATEWAY_WORKER_SCALE_DOWN_STABILIZATION_SECONDS=600 AI_GATEWAY_WORKER_DRAIN_TIMEOUT_SECONDS=600 AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT=%q AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT=%q AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT=%q AI_GATEWAY_DATABASE_MAX_CONNS=%q AI_GATEWAY_API_DATABASE_MAX_CONNS=%q AI_GATEWAY_DATABASE_CRITICAL_MAX_CONNS=4 AI_GATEWAY_DATABASE_RIVER_MAX_CONNS=%q AI_GATEWAY_API_DATABASE_RIVER_MAX_CONNS=%q AI_GATEWAY_DATABASE_MIN_IDLE_CONNS=4 AI_GATEWAY_DATABASE_MAX_CONN_IDLE_SECONDS=%q AI_GATEWAY_MEDIA_MATERIALIZATION_CONCURRENCY=%q AI_GATEWAY_MEDIA_REQUEST_CONCURRENCY=%q AI_GATEWAY_WORKER_TARGET_OUTSTANDING_PER_REPLICA=%q AI_GATEWAY_WORKER_MEMORY_REQUEST_MIB=%q AI_GATEWAY_WORKER_CPU_REQUEST_MILLICORES=%q %q capacity' \ + 'AI_GATEWAY_WORKER_REPLICAS_NINGBO=%q AI_GATEWAY_WORKER_REPLICAS_HONGKONG=%q AI_GATEWAY_WORKER_MIN_REPLICAS_NINGBO=%q AI_GATEWAY_WORKER_MIN_REPLICAS_HONGKONG=%q AI_GATEWAY_WORKER_MAX_REPLICAS_NINGBO=%q AI_GATEWAY_WORKER_MAX_REPLICAS_HONGKONG=%q AI_GATEWAY_WORKER_AUTOSCALING_ENABLED=true AI_GATEWAY_WORKER_SCALE_UP_WINDOW_SECONDS=20 AI_GATEWAY_WORKER_SCALE_DOWN_STABILIZATION_SECONDS=600 AI_GATEWAY_WORKER_DRAIN_TIMEOUT_SECONDS=600 AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT=%q AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT=%q AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT=%q AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE=%q AI_GATEWAY_DATABASE_MAX_CONNS=%q AI_GATEWAY_API_DATABASE_MAX_CONNS=%q AI_GATEWAY_DATABASE_CRITICAL_MAX_CONNS=4 AI_GATEWAY_DATABASE_RIVER_MAX_CONNS=%q AI_GATEWAY_API_DATABASE_RIVER_MAX_CONNS=%q AI_GATEWAY_DATABASE_MIN_IDLE_CONNS=4 AI_GATEWAY_DATABASE_MAX_CONN_IDLE_SECONDS=%q AI_GATEWAY_MEDIA_MATERIALIZATION_CONCURRENCY=%q AI_GATEWAY_MEDIA_REQUEST_CONCURRENCY=%q AI_GATEWAY_WORKER_TARGET_OUTSTANDING_PER_REPLICA=%q AI_GATEWAY_WORKER_MEMORY_REQUEST_MIB=%q AI_GATEWAY_WORKER_CPU_REQUEST_MILLICORES=%q %q capacity' \ "$AI_GATEWAY_ACCEPTANCE_AUTOSCALING_MIN_REPLICAS_NINGBO" \ "$AI_GATEWAY_ACCEPTANCE_AUTOSCALING_MIN_REPLICAS_HONGKONG" \ "$AI_GATEWAY_ACCEPTANCE_AUTOSCALING_MIN_REPLICAS_NINGBO" \ "$AI_GATEWAY_ACCEPTANCE_AUTOSCALING_MIN_REPLICAS_HONGKONG" \ "$max_replicas_ningbo" \ "$max_replicas_hongkong" \ - "$slots" "$global" "$global" "$pool" \ + "$slots" "$global" "$global" "$AI_GATEWAY_ACCEPTANCE_ASYNC_ADMISSION_MICROBATCH_SIZE" "$pool" \ "$AI_GATEWAY_ACCEPTANCE_API_DATABASE_MAX_CONNS" \ "$AI_GATEWAY_ACCEPTANCE_WORKER_RIVER_MAX_CONNS" \ "$AI_GATEWAY_ACCEPTANCE_API_RIVER_MAX_CONNS" \ diff --git a/tests/release/cluster-release-helper-test.sh b/tests/release/cluster-release-helper-test.sh index 2aa216b..2a06b49 100755 --- a/tests/release/cluster-release-helper-test.sh +++ b/tests/release/cluster-release-helper-test.sh @@ -45,6 +45,7 @@ AI_GATEWAY_WORKER_REPLICAS_HONGKONG=1 AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT=24 AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT=48 AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT=48 +AI_GATEWAY_ASYNC_ADMISSION_MICROBATCH_SIZE=8 AI_GATEWAY_DATABASE_MAX_CONNS=32 AI_GATEWAY_DATABASE_CRITICAL_MAX_CONNS=4 AI_GATEWAY_DATABASE_RIVER_MAX_CONNS=4