板材、配件导入优化

This commit is contained in:
gaoqr
2025-03-25 10:35:16 +08:00
parent d1450061d0
commit 8fde55790e
11 changed files with 512 additions and 62 deletions
@@ -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, "配件导入异常中断,请稍后再试");
}
@@ -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());
}
// 数据长度判断
@@ -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<PartsDO> {
}
@@ -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<PlateGoodDO> {
}
@@ -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;
}
}
@@ -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<PartsImportExcelVO> importParts, boolean isUpdateSupport, Long organId) {
// 组织id
organId = getOrganId(BASE_PARTS_IMPORT_PERMISSION, getLoginUser().getId(), organId);
// 批量插入集合
List<PartsDO> insertList = new ArrayList<>();
List<PartsDO> insertList = new CopyOnWriteArrayList<>();
// 批量更新集合
List<PartsDO> updateList = new ArrayList<>();
List<PartsDO> updateList = new CopyOnWriteArrayList<>();
Long finalOrganId = organId;
// 每100条 goodsId in查询一次
List<List<PartsImportExcelVO>> 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<List<PartsImportExcelVO>> batches, Long organId,
boolean isUpdateSupport, List<PartsDO> insertList,
List<PartsDO> updateList) {
// 查询过程中是否异常了
AtomicBoolean hasError = new AtomicBoolean(false);
// 计数器
CountDownLatch latch = new CountDownLatch(batches.size());
// 主线程的data_code
String peek = DynamicDataSourceContextHolder.peek();
for (List<PartsImportExcelVO> batch : batches) {
threadPoolTaskExecutor.submit(() -> {
threadPoolTaskExecutor.submit(() -> {
try {
// 子线程同步主线程的多数据源
DynamicDataSourceContextHolder.push(peek);
List<String> goodsIds = batch.stream().map(PartsImportExcelVO::getGoodsId).collect(Collectors.toList());
Map<String, PartsDO> partsMap = partsMapper.selectList(
new LambdaQueryWrapper<PartsDO>().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<PartsDO> dataList, int batchSize, Consumer<List<PartsDO>> 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<PartsDO> subList = dataList.subList(i, end);
if (!subList.isEmpty()) {
processor.accept(subList);
}
}
}
/**
* 拆分导入数据为batchSize一组
*
* @param list
* @param batchSize
* @return
*/
private List<List<PartsImportExcelVO>> createBatches(List<PartsImportExcelVO> list, int batchSize) {
List<List<PartsImportExcelVO>> 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<PartsDO> getPartAll(String category) {
return partsMapper.getAllByOrganId(OrganContextHolder.getOrganId(), category);
@@ -81,7 +81,7 @@ public interface PlateManageService {
* @param isUpdateSupport 是否支持更新
* @return 导入结果
*/
PlateImportRespVO importPlateList(List<PlateImportExcelVO> importUsers, boolean isUpdateSupport , Long organId);
void importPlateList(List<PlateImportExcelVO> importUsers, boolean isUpdateSupport , Long organId);
/**
* 生产导入文件模板下载
@@ -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<PlateImportExcelVO> importPlates, boolean isUpdateSupport, Long organId) {
public void importPlateList(List<PlateImportExcelVO> 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<PlateGoodDO> insertList = new ArrayList<>();
// 批量更新集合
List<PlateGoodDO> 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<PlateGoodDO> insertList = new CopyOnWriteArrayList<>();
// 批量更新集合
List<PlateGoodDO> updateList = new CopyOnWriteArrayList<>();
// 每100条 goodsId in查询一次
List<List<PlateImportExcelVO>> 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<List<PlateImportExcelVO>> batches, Long organId,
boolean isUpdateSupport, List<PlateGoodDO> insertList,
List<PlateGoodDO> updateList) {
// 查询过程中是否异常了
AtomicBoolean hasError = new AtomicBoolean(false);
// 计数器
CountDownLatch latch = new CountDownLatch(batches.size());
// 主线程的data_code
String peek = DynamicDataSourceContextHolder.peek();
for (List<PlateImportExcelVO> batch : batches) {
threadPoolTaskExecutor.submit(() -> {
threadPoolTaskExecutor.submit(() -> {
try {
// 子线程同步主线程的多数据源
DynamicDataSourceContextHolder.push(peek);
List<String> goodsIds = batch.stream().map(PlateImportExcelVO::getGoodsId).collect(Collectors.toList());
Map<String, PlateGoodDO> plateGoodMap = plateGoodMapper.selectList(
new LambdaQueryWrapper<PlateGoodDO>().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<PlateGoodDO> dataList, int batchSize, Consumer<List<PlateGoodDO>> 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<PlateGoodDO> subList = dataList.subList(i, end);
if (!subList.isEmpty()) {
processor.accept(subList);
}
}
}
/**
* 拆分导入数据为batchSize一组
*
* @param list
* @param batchSize
* @return
*/
private List<List<PlateImportExcelVO>> createBatches(List<PlateImportExcelVO> list, int batchSize) {
List<List<PlateImportExcelVO>> 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