perf(worker): 批量填充异步执行槽

跨地域同步复制下逐任务准入会为每个 River Job 支付一次事务提交 RTT,导致 P24 实际只能维持少量运行任务。调度器现在按全局容量窗口准备 FIFO 批次,在同一 PostgreSQL 事务内逐项校验队首、创建租约并插入唯一 River Job,一次提交即可填满可用槽;任一 Hook 失败时整批原子回滚。\n\n验收压力采样同时排除正在终止的 Ready Pod,避免滚动切换期间误连已移除容器。\n\n验证:真实 PostgreSQL 批量提交及整批回滚集成测试、Go 全量测试、go vet、runner/store race、gofmt、bash -n、ShellCheck、迁移安全检查通过。
This commit is contained in:
2026-07-31 07:59:57 +08:00
parent 14d13ad3a7
commit 5dd7765ae3
4 changed files with 364 additions and 82 deletions
+99 -38
View File
@@ -3,6 +3,7 @@ package runner
import (
"context"
"errors"
"fmt"
"hash/fnv"
"strings"
"time"
@@ -279,25 +280,42 @@ func (s *Service) tryTaskAdmissionWithAdmittedHook(
waiterID string,
onAdmitted func(pgx.Tx) error,
) (store.TaskAdmissionResult, error) {
input := taskAdmissionInput(task, plan, waiterID)
current, currentErr := s.store.GetTaskAdmission(ctx, task.ID)
if currentErr != nil && !errors.Is(currentErr, pgx.ErrNoRows) {
return store.TaskAdmissionResult{}, currentErr
}
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 {
return store.TaskAdmissionResult{}, rebindErr
}
s.observeTaskAdmission("candidate_migrated")
input, err := s.prepareTaskAdmissionInput(ctx, task, plan, waiterID)
if err != nil {
return store.TaskAdmissionResult{}, err
}
var result store.TaskAdmissionResult
var err error
if onAdmitted == nil {
result, err = s.store.TryTaskAdmission(ctx, input)
} else {
result, err = s.store.TryTaskAdmissionWithAdmittedHook(ctx, input, onAdmitted)
}
s.observeTaskAdmissionAttempt(result, err)
return result, err
}
func (s *Service) prepareTaskAdmissionInput(
ctx context.Context,
task store.GatewayTask,
plan taskAdmissionPlan,
waiterID string,
) (store.TaskAdmissionInput, error) {
input := taskAdmissionInput(task, plan, waiterID)
current, currentErr := s.store.GetTaskAdmission(ctx, task.ID)
if currentErr != nil && !errors.Is(currentErr, pgx.ErrNoRows) {
return store.TaskAdmissionInput{}, currentErr
}
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 {
return store.TaskAdmissionInput{}, rebindErr
}
s.observeTaskAdmission("candidate_migrated")
}
return input, nil
}
func (s *Service) observeTaskAdmissionAttempt(result store.TaskAdmissionResult, err error) {
if result.NewlyAdmitted {
s.observeTaskAdmission("admitted")
s.observeTaskAdmissionWait(time.Since(result.Admission.EnqueuedAt))
@@ -309,7 +327,6 @@ func (s *Service) tryTaskAdmissionWithAdmittedHook(
case errors.Is(err, store.ErrQueueTimeout):
s.observeTaskAdmission("timeout")
}
return result, err
}
func taskAdmissionInput(task store.GatewayTask, plan taskAdmissionPlan, waiterID string) store.TaskAdmissionInput {
@@ -603,25 +620,61 @@ func (s *Service) SubmitAsyncTask(ctx context.Context, task store.GatewayTask) e
return nil
}
func (s *Service) dispatchWaitingAsyncTask(ctx context.Context, task store.GatewayTask) (bool, error) {
user := authUserFromTask(task)
plan, err := s.buildTaskAdmissionPlan(ctx, task, user)
if err != nil {
return false, err
func (s *Service) dispatchWaitingAsyncTasks(ctx context.Context, tasks []store.GatewayTask) (bool, error) {
if len(tasks) == 0 {
return false, nil
}
if !plan.Eligible {
if err := s.EnqueueAsyncTask(ctx, task); err != nil {
inputs := make([]store.TaskAdmissionInput, 0, len(tasks))
tasksByID := make(map[string]store.GatewayTask, len(tasks))
for _, task := range tasks {
user := authUserFromTask(task)
plan, err := s.buildTaskAdmissionPlan(ctx, task, user)
if err != nil {
return false, err
}
return true, nil
if !plan.Eligible {
if err := s.EnqueueAsyncTask(ctx, task); err != nil {
return false, err
}
continue
}
input, err := s.prepareTaskAdmissionInput(ctx, task, plan, "")
if err != nil {
return false, err
}
inputs = append(inputs, input)
tasksByID[task.ID] = task
}
result, err := s.tryTaskAdmissionWithAdmittedHook(ctx, task, plan, "", func(tx pgx.Tx) error {
return s.enqueueAsyncTaskTx(ctx, tx, task.ID, asyncTaskInsertOpts(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
}
return result.Admitted, nil
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
}
if !outcome.Result.Admitted {
saturated = true
break
}
}
return saturated, nil
}
func (s *Service) cancelAsyncSubmissionIfDisconnected(ctx context.Context, taskID string) bool {
@@ -650,11 +703,19 @@ func (s *Service) dispatchWaitingAsyncAdmissions(ctx context.Context) {
case <-s.asyncAdmissionWake:
case <-ticker.C:
}
taskIDs, err := s.store.ListWaitingAsyncAdmissionTaskIDs(ctx, 1000)
batchLimit := s.cfg.AsyncWorkerHardLimit
if batchLimit <= 0 {
batchLimit = 64
}
if batchLimit > 128 {
batchLimit = 128
}
taskIDs, err := s.store.ListWaitingAsyncAdmissionTaskIDs(ctx, batchLimit)
if err != nil {
s.logger.Warn("list waiting async admissions failed", "error", err)
continue
}
tasks := make([]store.GatewayTask, 0, len(taskIDs))
for _, taskID := range taskIDs {
task, err := s.store.GetTask(ctx, taskID)
if err != nil {
@@ -666,18 +727,18 @@ func (s *Service) dispatchWaitingAsyncAdmissions(ctx context.Context) {
if task.Status != "queued" {
continue
}
admitted, err := s.dispatchWaitingAsyncTask(ctx, task)
if err != nil {
s.logger.Warn("dispatch waiting async admission failed", "taskID", taskID, "error", err)
continue
}
if !admitted {
// Every asynchronous task shares the global worker-capacity FIFO
// scope. Once its current head cannot be admitted, later tasks
// cannot make progress until a lease is released.
s.observeTaskAdmission("dispatch_saturated")
break
}
tasks = append(tasks, task)
}
saturated, err := s.dispatchWaitingAsyncTasks(ctx, tasks)
if err != nil {
s.logger.Warn("dispatch waiting async admission batch failed", "tasks", len(tasks), "error", err)
continue
}
if saturated {
// Every asynchronous task shares the global worker-capacity FIFO
// scope. Once its current head cannot be admitted, later tasks
// cannot make progress until a lease is released.
s.observeTaskAdmission("dispatch_saturated")
}
}
}