perf(worker): 以有界微批次提升准入吞吐

原因:线上 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 渲染及人工发布脚本均通过。
This commit is contained in:
2026-08-01 20:26:47 +08:00
parent 92e328a575
commit c89c56ca65
13 changed files with 282 additions and 31 deletions
+5
View File
@@ -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 ||
+25
View File
@@ -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) {
+39 -25
View File
@@ -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
@@ -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 {
@@ -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(&microbatchAdmissions, &microbatchLeases, &microbatchRiverJobs); 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)