diff --git a/service/model_task_end_service.go b/service/model_task_end_service.go index 9e4e781..1c6d9b9 100644 --- a/service/model_task_end_service.go +++ b/service/model_task_end_service.go @@ -126,6 +126,13 @@ func (s *modelTaskEndService) handleSingleTask(ctx context.Context, item *entity UserName: item.Creator, TenantId: item.TenantId, }) + // 任务处理结束(无论成功失败)都释放 Redis 锁,避免失败路径残留锁、 + // 在锁 TTL(1200s)内阻塞任务被其他 worker 重新获取 + defer func() { + if _, delErr := g.Redis().Del(asyncCtx, "model_video_task:"+gconv.String(item.Id)); delErr != nil { + g.Log().Errorf(asyncCtx, "清理任务锁失败: %v", delErr) + } + }() // 按 modelId 现查模型配置(异步映射/token 映射/计费规则不随任务快照,任务完成时取当前配置) modelInfo, err := dao.ModelManage.GetNotTenantId(asyncCtx, &dto.GetModelManageReq{Id: item.ModelId}) if err != nil { @@ -241,11 +248,6 @@ func (s *modelTaskEndService) handleSingleTask(ctx context.Context, item *entity if err != nil { g.Log().Errorf(asyncCtx, "保存视频任务结果失败: %v", err) } - // 删除redis视频任务 - _, err = g.Redis().Del(asyncCtx, "model_video_task:"+gconv.String(item.Id)) - if err != nil { - return - } // 发布消息 if err = TaskMsgPublish(asyncCtx, item.MsgTopic, docMsg); err != nil { g.Log().Errorf(asyncCtx, "模型消息发布失败: %v", err)