| src/main/java/com/zy/core/plugin/GslProcess.java | ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史 | |
| src/main/java/com/zy/core/plugin/NormalProcess.java | ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史 | |
| src/main/java/com/zy/core/plugin/XiaosongProcess.java | ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史 | |
| src/main/java/com/zy/core/task/MainProcessLane.java | ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史 | |
| src/main/java/com/zy/core/utils/StationOperateProcessUtils.java | ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史 |
src/main/java/com/zy/core/plugin/GslProcess.java
@@ -4,7 +4,6 @@ import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper; import com.core.common.Cools; import com.zy.asrs.entity.WrkLastno; import com.zy.asrs.utils.Utils; import com.zy.asrs.entity.BasDevp; import com.zy.asrs.service.BasDevpService; import com.zy.asrs.service.WrkLastnoService; @@ -60,8 +59,8 @@ @Override public void run() { //检测入库站是否有任务生成,并启动入库 checkInStationHasTask(); //检测入库站是否有任务生成,并按站点 lane 异步启动入库 stationOperateProcessUtils.submitStationEnableInTasks(DISPATCH_INTERVAL_MS); //请求生成入库任务,保留按站点 lane 串行提交 generateStoreWrkFile(); @@ -90,74 +89,6 @@ InTaskApplyRequest request = StoreInTaskPolicy.super.buildApplyRequest(context); request.getExtraParams().put("weight", context.getStationProtocol().getWeight()); return request; } //检测入库站是否有任务生成,并启动入库 private synchronized void checkInStationHasTask() { List<BasDevp> basDevps = basDevpService.list(new QueryWrapper<>()); for (BasDevp basDevp : basDevps) { StationThread stationThread = (StationThread) SlaveConnection.get(SlaveType.Devp, basDevp.getDevpNo()); if(stationThread == null){ continue; } Map<Integer, StationProtocol> stationMap = stationThread.getStatusMap(); List<StationObjModel> list = basDevp.getInStationList$(); for (StationObjModel entity : list) { Integer stationId = entity.getStationId(); if(!stationMap.containsKey(stationId)){ continue; } StationProtocol stationProtocol = stationMap.get(stationId); if (stationProtocol == null) { continue; } Object lock = redisUtil.get(RedisKeyType.GENERATE_ENABLE_IN_STATION_DATA_LIMIT.key + stationId); if(lock != null){ continue; } //满足自动、无物、工作号0,生成入库数据 if (stationProtocol.isAutoing() && stationProtocol.isLoading() && stationProtocol.getTaskNo() == 0 && stationProtocol.isEnableIn() ) { if (stationProtocol.getIoMode() != 1) { continue;//不属于入库模式 } //启动入库,删除条码站退回限制 Integer backStationId = entity.getBarcodeStation().getStationId(); String lockKey = RedisKeyType.GENERATE_STATION_BACK_LIMIT.key + backStationId; if (redisUtil.hasKey(lockKey)) { StationProtocol backStationProtocol = stationMap.get(backStationId); if (backStationProtocol == null) { continue; } if (backStationProtocol.isAutoing() && !backStationProtocol.isLoading() && stationProtocol.getTaskNo() == 0 ) { //条码站自动、无物、工作号0。删除条码站退回限制 redisUtil.del(lockKey); } } StationCommand command = stationThread.getCommand(StationCommandType.MOVE, commonService.getWorkNo(WrkIoType.ENABLE_IN.id), stationId, backStationId, 0); stationCommandDispatcher.dispatch(basDevp.getDevpNo(), command, "gsl-process", "enable-in"); if (entity.getBarcodeStation() != null && entity.getBarcodeStation().getStationId() != null) { Utils.precomputeInTaskEnableRow(entity.getBarcodeStation().getStationId()); } redisUtil.set(RedisKeyType.GENERATE_ENABLE_IN_STATION_DATA_LIMIT.key + stationId, "lock", 15); News.info("{}站点启动入库成功,数据包:{}", stationId, JSON.toJSONString(command)); } } } } private void generateStoreWrkFile() { src/main/java/com/zy/core/plugin/NormalProcess.java
@@ -1,22 +1,12 @@ package com.zy.core.plugin; import com.alibaba.fastjson.JSON; import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper; import com.core.common.Cools; import com.zy.asrs.utils.Utils; import com.zy.asrs.entity.BasDevp; import com.zy.asrs.service.BasDevpService; import com.zy.common.service.CommonService; import com.zy.common.utils.RedisUtil; import com.zy.core.News; import com.zy.core.cache.SlaveConnection; import com.zy.core.dispatch.StationCommandDispatcher; import com.zy.core.enums.RedisKeyType; import com.zy.core.enums.SlaveType; import com.zy.core.enums.StationCommandType; import com.zy.core.enums.WrkIoType; import com.zy.core.model.StationObjModel; import com.zy.core.model.command.StationCommand; import com.zy.core.model.protocol.StationProtocol; import com.zy.core.plugin.api.MainProcessPluginApi; import com.zy.core.plugin.store.StoreInTaskGenerationService; @@ -41,20 +31,14 @@ @Autowired private StationOperateProcessUtils stationOperateProcessUtils; @Autowired private CommonService commonService; @Autowired private BasDevpService basDevpService; @Autowired private RedisUtil redisUtil; @Autowired private StoreInTaskGenerationService storeInTaskGenerationService; @Autowired private StationCommandDispatcher stationCommandDispatcher; @Override public void run() { //检测入库站是否有任务生成,并启动入库 checkInStationHasTask(); stationOperateProcessUtils.submitStationEnableInTasks(DISPATCH_INTERVAL_MS); //请求生成入库任务 generateStoreWrkFile(); @@ -106,52 +90,6 @@ storeInTaskGenerationService.submitGenerateStoreTask(this, basDevp, stationObjModel, 0L, () -> storeInTaskGenerationService.generate(this, basDevp, stationObjModel)); } } } //检测入库站是否有任务生成,并启动入库 private synchronized void checkInStationHasTask() { List<BasDevp> basDevps = basDevpService.list(new QueryWrapper<>()); for (BasDevp basDevp : basDevps) { StationThread stationThread = (StationThread) SlaveConnection.get(SlaveType.Devp, basDevp.getDevpNo()); if(stationThread == null){ continue; } Map<Integer, StationProtocol> stationMap = stationThread.getStatusMap(); List<StationObjModel> list = basDevp.getInStationList$(); for (StationObjModel entity : list) { Integer stationId = entity.getStationId(); if(!stationMap.containsKey(stationId)){ continue; } StationProtocol stationProtocol = stationMap.get(stationId); if (stationProtocol == null) { continue; } Object lock = redisUtil.get(RedisKeyType.GENERATE_ENABLE_IN_STATION_DATA_LIMIT.key + stationId); if(lock != null){ continue; } //满足自动、无物、工作号0,生成入库数据 if (stationProtocol.isAutoing() && stationProtocol.isLoading() && stationProtocol.getTaskNo() == 0 && stationProtocol.isEnableIn() ) { StationCommand command = stationThread.getCommand(StationCommandType.MOVE, commonService.getWorkNo(WrkIoType.ENABLE_IN.id), stationId, entity.getBarcodeStation().getStationId(), 0); stationCommandDispatcher.dispatch(basDevp.getDevpNo(), command, "normal-process", "enable-in"); if (entity.getBarcodeStation() != null && entity.getBarcodeStation().getStationId() != null) { Utils.precomputeInTaskEnableRow(entity.getBarcodeStation().getStationId()); } redisUtil.set(RedisKeyType.GENERATE_ENABLE_IN_STATION_DATA_LIMIT.key + stationId, "lock", 15); News.info("{}站点启动入库成功,数据包:{}", stationId, JSON.toJSONString(command)); } } } } src/main/java/com/zy/core/plugin/XiaosongProcess.java
@@ -1,9 +1,7 @@ package com.zy.core.plugin; import com.alibaba.fastjson.JSON; import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper; import com.core.common.Cools; import com.zy.asrs.utils.Utils; import com.zy.asrs.entity.BasDevp; import com.zy.asrs.service.BasDevpService; import com.zy.common.service.CommonService; @@ -57,7 +55,7 @@ @Override public void run() { //检测入库站是否有任务生成,并启动入库 checkInStationHasTask(); stationOperateProcessUtils.submitStationEnableInTasks(DISPATCH_INTERVAL_MS); //请求生成入库任务 generateStoreWrkFile(); @@ -115,51 +113,6 @@ storeInTaskGenerationService.submitGenerateStoreTask(this, basDevp, stationObjModel, 0L, () -> storeInTaskGenerationService.generate(this, basDevp, stationObjModel)); } } } //检测入库站是否有任务生成,并启动入库 private synchronized void checkInStationHasTask() { List<BasDevp> basDevps = basDevpService.list(new QueryWrapper<>()); for (BasDevp basDevp : basDevps) { StationThread stationThread = (StationThread) SlaveConnection.get(SlaveType.Devp, basDevp.getDevpNo()); if(stationThread == null){ continue; } Map<Integer, StationProtocol> stationMap = stationThread.getStatusMap(); List<StationObjModel> list = basDevp.getInStationList$(); for (StationObjModel entity : list) { Integer stationId = entity.getStationId(); if(!stationMap.containsKey(stationId)){ continue; } StationProtocol stationProtocol = stationMap.get(stationId); if (stationProtocol == null) { continue; } Object lock = redisUtil.get(RedisKeyType.GENERATE_ENABLE_IN_STATION_DATA_LIMIT.key + stationId); if(lock != null){ continue; } if (stationProtocol.isAutoing() && stationProtocol.isLoading() && stationProtocol.getTaskNo() == 0 && stationProtocol.isEnableIn() ) { StationCommand command = stationThread.getCommand(StationCommandType.MOVE, commonService.getWorkNo(WrkIoType.ENABLE_IN.id), stationId, entity.getBarcodeStation().getStationId(), 0); stationCommandDispatcher.dispatch(basDevp.getDevpNo(), command, "xiaosong-process", "enable-in"); if (entity.getBarcodeStation() != null && entity.getBarcodeStation().getStationId() != null) { Utils.precomputeInTaskEnableRow(entity.getBarcodeStation().getStationId()); } redisUtil.set(RedisKeyType.GENERATE_ENABLE_IN_STATION_DATA_LIMIT.key + stationId, "lock", 15); News.info("{}站点启动入库成功,数据包:{}", stationId, JSON.toJSONString(command)); } } } } src/main/java/com/zy/core/task/MainProcessLane.java
@@ -7,6 +7,7 @@ DUAL_CRN_IO("dual-crn-io-"), DUAL_CRN_IO_FINISH("dual-crn-io-finish-"), STATION("station"), STATION_ENABLE_IN("station-enable-in-"), STATION_IN("station-in-"), STATION_OUT("station-out-"), DUAL_STATION_OUT("dual-station-out-"), src/main/java/com/zy/core/utils/StationOperateProcessUtils.java
@@ -1,16 +1,24 @@ package com.zy.core.utils; import com.alibaba.fastjson.JSON; import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper; import com.zy.asrs.entity.BasCrnp; import com.zy.asrs.entity.BasDevp; import com.zy.asrs.entity.WrkMast; import com.zy.asrs.utils.Utils; import com.zy.asrs.service.*; import com.zy.common.service.CommonService; import com.zy.common.utils.RedisUtil; import com.zy.core.News; import com.zy.core.cache.SlaveConnection; import com.zy.core.dispatch.StationCommandDispatcher; import com.zy.core.enums.RedisKeyType; import com.zy.core.enums.SlaveType; import com.zy.core.enums.StationCommandType; import com.zy.core.enums.WrkIoType; import com.zy.core.enums.WrkStsType; import com.zy.core.model.StationObjModel; import com.zy.core.model.command.StationCommand; import com.zy.core.model.protocol.StationProtocol; import com.zy.core.task.MainProcessLane; import com.zy.core.task.MainProcessTaskSubmitter; @@ -26,10 +34,11 @@ import java.util.Date; import java.util.List; import java.util.Map; import java.util.Objects; @Component public class StationOperateProcessUtils { private static final String STATION_COMMAND_SOURCE = "station-operate-process"; @Autowired private WrkMastService wrkMastService; @Autowired @@ -52,12 +61,106 @@ private StationOutboundDecisionSupport stationOutboundDecisionSupport; @Autowired private BasCrnpService basCrnpService; @Autowired private CommonService commonService; @Autowired private RedisUtil redisUtil; @Autowired private StationCommandDispatcher stationCommandDispatcher; public void submitStationEnableInTasks(long minIntervalMs) { submitStationEnableInTasks(MainProcessLane.STATION_ENABLE_IN, minIntervalMs); } public void submitStationEnableInTasks(MainProcessLane lane, long minIntervalMs) { List<BasDevp> basDevps = basDevpService.list(new QueryWrapper<>()); for (BasDevp basDevp : basDevps) { StationThread stationThread = (StationThread) SlaveConnection.get(SlaveType.Devp, basDevp.getDevpNo()); if (stationThread == null) { continue; } Map<Integer, StationProtocol> stationMap = stationThread.getStatusMap(); if (stationMap == null || stationMap.isEmpty()) { continue; } for (StationObjModel stationObjModel : basDevp.getInStationList$()) { Integer stationId = stationObjModel == null ? null : stationObjModel.getStationId(); if (stationId == null || !stationMap.containsKey(stationId)) { continue; } mainProcessTaskSubmitter.submitKeyedSerialTask( lane, stationId, "stationEnableInExecute", minIntervalMs, () -> stationEnableInExecute(basDevp, stationObjModel) ); } } } // 执行单个站点的入库任务下发 public void stationInExecute(BasDevp basDevp, StationObjModel stationObjModel) { stationRegularDispatchProcessor.stationInExecute(basDevp, stationObjModel); } // 执行单个站点的启动入库下发 public void stationEnableInExecute(BasDevp basDevp, StationObjModel stationObjModel) { if (basDevp == null || stationObjModel == null || stationObjModel.getStationId() == null) { return; } StationThread stationThread = (StationThread) SlaveConnection.get(SlaveType.Devp, basDevp.getDevpNo()); if (stationThread == null) { return; } Map<Integer, StationProtocol> stationMap = stationThread.getStatusMap(); if (stationMap == null || stationMap.isEmpty()) { return; } Integer stationId = stationObjModel.getStationId(); if (!stationMap.containsKey(stationId)) { return; } StationProtocol stationProtocol = stationMap.get(stationId); if (stationProtocol == null) { return; } Object lock = redisUtil.get(RedisKeyType.GENERATE_ENABLE_IN_STATION_DATA_LIMIT.key + stationId); if (lock != null) { return; } if (!stationProtocol.isAutoing() || !stationProtocol.isLoading() || stationProtocol.getTaskNo() != 0 || !stationProtocol.isEnableIn()) { return; } Integer backStationId = stationObjModel.getBarcodeStation() == null ? null : stationObjModel.getBarcodeStation().getStationId(); if (backStationId == null) { return; } StationCommand command = stationThread.getCommand( StationCommandType.MOVE, commonService.getWorkNo(WrkIoType.ENABLE_IN.id), stationId, backStationId, 0 ); stationCommandDispatcher.dispatch(basDevp.getDevpNo(), command, STATION_COMMAND_SOURCE, "enable-in"); Utils.precomputeInTaskEnableRow(backStationId); redisUtil.set(RedisKeyType.GENERATE_ENABLE_IN_STATION_DATA_LIMIT.key + stationId, "lock", 15); News.info("{}站点启动入库成功,数据包:{}", stationId, JSON.toJSONString(command)); } // 执行单个出库任务对应的输送站点下发 public void crnStationOutExecute(WrkMast wrkMast) { stationOutboundDispatchProcessor.crnStationOutExecute(wrkMast);