wangzhibo
4 天以前 c19f1f7e21c93d1f13b4f44a8a615f65af577559
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
import { Provide, Inject } from '@midwayjs/core';
import { InjectEntityModel } from '@midwayjs/typeorm';
import { Repository, In } from 'typeorm';
import { ReportJoblogEntity } from '../entity/joblog';
import { TaskServiceSorter } from './taskservicesorter'
 
 
 
@Provide()
export class ExecutorSorterDaily {
 
    @InjectEntityModel(ReportJoblogEntity)
    reportJoblogEntity: Repository<ReportJoblogEntity>;
 
    @Inject()
    taskServiceSorter: TaskServiceSorter;
 
    async execute() {
        // 捞出所有待处理任务
        const pendingSiteJobs = await this.reportJoblogEntity.find({
            where: {
                status: In(['INIT', 'FAILED']),
                jobName: In(['SORTER_DAILY']),
            },
            order: {
                dataDate: 'ASC',
            },
        });
 
        for (const job of pendingSiteJobs) {
 
            // 锁定任务状态为 RUNNING(轻量级乐观锁锁机制)
            const lockResult = await this.reportJoblogEntity.update(
                { id: job.id, status: job.status }, // 带有原始状态检查,防止被别的线程抢跑
                { status: 'RUNNING', startTime: new Date() }
            );
 
            if (lockResult.affected === 0) {
                continue; // 没抢到锁,说明别的定时器在跑它,跳过
            }
 
            // 真正开始跑数据
            try {
                let recordsProcessed = 0;
 
 
                recordsProcessed = await this.taskServiceSorter.aggregateDailyData(job.dataDate) || 0;
 
 
                // 成功:更新状态
                await this.reportJoblogEntity.update(job.id, {
                    status: 'SUCCESS',
                    recordsProcessed: recordsProcessed,
                    endTime: new Date(),
                    errorMessage: '',
                });
 
            } catch (error) {
                // 失败:记录堆栈,等待下一小时重试
 
                var errorMessage = '';
                if (error instanceof Error) {
                    errorMessage = error.stack ?? error.message;
                }
                if (typeof error === 'string') {
                    errorMessage = error;
                }
                try {
                    errorMessage = JSON.stringify(error);
                } catch {
                    errorMessage = String(error);
                }
 
                await this.reportJoblogEntity.update(job.id, {
                    status: 'FAILED',
                    endTime: new Date(),
                    errorMessage: errorMessage,
                    failedCount: () => 'failedCount + 1',
                });
            }
        }
    }
 
}