wangzhibo
2026-07-18 167510c1019f2c72ce9a544daec91aa5fe5409a5
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
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
import { Provide, Inject } from '@midwayjs/core';
import { BaseService } from '@cool-midway/core';
import { InjectEntityModel } from '@midwayjs/typeorm';
import { Repository, In } from 'typeorm';
import { ReportJoblogEntity } from '../entity/joblog';
import { TaskSiteService } from './tasksiteservice'
import { TaskDepartmentService } from './taskdepartmentservice'
 
/**
 *  定时任务-汇总站点/各个单位的各类垃圾汇总数据
 *  总入口,创建和调度任务,具体执行调用TaskSiteService 和  TaskDepartmentService
 */
 
 
@Provide()
export class ReportTaskService extends BaseService {
  @InjectEntityModel(ReportJoblogEntity)
  reportJoblogEntity: Repository<ReportJoblogEntity>;
 
  // 注入你具体的业务 Service,用于跑真正的汇总 SQL
  @Inject()
  taskSiteService: TaskSiteService;
 
  @Inject()
  departmentCubeService: TaskDepartmentService;
 
  /**
   * 外部定时任务(Crontab)每小时调用的总入口
   * 建议表达式: 0 5 * * * * (每小时的第5分钟执行)
   */
  async clockIn() {
    // 1. 自动看门狗:确保昨天和今天的任务实例已经在表里初始化了
    await this.ensureJobInstancesCreated();
 
    // 2. 执行引擎:捞出需要执行或重试的任务
    await this.executePendingJobs();
 
    return '任务执行成功:ReportTaskService';
  }
 
  /**
   * 检查并初始化任务实例(守护动作)
   */
  private async ensureJobInstancesCreated() {
    const now = new Date();
    const year = now.getFullYear();
    const month = String(now.getMonth() + 1).padStart(2, '0');
    const day = String(now.getDate()).padStart(2, '0');
    const yesterdayDate = new Date(now.getTime() - 24 * 60 * 60 * 1000);
    const yYear = yesterdayDate.getFullYear();
    const yMonth = String(yesterdayDate.getMonth() + 1).padStart(2, '0');
    const yDay = String(yesterdayDate.getDate()).padStart(2, '0');
 
    const todayStr = `${year}-${month}-${day}`;
    const yesterdayStr = `${yYear}-${yMonth}-${yDay}`;
 
    // 我们需要守护的目标日期和任务类型组合
    const targets = [
      { date: yesterdayStr, type: 'SITE_DAILY',isToday: false},
      { date: yesterdayStr, type: 'DEPT_DAILY',isToday: false},
      { date: todayStr, type: 'SITE_DAILY' ,isToday: true},
      { date: todayStr, type: 'DEPT_DAILY' ,isToday: true},
    ];
 
    for (const target of targets) {
      try {
        // 利用 try-catch 和数据库唯一索引,防并发同时也能实现“不存在则创建”
        const exist = await this.reportJoblogEntity.findOneBy({
          jobName: target.type,
          dataDate: target.date,
        });
        if (!exist) {
          const job = new ReportJoblogEntity();
          job.jobName = target.type;
          job.dataDate = target.date;
          job.status = 'INIT';
          await this.reportJoblogEntity.save(job);
        }else if (target.isToday && exist.status === 'SUCCESS') {
          // 2. 【核心修复】如果是今天的任务,且已经是 SUCCESS 状态,强行洗回 INIT 参与本小时的重算
          await this.reportJoblogEntity.update(exist.id, {
            status: 'INIT'
          });
        }
      } catch (err) {
        // 并发写入时可能触发唯一索引冲突,直接忽略即可
      }
    }
  }
 
  /**
   * 捞出 INIT 和 FAILED 的任务,并串行严格按顺序执行
   */
  private async executePendingJobs() {
    // 捞出所有待处理任务
    const pendingJobs = await this.reportJoblogEntity.find({
      where: {
        status: In(['INIT', 'FAILED']),
      },
      order: {
        dataDate: 'ASC', // 先算历史,后算今天
        // 确保在同一天内,SITE_DAILY 先跑,DEPT_DAILY 后跑
        // 在 JS 数组中我们可以进一步控制排序,或者由业务逻辑严格控制依赖
      },
    });
 
    for (const job of pendingJobs) {
      // 严格依赖检查:如果是部门日表任务,必须确保同天的站点日表任务已经是 SUCCESS
      if (job.jobName === 'DEPT_DAILY') {
        const siteJob = await this.reportJoblogEntity.findOneBy({
          dataDate: job.dataDate,
          jobName: 'SITE_DAILY',
        });
        if (!siteJob || siteJob.status !== 'SUCCESS') {
          // 源头明细还没算好/或还在运行,部门表任务必须等待,跳过本次执行
          continue;
        }
      }
 
      // 锁定任务状态为 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 ;
 
        if (job.jobName === 'SITE_DAILY') {
          // 调用你写的 站点日表 提取逻辑
          recordsProcessed = await this.taskSiteService.aggregateDailyData(job.dataDate) || 0;
        } else if (job.jobName === 'DEPT_DAILY') {
          // 调用你写的 部门树轰炸 提取逻辑
          recordsProcessed = await this.departmentCubeService.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',
        });
      }
    }
  }
}