fix(worker): 容量饱和后停止重复准入扫描
双 Worker 每秒遍历全部 waiting 任务,在全局执行槽已满时仍重复探测,导致 River 执行协程长时间等待同一 worker_capacity 进程锁。 异步 FIFO 头部未准入时立即结束本轮扫描,等待租约释放通知后再继续;任务级错误仍跳过并处理后续任务。 验证:Go 全量测试、runner race、go vet、迁移安全检查通过。
This commit is contained in:
@@ -603,19 +603,25 @@ func (s *Service) SubmitAsyncTask(ctx context.Context, task store.GatewayTask) e
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Service) dispatchWaitingAsyncTask(ctx context.Context, task store.GatewayTask) error {
|
func (s *Service) dispatchWaitingAsyncTask(ctx context.Context, task store.GatewayTask) (bool, error) {
|
||||||
user := authUserFromTask(task)
|
user := authUserFromTask(task)
|
||||||
plan, err := s.buildTaskAdmissionPlan(ctx, task, user)
|
plan, err := s.buildTaskAdmissionPlan(ctx, task, user)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return false, err
|
||||||
}
|
}
|
||||||
if !plan.Eligible {
|
if !plan.Eligible {
|
||||||
return s.EnqueueAsyncTask(ctx, task)
|
if err := s.EnqueueAsyncTask(ctx, task); err != nil {
|
||||||
|
return false, err
|
||||||
|
}
|
||||||
|
return true, nil
|
||||||
}
|
}
|
||||||
_, err = s.tryTaskAdmissionWithAdmittedHook(ctx, task, plan, "", func(tx pgx.Tx) error {
|
result, err := s.tryTaskAdmissionWithAdmittedHook(ctx, task, plan, "", func(tx pgx.Tx) error {
|
||||||
return s.enqueueAsyncTaskTx(ctx, tx, task.ID, asyncTaskInsertOpts(task))
|
return s.enqueueAsyncTaskTx(ctx, tx, task.ID, asyncTaskInsertOpts(task))
|
||||||
})
|
})
|
||||||
return err
|
if err != nil {
|
||||||
|
return false, err
|
||||||
|
}
|
||||||
|
return result.Admitted, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Service) cancelAsyncSubmissionIfDisconnected(ctx context.Context, taskID string) bool {
|
func (s *Service) cancelAsyncSubmissionIfDisconnected(ctx context.Context, taskID string) bool {
|
||||||
@@ -660,8 +666,17 @@ func (s *Service) dispatchWaitingAsyncAdmissions(ctx context.Context) {
|
|||||||
if task.Status != "queued" {
|
if task.Status != "queued" {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
if err := s.dispatchWaitingAsyncTask(ctx, task); err != nil {
|
admitted, err := s.dispatchWaitingAsyncTask(ctx, task)
|
||||||
|
if err != nil {
|
||||||
s.logger.Warn("dispatch waiting async admission failed", "taskID", taskID, "error", err)
|
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
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user