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 TaskDepartmentService 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;
|
}
|
}
|