| | |
| | | package com.zy.core.plugin; |
| | | |
| | | import com.alibaba.fastjson.JSON; |
| | | import com.alibaba.fastjson.JSONObject; |
| | | import com.alibaba.fastjson.serializer.SerializerFeature; |
| | | import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper; |
| | | import com.core.common.Cools; |
| | |
| | | import com.zy.asrs.entity.*; |
| | | import com.zy.asrs.service.*; |
| | | import com.zy.common.entity.FindCrnNoResult; |
| | | import com.zy.common.model.StartupDto; |
| | | import com.zy.common.service.CommonService; |
| | | import com.zy.common.utils.RedisUtil; |
| | | import com.zy.core.News; |
| | |
| | | import com.zy.core.plugin.store.StoreInTaskContext; |
| | | import com.zy.core.plugin.store.StoreInTaskGenerationService; |
| | | import com.zy.core.plugin.store.StoreInTaskPolicy; |
| | | import com.zy.core.properties.SystemProperties; |
| | | import com.zy.core.task.MainProcessAsyncTaskScheduler; |
| | | import com.zy.core.thread.CrnThread; |
| | | import com.zy.core.thread.DualCrnThread; |
| | | import com.zy.core.thread.StationThread; |
| | | import com.zy.core.utils.CrnOperateProcessUtils; |
| | | import com.zy.core.utils.DualCrnOperateProcessUtils; |
| | | import com.zy.core.utils.StationOperateProcessUtils; |
| | | import com.zy.core.utils.WmsOperateUtils; |
| | | import com.zy.system.entity.Config; |
| | | import com.zy.system.service.ConfigService; |
| | | import lombok.extern.slf4j.Slf4j; |
| | |
| | | |
| | | import java.util.*; |
| | | import java.util.concurrent.ConcurrentHashMap; |
| | | import java.util.concurrent.ExecutorService; |
| | | import java.util.concurrent.Executors; |
| | | import java.util.concurrent.Future; |
| | | import java.util.concurrent.TimeUnit; |
| | | import java.util.concurrent.TimeoutException; |
| | | |
| | | @Slf4j |
| | | @Component |
| | | public class FakeProcess implements MainProcessPluginApi, StoreInTaskPolicy { |
| | | private static final String CRN_TASK_LANE = "fake-crn"; |
| | | private static final String STATION_TASK_LANE = "fake-station"; |
| | | private static final String DUAL_CRN_TASK_LANE = "fake-dual-crn"; |
| | | private static final String GENERATE_STORE_TASK_LANE_PREFIX = "fake-generate-store-"; |
| | | private static final String ASYNC_TASK_LANE_PREFIX = "fake-async-"; |
| | | private static final long MAIN_DISPATCH_INTERVAL_MS = 200L; |
| | | private static final long ASYNC_DISPATCH_INTERVAL_MS = 50L; |
| | | private static final long TASK_SLOW_LOG_THRESHOLD_MS = 1000L; |
| | | |
| | | private static final long METHOD_TIMEOUT_MS = 15000; // 15秒超时 |
| | | private static final ExecutorService timeoutExecutor = Executors.newCachedThreadPool(); |
| | | |
| | | private static Map<Integer, Long> stationStayTimeMap = new ConcurrentHashMap<>(); |
| | | private static final Map<Integer, Long> stationStayTimeMap = new ConcurrentHashMap<>(); |
| | | private static volatile String enableFake = "N"; |
| | | private static volatile String fakeRealTaskRequestWms = "N"; |
| | | private static volatile String fakeGenerateInTask = "Y"; |
| | | private static volatile String fakeGenerateOutTask = "Y"; |
| | | |
| | | private Thread asyncFakeRunThread = null; |
| | | |
| | | @Autowired |
| | | private WrkMastService wrkMastService; |
| | |
| | | @Autowired |
| | | private StationOperateProcessUtils stationOperateProcessUtils; |
| | | @Autowired |
| | | private WmsOperateUtils wmsOperateUtils; |
| | | @Autowired |
| | | private WrkAnalysisService wrkAnalysisService; |
| | | @Autowired |
| | | private DualCrnOperateProcessUtils dualCrnOperateProcessUtils; |
| | |
| | | private StoreInTaskGenerationService storeInTaskGenerationService; |
| | | @Autowired |
| | | private StationCommandDispatcher stationCommandDispatcher; |
| | | |
| | | /** |
| | | * 带超时保护执行方法 |
| | | * |
| | | * @param taskName 任务名称(用于日志) |
| | | * @param task 要执行的任务 |
| | | */ |
| | | private void executeWithTimeout(String taskName, Runnable task) { |
| | | Future<?> future = timeoutExecutor.submit(task); |
| | | try { |
| | | future.get(METHOD_TIMEOUT_MS, TimeUnit.MILLISECONDS); |
| | | } catch (TimeoutException e) { |
| | | // 使用 cancel(false) 不发送中断信号,避免 RedisCommandInterruptedException |
| | | // 任务会在后台继续执行直到完成,但主循环不会等待 |
| | | future.cancel(false); |
| | | News.error("[WCS Warning] 方法执行超时,主循环已跳过: {}, 超时时间: {}ms (任务仍在后台运行)", taskName, METHOD_TIMEOUT_MS); |
| | | } catch (Exception e) { |
| | | News.error("[WCS Error] 方法执行异常: {}, 异常: {}", taskName, e.getMessage()); |
| | | } |
| | | } |
| | | @Autowired |
| | | private MainProcessAsyncTaskScheduler mainProcessAsyncTaskScheduler; |
| | | |
| | | @Override |
| | | public void run() { |
| | | long startTime = System.currentTimeMillis(); |
| | | asyncFakeRun(); |
| | | refreshFakeConfig(); |
| | | |
| | | // 请求生成入库任务 |
| | | this.generateStoreWrkFile(); |
| | | // 仿真相关任务按各自 lane 串行提交,互不阻塞 |
| | | submitAsyncTask("checkInStationHasTask", ASYNC_DISPATCH_INTERVAL_MS, this::checkInStationHasTask); |
| | | submitAsyncTask("generateFakeInTask", ASYNC_DISPATCH_INTERVAL_MS, this::generateFakeInTask); |
| | | submitAsyncTask("generateFakeOutTask", ASYNC_DISPATCH_INTERVAL_MS, this::generateFakeOutTask); |
| | | submitAsyncTask("calcAllStationStayTime", ASYNC_DISPATCH_INTERVAL_MS, this::calcAllStationStayTime); |
| | | submitAsyncTask("checkOutStationStayTimeOut", ASYNC_DISPATCH_INTERVAL_MS, this::checkOutStationStayTimeOut); |
| | | submitAsyncTask("checkInStationCrnTake", ASYNC_DISPATCH_INTERVAL_MS, this::checkInStationCrnTake); |
| | | submitStationTask("checkStationRunBlock", ASYNC_DISPATCH_INTERVAL_MS, stationOperateProcessUtils::checkStationRunBlock); |
| | | |
| | | // 执行堆垛机任务 |
| | | crnOperateUtils.crnIoExecute(); |
| | | // 堆垛机任务执行完成-具备仿真能力 |
| | | this.crnIoExecuteFinish(); |
| | | // 执行输送站点入库任务 |
| | | stationOperateProcessUtils.stationInExecute(); |
| | | // 执行输送站点出库任务 |
| | | stationOperateProcessUtils.crnStationOutExecute(); |
| | | // 检测出库排序 |
| | | stationOperateProcessUtils.checkStationOutOrder(); |
| | | // 监控绕圈站点 |
| | | stationOperateProcessUtils.watchCircleStation(); |
| | | // 请求生成入库任务,按站点拆成独立串行通道,避免单站阻塞整轮扫描 |
| | | generateStoreWrkFile(); |
| | | |
| | | // 执行双工位堆垛机任务 |
| | | dualCrnOperateProcessUtils.dualCrnIoExecute(); |
| | | // 双工位堆垛机任务执行完成 |
| | | dualCrnOperateProcessUtils.dualCrnIoExecuteFinish(); |
| | | // 堆垛机、输送线、双工位堆垛机任务分别进入各自串行通道逐个执行 |
| | | submitCrnTask("crnIoExecute", MAIN_DISPATCH_INTERVAL_MS, crnOperateUtils::crnIoExecute); |
| | | submitCrnTask("crnIoExecuteFinish", MAIN_DISPATCH_INTERVAL_MS, this::crnIoExecuteFinish); |
| | | submitStationTask("stationInExecute", MAIN_DISPATCH_INTERVAL_MS, stationOperateProcessUtils::stationInExecute); |
| | | submitStationTask("crnStationOutExecute", MAIN_DISPATCH_INTERVAL_MS, stationOperateProcessUtils::crnStationOutExecute); |
| | | submitStationTask("checkStationOutOrder", MAIN_DISPATCH_INTERVAL_MS, stationOperateProcessUtils::checkStationOutOrder); |
| | | submitStationTask("watchCircleStation", MAIN_DISPATCH_INTERVAL_MS, stationOperateProcessUtils::watchCircleStation); |
| | | submitDualCrnTask("dualCrnIoExecute", MAIN_DISPATCH_INTERVAL_MS, dualCrnOperateProcessUtils::dualCrnIoExecute); |
| | | submitDualCrnTask("dualCrnIoExecuteFinish", MAIN_DISPATCH_INTERVAL_MS, dualCrnOperateProcessUtils::dualCrnIoExecuteFinish); |
| | | |
| | | News.info("[WCS Debug] 主线程Run执行完成,耗时:{}ms", System.currentTimeMillis() - startTime); |
| | | } |
| | | |
| | | public void asyncFakeRun() { |
| | | if (asyncFakeRunThread != null) { |
| | | return; |
| | | } |
| | | |
| | | asyncFakeRunThread = new Thread(() -> { |
| | | while (!Thread.currentThread().isInterrupted()) { |
| | | try { |
| | | Config enableFakeConfig = configService |
| | | .getOne(new QueryWrapper<Config>().eq("code", "enableFake")); |
| | | private void refreshFakeConfig() { |
| | | Config enableFakeConfig = configService.getOne(new QueryWrapper<Config>().eq("code", "enableFake")); |
| | | if (enableFakeConfig != null) { |
| | | enableFake = enableFakeConfig.getValue(); |
| | | } |
| | |
| | | if (fakeGenerateOutTaskConfig != null) { |
| | | fakeGenerateOutTask = fakeGenerateOutTaskConfig.getValue(); |
| | | } |
| | | } |
| | | |
| | | // 系统运行状态判断 |
| | | if (!SystemProperties.WCS_RUNNING_STATUS.get()) { |
| | | private void submitCrnTask(String taskName, long minIntervalMs, Runnable task) { |
| | | submitProcessTask(CRN_TASK_LANE, taskName, minIntervalMs, task); |
| | | } |
| | | |
| | | private void submitStationTask(String taskName, long minIntervalMs, Runnable task) { |
| | | submitProcessTask(STATION_TASK_LANE, taskName, minIntervalMs, task); |
| | | } |
| | | |
| | | private void submitDualCrnTask(String taskName, long minIntervalMs, Runnable task) { |
| | | submitProcessTask(DUAL_CRN_TASK_LANE, taskName, minIntervalMs, task); |
| | | } |
| | | |
| | | private void submitAsyncTask(String taskName, long minIntervalMs, Runnable task) { |
| | | submitProcessTask(ASYNC_TASK_LANE_PREFIX + taskName, taskName, minIntervalMs, task); |
| | | } |
| | | |
| | | private void submitGenerateStoreTasks() { |
| | | List<BasDevp> basDevps = basDevpService.list(new QueryWrapper<>()); |
| | | for (BasDevp basDevp : basDevps) { |
| | | List<StationObjModel> barcodeStations = getBarcodeStations(basDevp); |
| | | for (StationObjModel stationObjModel : barcodeStations) { |
| | | Integer stationId = stationObjModel == null ? null : stationObjModel.getStationId(); |
| | | if (stationId == null) { |
| | | continue; |
| | | } |
| | | |
| | | // 检测入库站是否有任务生成,并仿真生成模拟入库站点数据 |
| | | checkInStationHasTask(); |
| | | // 生成仿真模拟入库任务 |
| | | generateFakeInTask(); |
| | | // 生成仿真模拟出库任务 |
| | | generateFakeOutTask(); |
| | | // 计算所有站点停留时间 |
| | | calcAllStationStayTime(); |
| | | // 检测出库站点停留是否超时 |
| | | checkOutStationStayTimeOut(); |
| | | // 检测入库站点堆垛机是否取走货物 |
| | | checkInStationCrnTake(); |
| | | |
| | | // 检测输送站点是否运行堵塞 |
| | | stationOperateProcessUtils.checkStationRunBlock(); |
| | | |
| | | // 间隔 |
| | | Thread.sleep(50); |
| | | } catch (InterruptedException ie) { |
| | | Thread.currentThread().interrupt(); |
| | | break; |
| | | } catch (Exception e) { |
| | | e.printStackTrace(); |
| | | submitGenerateStoreTask(stationId, MAIN_DISPATCH_INTERVAL_MS, |
| | | () -> storeInTaskGenerationService.generate(this, basDevp, stationObjModel)); |
| | | } |
| | | } |
| | | }); |
| | | asyncFakeRunThread.setName("asyncFakeRunProcess"); |
| | | asyncFakeRunThread.setDaemon(true); |
| | | asyncFakeRunThread.start(); |
| | | } |
| | | |
| | | private void submitGenerateStoreTask(Integer stationId, long minIntervalMs, Runnable task) { |
| | | submitProcessTask(GENERATE_STORE_TASK_LANE_PREFIX + stationId, "generateStoreWrkFile", minIntervalMs, task); |
| | | } |
| | | |
| | | private void submitProcessTask(String laneName, String taskName, long minIntervalMs, Runnable task) { |
| | | mainProcessAsyncTaskScheduler.submit( |
| | | laneName, |
| | | taskName, |
| | | minIntervalMs, |
| | | TASK_SLOW_LOG_THRESHOLD_MS, |
| | | task |
| | | ); |
| | | } |
| | | |
| | | // 检测入库站是否有任务生成,并仿真生成模拟入库站点数据 |
| | | private synchronized void checkInStationHasTask() { |
| | | private void checkInStationHasTask() { |
| | | if (!enableFake.equals("Y")) { |
| | | return; |
| | | } |
| | |
| | | } |
| | | |
| | | // 生成仿真模拟入库任务 |
| | | private synchronized void generateFakeInTask() { |
| | | private void generateFakeInTask() { |
| | | if (!enableFake.equals("Y")) { |
| | | return; |
| | | } |
| | |
| | | } |
| | | |
| | | // 生成仿真模拟出库任务 |
| | | private synchronized void generateFakeOutTask() { |
| | | private void generateFakeOutTask() { |
| | | try { |
| | | if (!enableFake.equals("Y")) { |
| | | return; |
| | |
| | | } |
| | | } |
| | | } catch (Exception e) { |
| | | e.printStackTrace(); |
| | | log.error("generateFakeOutTask execute error", e); |
| | | } |
| | | } |
| | | |
| | |
| | | * 请求生成入库任务 |
| | | * 入库站,根据条码扫描生成入库工作档 |
| | | */ |
| | | public synchronized void generateStoreWrkFile() { |
| | | storeInTaskGenerationService.generate(this); |
| | | public void generateStoreWrkFile() { |
| | | submitGenerateStoreTasks(); |
| | | } |
| | | |
| | | @Override |
| | |
| | | } |
| | | |
| | | // 计算所有站点停留时间 |
| | | public synchronized void calcAllStationStayTime() { |
| | | public void calcAllStationStayTime() { |
| | | List<BasDevp> basDevps = basDevpService.list(new QueryWrapper<>()); |
| | | for (BasDevp basDevp : basDevps) { |
| | | StationThread stationThread = (StationThread) SlaveConnection.get(SlaveType.Devp, basDevp.getDevpNo()); |
| | |
| | | } |
| | | |
| | | // 检测出库站点停留是否超时 |
| | | public synchronized void checkOutStationStayTimeOut() { |
| | | public void checkOutStationStayTimeOut() { |
| | | List<BasDevp> basDevps = basDevpService.list(new QueryWrapper<>()); |
| | | for (BasDevp basDevp : basDevps) { |
| | | List<StationObjModel> outStationList = basDevp.getOutStationList$(); |
| | |
| | | } |
| | | |
| | | // 检测入库站点堆垛机是否取走货物 |
| | | public synchronized void checkInStationCrnTake() { |
| | | public void checkInStationCrnTake() { |
| | | List<BasCrnp> basCrnps = basCrnpService.list(new QueryWrapper<>()); |
| | | for (BasCrnp basCrnp : basCrnps) { |
| | | List<StationObjModel> inStationList = basCrnp.getInStationList$(); |
| | |
| | | } |
| | | } |
| | | |
| | | private synchronized void checkInStationListCrnTake(List<StationObjModel> inStationList) { |
| | | private void checkInStationListCrnTake(List<StationObjModel> inStationList) { |
| | | for (StationObjModel stationObjModel : inStationList) { |
| | | Object lock = redisUtil |
| | | .get(RedisKeyType.CHECK_IN_STATION_STAY_TIME_OUT_LIMIT.key + stationObjModel.getStationId()); |
| | |
| | | } |
| | | |
| | | // 堆垛机任务执行完成 |
| | | public synchronized void crnIoExecuteFinish() { |
| | | public void crnIoExecuteFinish() { |
| | | List<BasCrnp> basCrnps = basCrnpService.list(new QueryWrapper<>()); |
| | | for (BasCrnp basCrnp : basCrnps) { |
| | | CrnThread crnThread = (CrnThread) SlaveConnection.get(SlaveType.Crn, basCrnp.getCrnNo()); |