wangzhibo
2026-07-28 920cec41cedfde27f89d7c19b78e03147b926ab5
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
import { Provide, Inject } from '@midwayjs/core';
import { BaseService } from '@cool-midway/core';
import { InjectEntityModel } from '@midwayjs/typeorm';
import { Repository } from 'typeorm';
import { ReportDailysiteEntity } from '../entity/dailysite';
import { ReportDailyDepartmentEntity } from '../entity/dailydepartment';
import { BaseSysDepartmentEntity } from '../../base/entity/sys/department'; 
import { BasicdataSiteEntity } from '../../basicdata/entity/site';
 
 
/**
 *  定时任务-子任务 汇总部门的各类垃圾数据
 */
@Provide()
export class TaskServiceDepartment extends BaseService {
  @InjectEntityModel(ReportDailysiteEntity)
  reportDailysiteEntity: Repository<ReportDailysiteEntity>;
 
  @InjectEntityModel(ReportDailyDepartmentEntity)
  reportDailyDepartmentEntity: Repository<ReportDailyDepartmentEntity>;
 
  @InjectEntityModel(BaseSysDepartmentEntity)
  baseSysDepartmentEntity: Repository<BaseSysDepartmentEntity>;
 
  @InjectEntityModel(BasicdataSiteEntity)
  basicdataSiteEntity: Repository<BasicdataSiteEntity>;
 
  /**
   * 部门日表无级递归汇总引擎入口
   * @param targetDate 目标计算日期格式 'YYYY-MM-DD'
   */
  async aggregateDailyData(targetDate: string) {
    // 1. 异步并行加载基础字典和当天已算好的站点日表数据,最大化I/O效率
    const [allDepts, allSites, dailySites] = await Promise.all([
      this.baseSysDepartmentEntity.find({ select: ['id', 'parentId', 'name'] }),
      this.basicdataSiteEntity.find({ select: ['id', 'departmentId'], where: { siteType: '收集点' } }),
      this.reportDailysiteEntity.findBy({ date: targetDate })
    ]);
 
    if (dailySites.length === 0) {
      // 如果当天没有任何站点产生流水,顺应“没数不存”原则,部门表也直接不存任何数据,完美收工
      return;
    }
 
    // 2. 建立 站点ID -> 部门ID 的映射快查表
    const siteToDeptMap = new Map<number, number>();
    allSites.forEach(s => {
      if (s.departmentId) siteToDeptMap.set(s.id, s.departmentId);
    });
 
    // 3. 构建部门树的关系链(用于无级向上追溯)
    const deptParentMap = new Map<number, number | null>();
    const deptNameMap = new Map<number, string>();
    allDepts.forEach(d => {
      deptParentMap.set(d.id, d.parentId || null);
      deptNameMap.set(d.id, d.name);
    });
 
    // 4. 定义需要累加的动态分类字段清单(必须和站点表算出来的50多个字段完全对齐)
    const typeSuffixes = [
      'Sw60', 'Sw61', 'Sw62', 'Sw63', 'Sw64',
      'Sw62001', 'Sw62002', 'Sw62003', 'Sw62004', 'Sw62005', 'Sw62006', 'Sw62007'
    ];
   
 
    // 5. 声明部门数据的内存聚合容器
    const deptAggMap = new Map<number, any>();
 
    // 初始化一个干净部门对象的辅助函数
    const initDeptData = (deptId: number) => ({
      date: targetDate,
      departmentId: deptId,
      departmentName: deptNameMap.get(deptId) || `未知部门(${deptId})`,
      weightTotal: 0,
      timeTotal: 0,
      carbonTotal: 0,
      estimatedSalesSw62: 0,
      estimatedSalesSw62001: 0,
      estimatedSalesSw62002: 0,
      estimatedSalesSw62003: 0,
      estimatedSalesSw62004: 0,
      estimatedSalesSw62005: 0,
      estimatedSalesSw62006: 0,
      estimatedSalesSw62007: 0,
      ...typeSuffixes.reduce((acc, suffix) => {
        acc[`weight${suffix}`] = 0;
        acc[`time${suffix}`] = 0;
        acc[`carbon${suffix}`] = 0;
        return acc;
      }, {})
    });
 
    // 6. 核心无级轰炸算法:遍历每一个有数据的站点,向它的所有上级父部门“冒泡”贡献数据
    for (const siteRow of dailySites) {
      const directDeptId = siteToDeptMap.get(siteRow.siteId);
      if (!directDeptId) continue; // 如果站点脱离了部门组织架构,跳过
 
      // 从当前直属部门开始,沿着 parentId 一路往上爬,直到爬到祖先根节点(parentId 为 null)
      let currentDeptId: number | null | undefined = directDeptId;
      
      while (currentDeptId !== null && currentDeptId !== undefined) {
        // 如果树中没有这个部门ID(异常数据防崩),中断退出
        if (!deptParentMap.has(currentDeptId)) break;
 
        // 如果内存中还没有这个部门的聚合空间,立即为其独立初始化
        if (!deptAggMap.has(currentDeptId)) {
          deptAggMap.set(currentDeptId, initDeptData(currentDeptId));
        }
 
        const deptAgg = deptAggMap.get(currentDeptId);
 
        // 核心轰炸:将该站点的各项数值,叠加到当前部门的所有对应指标上
        deptAgg.weightTotal += Number(siteRow.weightTotal) || 0;
        deptAgg.timeTotal += Number(siteRow.timeTotal) || 0;
        deptAgg.carbonTotal += Number(siteRow.carbonTotal) || 0;
 
        // 动态轰炸 50 多个子类字段
        typeSuffixes.forEach(suffix => {
          deptAgg[`weight${suffix}`] += Number(siteRow[`weight${suffix}`]) || 0;
          deptAgg[`time${suffix}`] += Number(siteRow[`time${suffix}`]) || 0;
          deptAgg[`carbon${suffix}`] += Number(siteRow[`carbon${suffix}`]) || 0;
 
          // 回收物子类的估算金额同步叠加
          if (suffix.startsWith('Sw62')) {
            deptAgg[`estimatedSales${suffix}`] += Number(siteRow[`estimatedSales${suffix}`]) || 0;
          }
        });
 
        // 【关键跃迁】:将指针移向它的亲生父节点,实现无级递归向上冒泡!
        currentDeptId = deptParentMap.get(currentDeptId);
      }
    }
 
    // 7. 将内存中所有被“轰炸”过、产生了有效数据的部门记录提取出来
    const finalDeptRows = Array.from(deptAggMap.values());
 
    if (finalDeptRows.length > 0) {
            // 1. 动态抓取所有的键,将主键、联合唯一键彻底从“待更新字段”中排除
            // 这样能确保 ON DUPLICATE KEY UPDATE 后面的赋值语句绝对干净
            const allUpdateFields = Object.keys(finalDeptRows[0]).filter(
                k => k !== 'date' && k !== 'departmentId' && k !== 'id'
            );
 
            // 2. 依然推荐采用 Chunk 分批,这是应对海量数据/多子字段最安全、执行效率最高的做法
            const chunkSize = 50;
            for (let i = 0; i < finalDeptRows.length; i += chunkSize) {
                const chunk = finalDeptRows.slice(i, i + chunkSize);
 
                await this.reportDailyDepartmentEntity.createQueryBuilder()
                    .insert()
                    .values(chunk)
                    // 💡 显式声明:当 ['date', 'departmentId'] 发生冲突时,强行把 allUpdateFields 里的字段全部更新一遍
                    .orUpdate(allUpdateFields, ['date', 'departmentId'])
                    .execute();
            }
        }
 
    return finalDeptRows.length;
  }
}