From 8fde55790ee1f0c1d04f7138a57449182e77f4d0 Mon Sep 17 00:00:00 2001 From: gaoqr <13665037151@163.com> Date: Tue, 25 Mar 2025 10:35:16 +0800 Subject: [PATCH] =?UTF-8?q?=E6=9D=BF=E6=9D=90=E3=80=81=E9=85=8D=E4=BB=B6?= =?UTF-8?q?=E5=AF=BC=E5=85=A5=E4=BC=98=E5=8C=96?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../core/injector/BaseBatchMapper.java | 8 + .../core/injector/MybatisPlusInjector.java | 1 + .../core/injector/UpdateBatchMethod.java | 80 +++++++ .../executor/enums/ErrorCodeConstants.java | 9 +- .../plate/vo/plate/PlateImportExcelVO.java | 2 +- .../mysql/orderParts/PartsBatchMapper.java | 15 ++ .../dal/mysql/plate/PlateGoodBatchMapper.java | 15 ++ .../ExecutorThreadPoolConfiguration.java | 30 +++ .../manage/parts/PartsServiceImpl.java | 201 +++++++++++++++-- .../manage/plate/PlateManageService.java | 2 +- .../manage/plate/PlateManageServiceImpl.java | 211 +++++++++++++++--- 11 files changed, 512 insertions(+), 62 deletions(-) create mode 100644 cf-framework/cf-spring-boot-starter-mybatis/src/main/java/com/cf/imes/framework/mybatis/core/injector/UpdateBatchMethod.java create mode 100644 cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/dal/mysql/orderParts/PartsBatchMapper.java create mode 100644 cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/dal/mysql/plate/PlateGoodBatchMapper.java create mode 100644 cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/framework/executor/config/ExecutorThreadPoolConfiguration.java diff --git a/cf-framework/cf-spring-boot-starter-mybatis/src/main/java/com/cf/imes/framework/mybatis/core/injector/BaseBatchMapper.java b/cf-framework/cf-spring-boot-starter-mybatis/src/main/java/com/cf/imes/framework/mybatis/core/injector/BaseBatchMapper.java index e3d3c4f49..048f5d699 100644 --- a/cf-framework/cf-spring-boot-starter-mybatis/src/main/java/com/cf/imes/framework/mybatis/core/injector/BaseBatchMapper.java +++ b/cf-framework/cf-spring-boot-starter-mybatis/src/main/java/com/cf/imes/framework/mybatis/core/injector/BaseBatchMapper.java @@ -19,4 +19,12 @@ public interface BaseBatchMapper extends BaseMapper { * @return */ int insertBatchSomeColumn(List entityList); + + /** + * 批量更新,update语句一次性提交,区别于update轮训执行 + * + * @param entityList + * @return + */ + int updateBatch(List entityList); } diff --git a/cf-framework/cf-spring-boot-starter-mybatis/src/main/java/com/cf/imes/framework/mybatis/core/injector/MybatisPlusInjector.java b/cf-framework/cf-spring-boot-starter-mybatis/src/main/java/com/cf/imes/framework/mybatis/core/injector/MybatisPlusInjector.java index bedeee563..67118ecd8 100644 --- a/cf-framework/cf-spring-boot-starter-mybatis/src/main/java/com/cf/imes/framework/mybatis/core/injector/MybatisPlusInjector.java +++ b/cf-framework/cf-spring-boot-starter-mybatis/src/main/java/com/cf/imes/framework/mybatis/core/injector/MybatisPlusInjector.java @@ -17,6 +17,7 @@ public class MybatisPlusInjector extends DefaultSqlInjector { public List getMethodList(Class mapperClass, TableInfo tableInfo) { List methodList = super.getMethodList(mapperClass, tableInfo); methodList.add(new InsertBatchSomeColumn(i -> i.getFieldFill() != FieldFill.UPDATE)); + methodList.add(new UpdateBatchMethod("updateBatch")); return methodList; } } diff --git a/cf-framework/cf-spring-boot-starter-mybatis/src/main/java/com/cf/imes/framework/mybatis/core/injector/UpdateBatchMethod.java b/cf-framework/cf-spring-boot-starter-mybatis/src/main/java/com/cf/imes/framework/mybatis/core/injector/UpdateBatchMethod.java new file mode 100644 index 000000000..44dd8ed8b --- /dev/null +++ b/cf-framework/cf-spring-boot-starter-mybatis/src/main/java/com/cf/imes/framework/mybatis/core/injector/UpdateBatchMethod.java @@ -0,0 +1,80 @@ +package com.cf.imes.framework.mybatis.core.injector; + +import com.baomidou.mybatisplus.core.injector.AbstractMethod; +import com.baomidou.mybatisplus.core.metadata.TableInfo; +import org.apache.ibatis.mapping.MappedStatement; +import org.apache.ibatis.mapping.SqlSource; + +/** + * 自定义批量更新 + * + * @author Gqr + * @since 2025/3/24 11:26 + */ +public class UpdateBatchMethod extends AbstractMethod { + + protected UpdateBatchMethod(String methodName) { + super(methodName); + } + + @Override + public MappedStatement injectMappedStatement(Class mapperClass, Class modelClass, TableInfo tableInfo) { + // 1. 生成动态 SET 子句(属性为空时不更新) + String setSql = generateConditionalSetSql(tableInfo); + + // 2. 处理乐观锁和逻辑删除的附加条件 + String additional = tableInfo.isWithVersion() + ? tableInfo.getVersionFieldInfo().getVersionOli("item", "item.") + : "" + tableInfo.getLogicDeleteSql(true, true); + + // 3. 构建完整 SQL + String sqlTemplate = ""; + String sqlResult = String.format( + sqlTemplate, + tableInfo.getTableName(), + setSql, + tableInfo.getKeyColumn(), + tableInfo.getKeyProperty(), + additional + ); + + // 4. 创建 SqlSource 和 MappedStatement + SqlSource sqlSource = languageDriver.createSqlSource(configuration, sqlResult, modelClass); + return addUpdateMappedStatement(mapperClass, modelClass, "updateBatch", sqlSource); + } + + /** + * 生成动态 SET 子句:属性为空时不更新 + */ + private String generateConditionalSetSql(TableInfo tableInfo) { + StringBuilder setSqlBuilder = new StringBuilder(""); + setSqlBuilder.append(""); + + String keyColumn = tableInfo.getKeyColumn(); + + // 遍历所有字段(排除主键、逻辑删除字段、版本字段、organid租户字段) + tableInfo.getFieldList().stream() + .filter(f -> !f.getColumn().equals(keyColumn) + && + !f.isLogicDelete() + && + (tableInfo.isWithVersion() ? !f.isVersion() : true) + && + !"organ_id".equals(f.getColumn()) + + ) + .forEach(field -> { + String column = field.getColumn(); + String property = "item." + field.getProperty(); + // 添加条件判断:属性不为空时更新 + setSqlBuilder.append( + String.format("%s = #{%s},", property, column, property) + ); + }); + + setSqlBuilder.append(""); + return setSqlBuilder.toString(); + } +} \ No newline at end of file diff --git a/cf-module-prod-executor/cf-module-prod-executor-api/src/main/java/com/cf/imes/module/executor/enums/ErrorCodeConstants.java b/cf-module-prod-executor/cf-module-prod-executor-api/src/main/java/com/cf/imes/module/executor/enums/ErrorCodeConstants.java index 0d61a360c..0faf33a0d 100644 --- a/cf-module-prod-executor/cf-module-prod-executor-api/src/main/java/com/cf/imes/module/executor/enums/ErrorCodeConstants.java +++ b/cf-module-prod-executor/cf-module-prod-executor-api/src/main/java/com/cf/imes/module/executor/enums/ErrorCodeConstants.java @@ -96,14 +96,19 @@ public class ErrorCodeConstants { public static final ErrorCode ORDER_IMPORT_ROOMCODE_NOT_EXIST_VALID_ERROR = new ErrorCode(1_002_031_014, "导入校验异常:板件{}下房间编码{}不存在于房间编码列表中"); public static final ErrorCode ORDER_IMPORT_BODYCODE_NOT_EXIST_VALID_ERROR = new ErrorCode(1_002_031_015, "导入校验异常:板件{}下柜体编码{}不存在于柜体编码列表中"); - // ========== manage配件 1_004_001_000 ========== + // ========== manage板材/配件 1_004_001_000 ========== public static final ErrorCode PARTS_NOT_EXISTS = new ErrorCode(1_004_001_000, "配件不存在"); public static final ErrorCode PARTS_GOOD_ID_EXISTS = new ErrorCode(1_004_001_001, "配件编号已存在"); public static final ErrorCode PARTS_UNIT_NOT_NULL = new ErrorCode(1_004_001_002, "配件-组件或五金或封边条,单位禁止为空"); public static final ErrorCode PARTS_EDGE_NOT_NULL = new ErrorCode(1_004_001_003, "配件-封边条颜色、宽厚禁止为空"); - public static final ErrorCode FILE_NULL = new ErrorCode(1_004_001_004, "上传文件为空"); public static final ErrorCode FILE_FORMAT_ERROR = new ErrorCode(1_004_001_005, "导入文件格式有误"); public static final ErrorCode FILE_EXCEED_SIZE = new ErrorCode(1_004_001_006, "文件大小超过 10M"); public static final ErrorCode FILE_CONTENT_NULL = new ErrorCode(1_004_001_007, "文件大小为空"); + public static final ErrorCode PARTS_IMPORT_QUERY_INTERRUPT_ERROR = new ErrorCode(1_004_001_008, "配件导入查询异常中断,请稍后再试"); + public static final ErrorCode PARTS_IMPORT_INTERRUPT_ERROR = new ErrorCode(1_004_001_009, "配件导入异常中断,请稍后再试"); + + // ========== manage板材 1_003_005_000 ========== + public static final ErrorCode PLATE_IMPORT_QUERY_INTERRUPT_ERROR = new ErrorCode(1_003_005_000, "配件导入查询异常中断,请稍后再试"); + public static final ErrorCode PLATE_IMPORT_INTERRUPT_ERROR = new ErrorCode(1_003_005_001, "配件导入异常中断,请稍后再试"); } diff --git a/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/controller/admin/manage/plate/vo/plate/PlateImportExcelVO.java b/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/controller/admin/manage/plate/vo/plate/PlateImportExcelVO.java index d2bdc8bdb..c392e9e33 100644 --- a/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/controller/admin/manage/plate/vo/plate/PlateImportExcelVO.java +++ b/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/controller/admin/manage/plate/vo/plate/PlateImportExcelVO.java @@ -75,7 +75,7 @@ public class PlateImportExcelVO { public boolean isValidGoods() { return !(goodsId == null || goodsId.isEmpty() || goodsName == null || goodsName.isEmpty() || material == null || - material.isEmpty() || color == null || color.isEmpty() ||thickness.isEmpty() ); + material.isEmpty() || color == null || color.isEmpty() || thickness == null || thickness.isEmpty()); } // 数据长度判断 diff --git a/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/dal/mysql/orderParts/PartsBatchMapper.java b/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/dal/mysql/orderParts/PartsBatchMapper.java new file mode 100644 index 000000000..0459c4898 --- /dev/null +++ b/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/dal/mysql/orderParts/PartsBatchMapper.java @@ -0,0 +1,15 @@ +package com.cf.imes.module.executor.dal.mysql.orderParts; + +import com.cf.imes.framework.mybatis.core.injector.BaseBatchMapper; +import com.cf.imes.module.executor.dal.dataobject.orderParts.PartsDO; +import org.apache.ibatis.annotations.Mapper; + +/** + * 配件批量mapper + * + * @author Gqr + * @since 2025/3/24 10:44 + */ +@Mapper +public interface PartsBatchMapper extends BaseBatchMapper { +} diff --git a/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/dal/mysql/plate/PlateGoodBatchMapper.java b/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/dal/mysql/plate/PlateGoodBatchMapper.java new file mode 100644 index 000000000..39e00dd5b --- /dev/null +++ b/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/dal/mysql/plate/PlateGoodBatchMapper.java @@ -0,0 +1,15 @@ +package com.cf.imes.module.executor.dal.mysql.plate; + +import com.cf.imes.framework.mybatis.core.injector.BaseBatchMapper; +import com.cf.imes.module.executor.dal.dataobject.plate.PlateGoodDO; +import org.apache.ibatis.annotations.Mapper; + +/** + * 板材批量mapper + * + * @author Gqr + * @since 2025/3/25 9:47 + */ +@Mapper +public interface PlateGoodBatchMapper extends BaseBatchMapper { +} diff --git a/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/framework/executor/config/ExecutorThreadPoolConfiguration.java b/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/framework/executor/config/ExecutorThreadPoolConfiguration.java new file mode 100644 index 000000000..545ee6e63 --- /dev/null +++ b/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/framework/executor/config/ExecutorThreadPoolConfiguration.java @@ -0,0 +1,30 @@ +package com.cf.imes.module.executor.framework.executor.config; + +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; + +import java.util.concurrent.ThreadPoolExecutor; + +/** + * @author Gqr + * @since 2025/3/24 15:21 + */ +@Configuration(proxyBeanMethods = false) +public class ExecutorThreadPoolConfiguration { + public static final String EXECUTOR_IMPOT_THREAD_POOL_TASK_EXECUTOR = "EXECUTOR_IMPORT_THREAD_POOL_TASK_EXECUTOR"; + + @Bean(EXECUTOR_IMPOT_THREAD_POOL_TASK_EXECUTOR) + public ThreadPoolTaskExecutor notifyThreadPoolTaskExecutor() { + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setCorePoolSize(8); // 设置核心线程数 + executor.setMaxPoolSize(8); // 设置最大线程数 + executor.setKeepAliveSeconds(60); // 设置空闲时间 + executor.setQueueCapacity(100); // 设置队列大小 + executor.setThreadNamePrefix("executor-import-Executor-"); // 配置线程池的前缀 + executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); + // 进行加载 + executor.initialize(); + return executor; + } +} diff --git a/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/service/manage/parts/PartsServiceImpl.java b/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/service/manage/parts/PartsServiceImpl.java index 004ef0313..8db7eb37d 100644 --- a/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/service/manage/parts/PartsServiceImpl.java +++ b/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/service/manage/parts/PartsServiceImpl.java @@ -1,5 +1,9 @@ package com.cf.imes.module.executor.service.manage.parts; +import cn.hutool.core.util.ObjectUtil; +import com.baomidou.dynamic.datasource.toolkit.DynamicDataSourceContextHolder; +import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; +import com.cf.imes.framework.common.exception.ServiceException; import com.cf.imes.framework.common.pojo.CommonResult; import com.cf.imes.framework.common.pojo.PageResult; import com.cf.imes.framework.common.util.object.BeanUtils; @@ -10,6 +14,7 @@ import com.cf.imes.module.executor.controller.admin.manage.parts.vo.PartsImportE import com.cf.imes.module.executor.controller.admin.manage.parts.vo.PartsPageReqVO; import com.cf.imes.module.executor.controller.admin.manage.parts.vo.PartsSaveReqVO; import com.cf.imes.module.executor.dal.dataobject.orderParts.PartsDO; +import com.cf.imes.module.executor.dal.mysql.orderParts.PartsBatchMapper; import com.cf.imes.module.executor.dal.mysql.orderParts.PartsMapper; import com.cf.imes.module.executor.enums.manage.parts.CategoryEnum; import com.cf.imes.module.executor.util.file.FileHelperUtil; @@ -17,6 +22,9 @@ import com.cf.imes.module.executor.util.number.NumberHelperUtil; import com.cf.imes.module.system.api.organ.OrganApi; import com.cf.imes.module.system.api.permission.PermissionApi; import com.cf.imes.module.system.enums.ErrorCodeConstants; +import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang.StringUtils; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; import org.springframework.validation.annotation.Validated; @@ -25,10 +33,18 @@ import javax.annotation.Resource; import javax.servlet.http.HttpServletResponse; import java.util.ArrayList; import java.util.List; +import java.util.Map; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.function.Consumer; +import java.util.stream.Collectors; + import static com.cf.imes.framework.common.exception.util.ServiceExceptionUtil.exception; import static com.cf.imes.framework.security.core.util.SecurityFrameworkUtils.getLoginUser; import static com.cf.imes.module.executor.enums.ErrorCodeConstants.*; +import static com.cf.imes.module.executor.framework.executor.config.ExecutorThreadPoolConfiguration.EXECUTOR_IMPOT_THREAD_POOL_TASK_EXECUTOR; /** @@ -38,11 +54,15 @@ import static com.cf.imes.module.executor.enums.ErrorCodeConstants.*; */ @Service @Validated +@Slf4j public class PartsServiceImpl implements PartsService { @Resource private PartsMapper partsMapper; + @Resource + private PartsBatchMapper partsBatchMapper; + @Resource private OrganApi organApi; @@ -52,12 +72,16 @@ public class PartsServiceImpl implements PartsService { @Resource private NumberHelperUtil numberUtil; + @Resource(name = EXECUTOR_IMPOT_THREAD_POOL_TASK_EXECUTOR) + private ThreadPoolTaskExecutor threadPoolTaskExecutor; + private final static String BASE_PARTS_IMPORT_PERMISSION = "base:accessory:import"; // 导入板材 private final static String BASE_PARTS_DELETE_PERMISSION = "base:accessory:delete"; // 删除板材信息 private final static String BASE_PARTS_CREATE_PERMISSION = "base:accessory:create"; // 创建板材信息 private final static String BASE_PARTS_UPDATE_PERMISSION = "base:accessory:update"; // 更新板材信息 private final static String BASE_PARTS_QUERY_PERMISSION = "base:accessory:query"; // 获得板材信息 + private static final String EMPTY_STRING = new String(); @Override @OrganIgnore @@ -87,14 +111,14 @@ public class PartsServiceImpl implements PartsService { Long organId = getOrganId(BASE_PARTS_UPDATE_PERMISSION, getLoginUser().getId(), updateReqVO.getOrganId()); // 校验id存在 - boolean index = validateExistsById(updateReqVO.getId(),organId); + boolean index = validateExistsById(updateReqVO.getId(), organId); if (!index) { throw exception(PARTS_NOT_EXISTS); } // 校验goodsId是否存在 - if (!validateExistsExceptGoodsId(updateReqVO.getId(),updateReqVO.getGoodsId(),organId)) { + if (!validateExistsExceptGoodsId(updateReqVO.getId(), updateReqVO.getGoodsId(), organId)) { throw exception(PARTS_GOOD_ID_EXISTS); } // 更新 @@ -179,40 +203,175 @@ public class PartsServiceImpl implements PartsService { @Override @OrganIgnore + @Transactional(rollbackFor = Exception.class) public void importPartsList(List importParts, boolean isUpdateSupport, Long organId) { // 组织id organId = getOrganId(BASE_PARTS_IMPORT_PERMISSION, getLoginUser().getId(), organId); // 批量插入集合 - List insertList = new ArrayList<>(); + List insertList = new CopyOnWriteArrayList<>(); // 批量更新集合 - List updateList = new ArrayList<>(); + List updateList = new CopyOnWriteArrayList<>(); - Long finalOrganId = organId; + // 每100条 goodsId in查询一次 + List> batches = createBatches(importParts, 100); + boolean hasError = processBatchesConcurrently(batches, organId, isUpdateSupport, insertList, updateList); - importParts.forEach(importPart -> { - // 判断如果不存在,在进行插入 - PartsDO existPlate = partsMapper.selectByGoodID(importPart.getGoodsId(), finalOrganId); - if (existPlate == null) { // 不存在 - insertList.add(BeanUtils.toBean(importPart, PartsDO.class).setOrganId(finalOrganId)); - return; - } - // 判断是否允许更新 - if (!isUpdateSupport) { - PartsDO partsDO = BeanUtils.toBean(importPart, PartsDO.class).setOrganId(finalOrganId).setId(existPlate.getId()); - updateList.add(partsDO); - } - }); -// 批量插入,批量更新 + if (hasError) { + throw new ServiceException(PARTS_IMPORT_QUERY_INTERRUPT_ERROR); + } + + // 每100条插入/更新一次,避免语句过大 if (!insertList.isEmpty()) { - partsMapper.insertBatch(insertList); + batchProcess(insertList, 100, partsBatchMapper::insertBatchSomeColumn); } if (!updateList.isEmpty()) { - partsMapper.updateBatch(updateList); + batchProcess(updateList, 100, partsBatchMapper::updateBatch); } } + /** + * 处理excel配件批次数据,多线程查询是新增还是更新 + * + * @param batches + * @param organId + * @param isUpdateSupport + * @param insertList + * @param updateList + * @return + */ + private boolean processBatchesConcurrently(List> batches, Long organId, + boolean isUpdateSupport, List insertList, + List updateList) { + // 查询过程中是否异常了 + AtomicBoolean hasError = new AtomicBoolean(false); + // 计数器 + CountDownLatch latch = new CountDownLatch(batches.size()); + // 主线程的data_code + String peek = DynamicDataSourceContextHolder.peek(); + + for (List batch : batches) { + threadPoolTaskExecutor.submit(() -> { + threadPoolTaskExecutor.submit(() -> { + try { + // 子线程同步主线程的多数据源 + DynamicDataSourceContextHolder.push(peek); + + List goodsIds = batch.stream().map(PartsImportExcelVO::getGoodsId).collect(Collectors.toList()); + Map partsMap = partsMapper.selectList( + new LambdaQueryWrapper().in(PartsDO::getGoodsId, goodsIds).eq(PartsDO::getOrganId, organId)) + .stream().collect(Collectors.toMap(PartsDO::getGoodsId, p -> p)); + + for (PartsImportExcelVO importPart : batch) { + PartsDO existPlate = partsMap.get(importPart.getGoodsId()); + if (existPlate == null) { + PartsDO partsDO = BeanUtils.toBean(importPart, PartsDO.class).setOrganId(organId); + partsDO.setIsComposite(ObjectUtil.isNull(partsDO.getIsComposite()) ? false : partsDO.getIsComposite()) + .setCategory(getEmpty(partsDO.getCategory())) + .setSubparts(getEmpty(partsDO.getType())) + .setColor(getEmpty(partsDO.getColor())) + .setFactory(getEmpty(partsDO.getFactory())) + .setRemark(getEmpty(partsDO.getRemark())) + .setBrand(getEmpty(partsDO.getBrand())) + .setModel(getEmpty(partsDO.getModel())) + .setSpec(getEmpty(partsDO.getSpec())) + .setMaterial(getEmpty(partsDO.getMaterial())) + .setLength(getZeroDouble(partsDO.getLength())) + .setWidth(getZeroDouble(partsDO.getWidth())) + .setThickness(getZeroDouble(partsDO.getThickness())) + .setPrice(getZeroDouble(partsDO.getPrice())) + .setUnit(getEmpty(partsDO.getUnit())) + .setOrganId(organId) + .setDeleted(false); + insertList.add(partsDO); + } else if (!isUpdateSupport) { + PartsDO partsDO = BeanUtils.toBean(importPart, PartsDO.class) + .setOrganId(organId) + .setId(existPlate.getId()); + updateList.add(partsDO); + } + } + } catch (Exception e) { + boolean b = hasError.get(); + if (!b) { + hasError.set(true); + } + } finally { + latch.countDown(); // 任务完成,计数器减1 + } + }); + }); + } + + try { + // 等待所有任务查询完成 + latch.await(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new ServiceException(PARTS_IMPORT_INTERRUPT_ERROR); + } + return hasError.get(); + } + + /** + * 每batchSize条数据处理一次DO + * + * @param dataList + * @param batchSize + * @param processor + */ + private void batchProcess(List dataList, int batchSize, Consumer> processor) { + if (dataList == null || dataList.isEmpty()) { + return; + } + + int total = dataList.size(); + for (int i = 0; i < total; i += batchSize) { + int end = Math.min(i + batchSize, total); + List subList = dataList.subList(i, end); + + if (!subList.isEmpty()) { + processor.accept(subList); + } + } + } + + /** + * 拆分导入数据为batchSize一组 + * + * @param list + * @param batchSize + * @return + */ + private List> createBatches(List list, int batchSize) { + List> batches = new ArrayList<>(); + for (int i = 0; i < list.size(); i += batchSize) { + batches.add(list.subList(i, Math.min(i + batchSize, list.size()))); + } + return batches; + } + + /** + * value为空返回EMPTY_STRING + * + * @param value + * @return + */ + private String getEmpty(String value) { + return StringUtils.defaultIfEmpty(value, EMPTY_STRING); + } + + /** + * value为空返回0.0 + * + * @param value + * @return + */ + private double getZeroDouble(Double value) { + return ObjectUtil.defaultIfNull(value, 0.0); + } + @Override public List getPartAll(String category) { return partsMapper.getAllByOrganId(OrganContextHolder.getOrganId(), category); diff --git a/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/service/manage/plate/PlateManageService.java b/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/service/manage/plate/PlateManageService.java index ccaecac11..3688d2559 100644 --- a/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/service/manage/plate/PlateManageService.java +++ b/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/service/manage/plate/PlateManageService.java @@ -81,7 +81,7 @@ public interface PlateManageService { * @param isUpdateSupport 是否支持更新 * @return 导入结果 */ - PlateImportRespVO importPlateList(List importUsers, boolean isUpdateSupport , Long organId); + void importPlateList(List importUsers, boolean isUpdateSupport , Long organId); /** * 生产导入文件模板下载 diff --git a/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/service/manage/plate/PlateManageServiceImpl.java b/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/service/manage/plate/PlateManageServiceImpl.java index 85cc121aa..d28bdf42b 100644 --- a/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/service/manage/plate/PlateManageServiceImpl.java +++ b/cf-module-prod-executor/cf-module-prod-executor-biz/src/main/java/com/cf/imes/module/executor/service/manage/plate/PlateManageServiceImpl.java @@ -1,21 +1,26 @@ package com.cf.imes.module.executor.service.manage.plate; +import cn.hutool.core.util.ObjectUtil; +import com.baomidou.dynamic.datasource.toolkit.DynamicDataSourceContextHolder; +import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; +import com.cf.imes.framework.common.exception.ServiceException; import com.cf.imes.framework.common.pojo.CommonResult; import com.cf.imes.framework.common.pojo.PageResult; import com.cf.imes.framework.common.util.object.BeanUtils; import com.cf.imes.framework.organ.core.aop.OrganIgnore; import com.cf.imes.framework.organ.core.context.OrganContextHolder; import com.cf.imes.module.executor.controller.admin.manage.plate.vo.plate.PlateImportExcelVO; -import com.cf.imes.module.executor.controller.admin.manage.plate.vo.plate.PlateImportRespVO; import com.cf.imes.module.executor.controller.admin.manage.plate.vo.plate.PlatePageReqVO; import com.cf.imes.module.executor.controller.admin.manage.plate.vo.plate.PlateSaveReqVO; import com.cf.imes.module.executor.dal.dataobject.plate.PlateGoodDO; +import com.cf.imes.module.executor.dal.mysql.plate.PlateGoodBatchMapper; import com.cf.imes.module.executor.dal.mysql.plate.PlateGoodMapper; -import com.cf.imes.module.executor.dal.mysql.plate.PlateMapper; import com.cf.imes.module.executor.util.file.FileHelperUtil; import com.cf.imes.module.system.api.organ.OrganApi; import com.cf.imes.module.system.api.permission.PermissionApi; import com.cf.imes.module.system.enums.ErrorCodeConstants; +import org.apache.commons.lang.StringUtils; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; import org.springframework.validation.annotation.Validated; @@ -23,11 +28,19 @@ import org.springframework.validation.annotation.Validated; import javax.annotation.Resource; import javax.servlet.http.HttpServletResponse; import java.util.ArrayList; -import java.util.LinkedHashMap; import java.util.List; +import java.util.Map; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.function.Consumer; +import java.util.stream.Collectors; import static com.cf.imes.framework.common.exception.util.ServiceExceptionUtil.exception; import static com.cf.imes.framework.security.core.util.SecurityFrameworkUtils.getLoginUser; +import static com.cf.imes.module.executor.enums.ErrorCodeConstants.PLATE_IMPORT_INTERRUPT_ERROR; +import static com.cf.imes.module.executor.enums.ErrorCodeConstants.PLATE_IMPORT_QUERY_INTERRUPT_ERROR; +import static com.cf.imes.module.executor.framework.executor.config.ExecutorThreadPoolConfiguration.EXECUTOR_IMPOT_THREAD_POOL_TASK_EXECUTOR; import static com.cf.imes.module.system.enums.ErrorCodeConstants.PLATE_GOODS_ID_NOT_SAME; import static com.cf.imes.module.system.enums.ErrorCodeConstants.REMAIN_PLATE_NOT_EXISTS; @@ -49,12 +62,20 @@ public class PlateManageServiceImpl implements PlateManageService { @Resource private PermissionApi permissionApi; + @Resource + private PlateGoodBatchMapper plateGoodBatchMapper; + + @Resource(name = EXECUTOR_IMPOT_THREAD_POOL_TASK_EXECUTOR) + private ThreadPoolTaskExecutor threadPoolTaskExecutor; + private final static String BASE_PLATE_IMPORT_PERMISSION = "base:plate:import"; // 导入板材 private final static String BASE_PLATE_DELETE_PERMISSION = "base:plate:delete"; // 删除板材信息 private final static String BASE_PLATE_CREATE_PERMISSION = "base:plate:create"; // 创建板材信息 private final static String BASE_PLATE_UPDATE_PERMISSION = "base:plate:update"; // 更新板材信息 private final static String BASE_PLATE_QUERY_PERMISSION = "base:plate:query"; // 获得板材信息 + private static final String EMPTY_STRING = new String(); + @Override @OrganIgnore public Long createPlate(PlateSaveReqVO createReqVO) { @@ -140,45 +161,161 @@ public class PlateManageServiceImpl implements PlateManageService { @Override @OrganIgnore @Transactional(rollbackFor = Exception.class) // 添加事务,异常则回滚所有导入 - public PlateImportRespVO importPlateList(List importPlates, boolean isUpdateSupport, Long organId) { + public void importPlateList(List importPlates, boolean isUpdateSupport, Long organId) { organId = getOrganId(BASE_PLATE_IMPORT_PERMISSION, getLoginUser().getId(), organId); - PlateImportRespVO respVO = PlateImportRespVO.builder().createPlateNames(new ArrayList<>()) - .updatePlateNames(new ArrayList<>()).failurePlateNames(new LinkedHashMap<>()).build(); -// 批量插入集合 - List insertList = new ArrayList<>(); -// 批量更新集合 - List updateList = new ArrayList<>(); - Long finalOrganId = organId; - importPlates.forEach(importPlate -> { - // 判断如果不存在,在进行插入 - PlateGoodDO existPlate = plateGoodMapper.selectByGoodID(importPlate.getGoodsId(), finalOrganId); - if (existPlate == null) { - respVO.getCreatePlateNames().add(importPlate.getGoodsName()); - insertList.add(BeanUtils.toBean(importPlate, PlateGoodDO.class) - .setOrganId(finalOrganId)); -// .setTexture(Boolean.valueOf(importPlate.getTexture())); - return; - } - // 判断是否允许更新 - if (!isUpdateSupport) { - respVO.getFailurePlateNames().put(importPlate.getGoodsName(), ErrorCodeConstants.PLATE_EXISTS.getMsg()); - PlateGoodDO plateDO = BeanUtils.toBean(importPlate, PlateGoodDO.class) - .setOrganId(finalOrganId) - .setId(existPlate.getId()); -// .setTexture(Boolean.valueOf(importPlate.getTexture())); - updateList.add(plateDO); - } - }); -// 批量插入,批量更新 - if (insertList.size() > 0) { - plateGoodMapper.insertBatch(insertList); + // 批量插入集合 + List insertList = new CopyOnWriteArrayList<>(); + // 批量更新集合 + List updateList = new CopyOnWriteArrayList<>(); + + // 每100条 goodsId in查询一次 + List> batches = createBatches(importPlates, 100); + boolean hasError = processBatchesConcurrently(batches, organId, isUpdateSupport, insertList, updateList); + + if (hasError) { + throw new ServiceException(PLATE_IMPORT_QUERY_INTERRUPT_ERROR); } - if (updateList.size() > 0) { - plateGoodMapper.updateBatch(updateList); + + // 每100条插入/更新一次,避免语句过大 + if (!insertList.isEmpty()) { + batchProcess(insertList, 100, plateGoodBatchMapper::insertBatchSomeColumn); } - return respVO; + if (!updateList.isEmpty()) { + batchProcess(updateList, 100, plateGoodBatchMapper::updateBatch); + } + } + + /** + * 处理excel配件批次数据,多线程查询是新增还是更新 + * + * @param batches + * @param organId + * @param isUpdateSupport + * @param insertList + * @param updateList + * @return + */ + private boolean processBatchesConcurrently(List> batches, Long organId, + boolean isUpdateSupport, List insertList, + List updateList) { + // 查询过程中是否异常了 + AtomicBoolean hasError = new AtomicBoolean(false); + // 计数器 + CountDownLatch latch = new CountDownLatch(batches.size()); + // 主线程的data_code + String peek = DynamicDataSourceContextHolder.peek(); + + for (List batch : batches) { + threadPoolTaskExecutor.submit(() -> { + threadPoolTaskExecutor.submit(() -> { + try { + // 子线程同步主线程的多数据源 + DynamicDataSourceContextHolder.push(peek); + + List goodsIds = batch.stream().map(PlateImportExcelVO::getGoodsId).collect(Collectors.toList()); + Map plateGoodMap = plateGoodMapper.selectList( + new LambdaQueryWrapper().in(PlateGoodDO::getGoodsId, goodsIds).eq(PlateGoodDO::getOrganId, organId)) + .stream().collect(Collectors.toMap(PlateGoodDO::getGoodsId, p -> p)); + + for (PlateImportExcelVO importPart : batch) { + PlateGoodDO existPlate = plateGoodMap.get(importPart.getGoodsId()); + if (existPlate == null) { + PlateGoodDO plateGoodDO = BeanUtils.toBean(importPart, PlateGoodDO.class).setOrganId(organId); + plateGoodDO.setColorBlack(EMPTY_STRING) + .setBrand(EMPTY_STRING) + .setSpec(getEmpty(plateGoodDO.getSpec())) + .setRemark(getEmpty(plateGoodDO.getRemark())) + .setTexture(ObjectUtil.isNull(plateGoodDO.getTexture()) ? false : plateGoodDO.getTexture()) + .setOrganId(organId) + .setDeleted(false); + insertList.add(plateGoodDO); + } else if (!isUpdateSupport) { + PlateGoodDO plateGoodDO = BeanUtils.toBean(importPart, PlateGoodDO.class) + .setOrganId(organId) + .setId(existPlate.getId()); + updateList.add(plateGoodDO); + } + } + } catch (Exception e) { + boolean b = hasError.get(); + if (!b) { + hasError.set(true); + } + } finally { + latch.countDown(); // 任务完成,计数器减1 + } + }); + }); + } + + try { + // 等待所有任务查询完成 + latch.await(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new ServiceException(PLATE_IMPORT_INTERRUPT_ERROR); + } + return hasError.get(); + } + + /** + * 每batchSize条数据处理一次DO + * + * @param dataList + * @param batchSize + * @param processor + */ + private void batchProcess(List dataList, int batchSize, Consumer> processor) { + if (dataList == null || dataList.isEmpty()) { + return; + } + + int total = dataList.size(); + for (int i = 0; i < total; i += batchSize) { + int end = Math.min(i + batchSize, total); + List subList = dataList.subList(i, end); + + if (!subList.isEmpty()) { + processor.accept(subList); + } + } + } + + /** + * 拆分导入数据为batchSize一组 + * + * @param list + * @param batchSize + * @return + */ + private List> createBatches(List list, int batchSize) { + List> batches = new ArrayList<>(); + for (int i = 0; i < list.size(); i += batchSize) { + batches.add(list.subList(i, Math.min(i + batchSize, list.size()))); + } + return batches; + } + + /** + * value为空返回EMPTY_STRING + * + * @param value + * @return + */ + private String getEmpty(String value) { + return StringUtils.defaultIfEmpty(value, EMPTY_STRING); + } + + /** + * value为空返回0.0 + * + * @param value + * @return + */ + private double getZeroDouble(Double value) { + return ObjectUtil.defaultIfNull(value, 0.0); } @Override