mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-29 05:56:38 +08:00
优化任务管理
This commit is contained in:
@@ -56,12 +56,9 @@ Admin dispatch -> ValidateAndNormalizePayload -> DispatchTask
|
||||
- 在 `internal/task/worker/worker.go` 的 Asynq mux 中添加 `mux.HandleFunc(task.YourAsynqTask, task.ProcessTask)`。
|
||||
- 所有业务任务都应交给 `task.ProcessTask`,由 executor 根据 task type 分发到 handler。
|
||||
|
||||
5. 如需 Cron 调度,补齐配置链路。
|
||||
- 在 `internal/task/scheduler/scheduler.go` 注册 cron。
|
||||
- 在 `internal/config/model.go` 添加 scheduler config 字段。
|
||||
- 在 `config.example.yaml` 添加对应配置项。
|
||||
- runtime 代码从 `config.Config` 读取配置,不直接读环境变量。
|
||||
- Scheduler 直接入队的任务可能没有 Admin 创建的 `TaskExecution` 记录;需要可见执行记录时,优先通过 Admin dispatch 触发。
|
||||
5. 如需 Cron 调度,系统默认定时任务必须通过 SQL 迁移(goose)初始化。
|
||||
- 确保任务在 `DispatchableTasks` 中已正确配置 `TaskMeta`。
|
||||
- 在 `internal/db/migrator/goose/postgres` 和 `sqlite` 下编写 migration 脚本,使用 `INSERT INTO schedules` 语句初始化任务,指定 `task_type` 和 `cron` 等字段。必须妥善处理冲突(如 `ON CONFLICT DO NOTHING`)以支持幂等。
|
||||
|
||||
6. 如改动 Admin API。
|
||||
- handler 放在 `internal/apps/admin/<module>/` 或现有 Admin task 模块内。
|
||||
|
||||
@@ -199,34 +199,21 @@ func StartWorker() error {
|
||||
|
||||
## Cron 调度和配置
|
||||
|
||||
如果任务需要定时运行,补齐 scheduler、config struct 和 `config.example.yaml`。
|
||||
系统默认的定时任务必须通过 Goose SQL 迁移初始化插入到 `schedules` 表。
|
||||
|
||||
```go
|
||||
const (
|
||||
cleanupDedupWindow = 23 * time.Hour
|
||||
cleanupMaxRetry = 3
|
||||
)
|
||||
在 `internal/db/migrator/goose/postgres` 下的示例:
|
||||
|
||||
if _, err = scheduler.Register(
|
||||
config.Config.Scheduler.CleanupUnusedUploadsTaskCron,
|
||||
asynq.NewTask(task.CleanupUnusedUploadsTask, nil),
|
||||
asynq.Unique(cleanupDedupWindow),
|
||||
asynq.MaxRetry(cleanupMaxRetry),
|
||||
); err != nil {
|
||||
return
|
||||
}
|
||||
```sql
|
||||
-- +goose Up
|
||||
INSERT INTO schedules (id, name, task_type, cron, payload, is_active, created_at, updated_at)
|
||||
VALUES (1, '清理未使用上传', 'cleanup_unused_uploads', '0 */2 * * *', '{}', TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||
ON CONFLICT (id) DO NOTHING;
|
||||
|
||||
-- +goose Down
|
||||
-- 根据业务需求决定是否需要在此删除
|
||||
```
|
||||
|
||||
```go
|
||||
type schedulerConfig struct {
|
||||
CleanupUnusedUploadsTaskCron string `mapstructure:"cleanup_unused_uploads_task_cron"`
|
||||
}
|
||||
```
|
||||
|
||||
```yaml
|
||||
scheduler:
|
||||
cleanup_unused_uploads_task_cron: "0 */2 * * *"
|
||||
```
|
||||
对于 `sqlite` 也可以使用类似的 `INSERT INTO ... ON CONFLICT(id) DO NOTHING` 语法。数据库更新后,后端会自动热重载调度器。
|
||||
|
||||
## Handler 测试
|
||||
|
||||
|
||||
@@ -379,7 +379,7 @@ export function TaskExecutionsManager() {
|
||||
{detailLoading ? <Spinner className="size-4" /> : <RefreshCw className="size-4" />}
|
||||
刷新详情
|
||||
</Button>
|
||||
{selectedExecution?.status === "failed" && selectedExecution.retryable && selectedExecution.retry_count < selectedExecution.max_retry && (
|
||||
{selectedExecution && selectedExecution.status !== "pending" && selectedExecution.status !== "running" && selectedExecution.retryable && selectedExecution.retry_count < selectedExecution.max_retry && (
|
||||
<Button onClick={handleRetryExecution} disabled={retrying}>
|
||||
{retrying ? <Spinner className="size-4" /> : <RotateCcw className="size-4" />}
|
||||
重试任务
|
||||
|
||||
@@ -187,7 +187,7 @@ func RetryTask(c *gin.Context) {
|
||||
switch {
|
||||
case strings.Contains(errMsg, "不存在"):
|
||||
c.JSON(http.StatusNotFound, util.Err(errMsg))
|
||||
case strings.Contains(errMsg, "只有失败") || strings.Contains(errMsg, "不支持重试") || strings.Contains(errMsg, "已达到最大重试"):
|
||||
case strings.Contains(errMsg, "只有非进行中") || strings.Contains(errMsg, "不支持重试") || strings.Contains(errMsg, "已达到最大重试"):
|
||||
c.JSON(http.StatusBadRequest, util.Err(errMsg))
|
||||
default:
|
||||
c.JSON(http.StatusInternalServerError, util.Err(fmt.Sprintf("%s: %v", TaskRetryFailed, err)))
|
||||
|
||||
@@ -8,7 +8,7 @@ const (
|
||||
errCreateTaskExecutionFailed = "创建任务执行记录失败: %w"
|
||||
errTaskEnqueueFailed = "任务入队失败: %w"
|
||||
errTaskExecutionNotFound = "任务执行记录不存在: %w"
|
||||
errRetryOnlyFailedTask = "只有失败的任务才能重试,当前状态: %s"
|
||||
errRetryOnlyFinishedTask = "只有非进行中的任务才能重试,当前状态: %s"
|
||||
errTaskNotRetryable = "该任务不支持重试"
|
||||
errTaskMaxRetryExceeded = "已达到最大重试次数 %d"
|
||||
errCreateRetryExecutionFailed = "创建重试任务执行记录失败: %w"
|
||||
|
||||
@@ -138,8 +138,8 @@ func RetryTask(ctx context.Context, id uint64) (string, error) {
|
||||
return "", fmt.Errorf(errTaskExecutionNotFound, err)
|
||||
}
|
||||
|
||||
if execution.Status != model.TaskExecutionStatusFailed {
|
||||
return "", fmt.Errorf(errRetryOnlyFailedTask, execution.Status)
|
||||
if execution.Status == model.TaskExecutionStatusPending || execution.Status == model.TaskExecutionStatusRunning {
|
||||
return "", fmt.Errorf(errRetryOnlyFinishedTask, execution.Status)
|
||||
}
|
||||
|
||||
if !execution.Retryable {
|
||||
@@ -299,6 +299,7 @@ func completeTaskExecution(ctx context.Context, execution *model.TaskExecution,
|
||||
span.RecordError(execErr)
|
||||
} else {
|
||||
execution.Status = model.TaskExecutionStatusSucceeded
|
||||
execution.ErrorMessage = "" // 清除历史重试失败遗留的错误信息
|
||||
if result != nil {
|
||||
execution.Result = result.Message
|
||||
if result.Detail != "" {
|
||||
|
||||
Reference in New Issue
Block a user