From 459540819d35d61445a52d727478c362721a8036 Mon Sep 17 00:00:00 2001 From: van Date: Tue, 29 Sep 2026 14:55:53 +0800 Subject: [PATCH] 1 --- .../jarvis/task/ServerResourceAlertTask.java | 288 +++++++++++++++++ .../src/main/resources/application-dev.yml | 16 + .../src/main/resources/application-prod.yml | 10 + .../config/JarvisGoofishProperties.java | 15 + .../service/goofish/GoofishOrderPipeline.java | 302 +++++++++++++++--- 5 files changed, 585 insertions(+), 46 deletions(-) create mode 100644 ruoyi-admin/src/main/java/com/ruoyi/jarvis/task/ServerResourceAlertTask.java diff --git a/ruoyi-admin/src/main/java/com/ruoyi/jarvis/task/ServerResourceAlertTask.java b/ruoyi-admin/src/main/java/com/ruoyi/jarvis/task/ServerResourceAlertTask.java new file mode 100644 index 0000000..53aab7c --- /dev/null +++ b/ruoyi-admin/src/main/java/com/ruoyi/jarvis/task/ServerResourceAlertTask.java @@ -0,0 +1,288 @@ +package com.ruoyi.jarvis.task; + +import com.alibaba.fastjson2.JSONObject; +import com.ruoyi.common.utils.Arith; +import com.ruoyi.framework.web.domain.Server; +import com.ruoyi.framework.web.domain.server.SysFile; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Component; +import org.springframework.util.StringUtils; + +import java.io.InputStream; +import java.io.OutputStream; +import java.net.HttpURLConnection; +import java.net.URL; +import java.nio.charset.StandardCharsets; +import java.time.LocalDateTime; +import java.time.format.DateTimeFormatter; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.HashSet; +import java.util.List; +import java.util.Locale; +import java.util.Set; +import java.util.concurrent.atomic.AtomicBoolean; + +/** + * 按服务监控同一套采样({@link Server#copyTo()})定时看 CPU 与磁盘。 + * 任一超过阈值时,经 wxSend {@code /wx/send}、消息类型 QL(企微应用 1000002)连发告警。 + */ +@Component +public class ServerResourceAlertTask { + + private static final Logger log = LoggerFactory.getLogger(ServerResourceAlertTask.class); + + private static final DateTimeFormatter TIME_FMT = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"); + + /** 内存盘、容器层等不代表物理硬盘,避免误报 */ + private static final Set PSEUDO_FS = new HashSet<>(Arrays.asList( + "tmpfs", "devtmpfs", "overlay", "squashfs", "ramfs", "sysfs", "proc", + "cgroup", "cgroup2", "devfs", "autofs", "nsfs", "tracefs", "debugfs", + "securityfs", "pstore", "bpf", "configfs", "mqueue", "hugetlbfs", + "fusectl", "binfmt_misc")); + + @Value("${jarvis.server.resource-alert.enabled:true}") + private boolean enabled; + + /** 与服务监控页磁盘标红线一致,超过即告警 */ + @Value("${jarvis.server.resource-alert.cpu-threshold:80}") + private double cpuThreshold; + + @Value("${jarvis.server.resource-alert.disk-threshold:80}") + private double diskThreshold; + + @Value("${jarvis.server.resource-alert.repeat-count:3}") + private int repeatCount; + + /** 持续超标时两次连发之间的最短间隔,避免每轮刷新都刷屏 */ + @Value("${jarvis.server.resource-alert.cooldown-minutes:30}") + private int cooldownMinutes; + + @Value("${jarvis.server.resource-alert.touser:LinPingFan}") + private String touser; + + /** wxSend 消息类型名,QL 对应企微应用 1000002 */ + @Value("${jarvis.server.resource-alert.message-type:QL}") + private String messageType; + + @Value("${jarvis.wecom.wxsend-base-url:}") + private String wxsendBaseUrl; + + @Value("${jarvis.wecom.wxsend-van-token:}") + private String vanToken; + + private final AtomicBoolean running = new AtomicBoolean(false); + + private volatile long lastAlertAtMs; + + @Scheduled(cron = "${jarvis.server.resource-alert.cron:0 */5 * * * ?}") + public void refreshAndAlert() { + if (!enabled) { + return; + } + if (!running.compareAndSet(false, true)) { + log.info("服务器资源监控上一轮尚未结束,跳过本次"); + return; + } + try { + sampleAndMaybeAlert(); + } catch (Exception e) { + log.error("服务器资源监控定时任务异常", e); + } finally { + running.set(false); + } + } + + private void sampleAndMaybeAlert() throws Exception { + Server server = new Server(); + server.copyTo(); + + double cpuFree = server.getCpu().getFree(); + double cpuUsage = clampPercent(Arith.sub(100D, cpuFree)); + boolean cpuHot = cpuUsage > cpuThreshold; + + List hotDisks = new ArrayList<>(); + double maxDisk = -1D; + String maxDiskDir = ""; + List files = server.getSysFiles(); + if (files != null) { + for (SysFile file : files) { + if (file == null || !isPhysicalDisk(file)) { + continue; + } + double usage = file.getUsage(); + if (usage > maxDisk) { + maxDisk = usage; + maxDiskDir = file.getDirName() == null ? "" : file.getDirName(); + } + if (usage > diskThreshold) { + hotDisks.add(formatDisk(file)); + } + } + } + boolean diskHot = !hotDisks.isEmpty(); + + String diskSummary = maxDisk < 0 ? "无物理盘" : String.format(Locale.ROOT, "%.2f%%(%s)", maxDisk, maxDiskDir); + log.info("服务器资源监控 CPU={}%(阈值{}%) 磁盘最高={}(阈值{}%)", + formatPercent(cpuUsage), formatPercent(cpuThreshold), diskSummary, formatPercent(diskThreshold)); + + if (!cpuHot && !diskHot) { + return; + } + long now = System.currentTimeMillis(); + long cooldownMs = Math.max(1, cooldownMinutes) * 60_000L; + if (lastAlertAtMs > 0 && now - lastAlertAtMs < cooldownMs) { + log.info("资源仍超标,处于告警冷却({} 分钟),本次不重复推送", cooldownMinutes); + return; + } + + String kind = cpuHot && diskHot ? "CPU与磁盘" : (cpuHot ? "CPU" : "磁盘"); + String body = buildBody(server, cpuUsage, cpuHot, hotDisks, diskHot); + int times = Math.max(1, repeatCount); + int sent = 0; + for (int i = 1; i <= times; i++) { + String title = "【资源告警 " + i + "/" + times + "】" + kind + "过高"; + if (postQl(title, body)) { + sent++; + } + if (i < times) { + Thread.sleep(400L); + } + } + if (sent > 0) { + lastAlertAtMs = System.currentTimeMillis(); + log.warn("服务器资源告警已推送 {}/{} 条,渠道 messageType={}(企微应用 1000002)", sent, times, messageType); + } else { + log.error("服务器资源已超标,但 {} 条告警均未发出,下轮将重试", times); + } + } + + private String buildBody(Server server, double cpuUsage, boolean cpuHot, List hotDisks, boolean diskHot) { + String name = server.getSys() == null ? "" : nullToEmpty(server.getSys().getComputerName()); + String ip = server.getSys() == null ? "" : nullToEmpty(server.getSys().getComputerIp()); + StringBuilder sb = new StringBuilder(); + sb.append("主机:").append(name); + if (StringUtils.hasText(ip)) { + sb.append("(").append(ip).append(")"); + } + sb.append("\n时间:").append(TIME_FMT.format(LocalDateTime.now())); + sb.append("\nCPU:").append(formatPercent(cpuUsage)).append("%"); + sb.append(cpuHot ? ",已超过阈值 " : ",未超过阈值 "); + sb.append(formatPercent(cpuThreshold)).append("%"); + if (diskHot) { + sb.append("\n磁盘超标:"); + for (String line : hotDisks) { + sb.append("\n- ").append(line); + } + } else { + sb.append("\n磁盘:未超过阈值 ").append(formatPercent(diskThreshold)).append("%"); + } + return sb.toString(); + } + + private boolean postQl(String title, String text) { + if (!StringUtils.hasText(wxsendBaseUrl)) { + log.error("资源告警未发送:未配置 jarvis.wecom.wxsend-base-url"); + return false; + } + if (!StringUtils.hasText(vanToken)) { + log.error("资源告警未发送:未配置 jarvis.wecom.wxsend-van-token"); + return false; + } + String base = wxsendBaseUrl.trim(); + if (base.endsWith("/")) { + base = base.substring(0, base.length() - 1); + } + String url = base + "/wx/send"; + JSONObject body = new JSONObject(); + body.put("title", title); + body.put("text", text); + body.put("touser", touser == null ? "" : touser.trim()); + body.put("vanToken", vanToken.trim()); + body.put("messageType", StringUtils.hasText(messageType) ? messageType.trim() : "QL"); + + HttpURLConnection conn = null; + try { + conn = (HttpURLConnection) new URL(url).openConnection(); + conn.setRequestMethod("POST"); + conn.setConnectTimeout(10000); + conn.setReadTimeout(20000); + conn.setDoOutput(true); + conn.setRequestProperty("Content-Type", "application/json;charset=UTF-8"); + byte[] bytes = body.toJSONString().getBytes(StandardCharsets.UTF_8); + try (OutputStream os = conn.getOutputStream()) { + os.write(bytes); + } + int code = conn.getResponseCode(); + InputStream is = code >= 200 && code < 300 ? conn.getInputStream() : conn.getErrorStream(); + String resp = readAll(is); + boolean ok = code >= 200 && code < 300 && resp != null && resp.contains("成功"); + if (!ok) { + log.error("资源告警推送失败 http={} url={} resp={}", code, url, resp); + return false; + } + log.info("资源告警已发出 title={} resp={}", title, resp); + return true; + } catch (Exception e) { + log.error("资源告警推送异常 url={} err={}", url, e.toString(), e); + return false; + } finally { + if (conn != null) { + conn.disconnect(); + } + } + } + + private static boolean isPhysicalDisk(SysFile file) { + String total = file.getTotal(); + if (!StringUtils.hasText(total) || !total.toUpperCase(Locale.ROOT).contains("GB")) { + return false; + } + String fsType = file.getSysTypeName(); + if (!StringUtils.hasText(fsType)) { + return true; + } + return !PSEUDO_FS.contains(fsType.trim().toLowerCase(Locale.ROOT)); + } + + private static String formatDisk(SysFile file) { + return nullToEmpty(file.getDirName()) + + " 已用 " + formatPercent(file.getUsage()) + "%" + + "(可用 " + nullToEmpty(file.getFree()) + + " / 共 " + nullToEmpty(file.getTotal()) + ")"; + } + + private static double clampPercent(double value) { + if (value < 0D) { + return 0D; + } + if (value > 100D) { + return 100D; + } + return Arith.round(value, 2); + } + + private static String formatPercent(double value) { + return String.format(Locale.ROOT, "%.2f", value); + } + + private static String nullToEmpty(String value) { + return value == null ? "" : value; + } + + private static String readAll(InputStream is) throws java.io.IOException { + if (is == null) { + return ""; + } + byte[] buf = new byte[4096]; + StringBuilder sb = new StringBuilder(); + int n; + while ((n = is.read(buf)) >= 0) { + sb.append(new String(buf, 0, n, StandardCharsets.UTF_8)); + } + return sb.toString(); + } +} diff --git a/ruoyi-admin/src/main/resources/application-dev.yml b/ruoyi-admin/src/main/resources/application-dev.yml index f3db124..41df13e 100644 --- a/ruoyi-admin/src/main/resources/application-dev.yml +++ b/ruoyi-admin/src/main/resources/application-dev.yml @@ -219,6 +219,16 @@ jarvis: order-delay-ms: 250 # 0=不限制;例如 40 可控制单轮最长耗时(余下下轮再扫) max-orders-per-round: 0 + # 服务监控定时刷新:CPU 或物理盘使用率超过阈值时,经企微应用 1000002(wxSend messageType=QL)连发告警 + resource-alert: + enabled: true + cron: "0 */5 * * * ?" + cpu-threshold: 80 + disk-threshold: 80 + repeat-count: 3 + cooldown-minutes: 30 + touser: "LinPingFan" + message-type: QL # 获取评论接口服务地址(后端转发,避免前端跨域) fetch-comments: base-url: http://192.168.8.60:5008 @@ -322,5 +332,11 @@ jarvis: # true=仅拉 auto-ship-order-statuses(省调用,其它状态依赖推送);false=时间窗内全状态(推荐,与本地对齐) pull-list-only-auto-ship-statuses: false auto-ship-batch-size: 20 + # 两次发货接口最小间隔(毫秒),避免连续/并发发货被平台返回「操作频繁」 + ship-min-interval-ms: 2000 + # 「操作频繁」、超时等可恢复错误的额外重试次数(不含首次) + ship-retry-times: 3 + # 重试基础等待(毫秒),第 n 次等待 n × 该值 + ship-retry-delay-ms: 2000 diff --git a/ruoyi-admin/src/main/resources/application-prod.yml b/ruoyi-admin/src/main/resources/application-prod.yml index 753e184..a122b33 100644 --- a/ruoyi-admin/src/main/resources/application-prod.yml +++ b/ruoyi-admin/src/main/resources/application-prod.yml @@ -216,6 +216,16 @@ jarvis: cron: "0 */5 * * * ?" order-delay-ms: 250 max-orders-per-round: 0 + # 服务监控定时刷新:CPU 或物理盘使用率超过阈值时,经企微应用 1000002(wxSend messageType=QL)连发告警 + resource-alert: + enabled: true + cron: "0 */5 * * * ?" + cpu-threshold: 80 + disk-threshold: 80 + repeat-count: 3 + cooldown-minutes: 30 + touser: "LinPingFan" + message-type: QL # 获取评论接口服务地址(后端转发) fetch-comments: base-url: http://192.168.8.60:5008 diff --git a/ruoyi-system/src/main/java/com/ruoyi/jarvis/config/JarvisGoofishProperties.java b/ruoyi-system/src/main/java/com/ruoyi/jarvis/config/JarvisGoofishProperties.java index 6b4a1b4..456ca38 100644 --- a/ruoyi-system/src/main/java/com/ruoyi/jarvis/config/JarvisGoofishProperties.java +++ b/ruoyi-system/src/main/java/com/ruoyi/jarvis/config/JarvisGoofishProperties.java @@ -71,4 +71,19 @@ public class JarvisGoofishProperties { /** 与 defaultShipExpressCode 配套的展示名称 */ private String defaultShipExpressName = "日日顺"; + + /** + * 两次发货接口调用的最小间隔(毫秒)。开放平台对连续发货会返回「操作频繁」。 + */ + private long shipMinIntervalMs = 2000L; + + /** + * 发货遇到「操作频繁」、超时等可恢复错误时的额外重试次数(不含首次)。 + */ + private int shipRetryTimes = 3; + + /** + * 发货重试基础等待(毫秒)。第 n 次重试等待 n × 该值。 + */ + private long shipRetryDelayMs = 2000L; } diff --git a/ruoyi-system/src/main/java/com/ruoyi/jarvis/service/goofish/GoofishOrderPipeline.java b/ruoyi-system/src/main/java/com/ruoyi/jarvis/service/goofish/GoofishOrderPipeline.java index 8e1be54..0c16f4f 100644 --- a/ruoyi-system/src/main/java/com/ruoyi/jarvis/service/goofish/GoofishOrderPipeline.java +++ b/ruoyi-system/src/main/java/com/ruoyi/jarvis/service/goofish/GoofishOrderPipeline.java @@ -29,6 +29,7 @@ import java.util.ArrayList; import java.util.Date; import java.util.List; import java.util.Objects; +import java.util.concurrent.TimeUnit; /** * 闲管家:推送/拉单后的落库、详情、关联京东单、同步运单、发货 */ @@ -39,6 +40,15 @@ public class GoofishOrderPipeline { private static final String REDIS_WAYBILL_KEY_PREFIX = "logistics:waybill:order:"; /** F-闲鱼等京东单:运单可按第三方单号=闲鱼单号写入 Redis,见 LogisticsServiceImpl 镜像写入 */ public static final String REDIS_WAYBILL_GOOFISH_ORDER_PREFIX = "logistics:waybill:goofish:"; + /** 同一闲鱼单发货互斥,避免定时扫描与京东运单回调同时打出发货 */ + private static final String REDIS_SHIP_ORDER_LOCK_PREFIX = "goofish:ship:lock:"; + /** 全局发货槽:同一时刻只允许一笔发货请求在途,结束后再保留最小间隔 */ + private static final String REDIS_SHIP_GLOBAL_LOCK = "goofish:ship:global"; + /** 发货扫描互斥,避免拉单后扫描与定时扫描叠在一起 */ + private static final String REDIS_SHIP_BATCH_LOCK = "goofish:ship:batch"; + /** 可恢复失败后的短冷却,避免同一轮扫描立刻再打一次 */ + private static final String REDIS_SHIP_COOLDOWN_PREFIX = "goofish:ship:cooldown:"; + private static final long SHIP_FAIL_COOLDOWN_SECONDS = 30L; @Resource private ErpGoofishOrderMapper erpGoofishOrderMapper; @@ -565,34 +575,7 @@ public class GoofishOrderPipeline { ship.setWaybillNo(waybill.trim()); ship.setExpressCode(expressCode); ship.setExpressName(expressName); - String resp = ship.getResponseBody(); - JSONObject r = JSON.parseObject(resp); - if (r != null && r.getIntValue("code") == 0) { - ErpGoofishOrder ok = new ErpGoofishOrder(); - ok.setId(row.getId()); - ok.setShipStatus(1); - ok.setShipTime(DateUtils.getNowDate()); - ok.setShipExpressCode(expressCode); - ok.setShipError(null); - ok.setUpdateTime(DateUtils.getNowDate()); - erpGoofishOrderMapper.update(ok); - row.setShipStatus(1); - if (goofishOrderChangeLogger != null) { - goofishOrderChangeLogger.append(row.getId(), row.getAppKey(), row.getOrderNo(), - GoofishOrderChangeLogger.TYPE_SHIP, "AUTO_SHIP", - "发货成功,运单 " + waybill.trim() - + ",订单状态 " + GoofishStatusLabels.orderStatusHuman(row.getOrderStatus()) - + ",退款状态 " + GoofishStatusLabels.refundStatusHuman(row.getRefundStatus())); - } - } else { - String msg = r != null ? r.getString("msg") : "unknown"; - patchShipError(row, msg != null ? msg : "发货接口返回失败"); - if (goofishOrderChangeLogger != null) { - goofishOrderChangeLogger.append(row.getId(), row.getAppKey(), row.getOrderNo(), - GoofishOrderChangeLogger.TYPE_SHIP, "AUTO_SHIP", - "发货失败:" + (row.getShipError() != null ? row.getShipError() : msg)); - } - } + submitShip(row, ship, expressCode, waybill.trim(), manualRetry); } catch (Exception ex) { patchShipError(row, ex.getMessage()); if (goofishOrderChangeLogger != null) { @@ -604,6 +587,220 @@ public class GoofishOrderPipeline { } } + /** + * 同一订单只允许一路发货;全局再串行调用开放平台,遇到「操作频繁」等可恢复错误按间隔重试。 + */ + private void submitShip(ErpGoofishOrder row, OrderShipRequest ship, String expressCode, String waybill, + boolean manualRetry) { + if (!manualRetry && inShipCooldown(row.getId())) { + log.info("闲管家发货冷却中,跳过本轮 orderNo={}", row.getOrderNo()); + return; + } + String orderLock = REDIS_SHIP_ORDER_LOCK_PREFIX + row.getId(); + Boolean got = stringRedisTemplate.opsForValue().setIfAbsent(orderLock, "1", 180, TimeUnit.SECONDS); + if (!Boolean.TRUE.equals(got)) { + log.info("闲管家发货跳过:该单正在发货 orderNo={}", row.getOrderNo()); + return; + } + try { + ErpGoofishOrder fresh = erpGoofishOrderMapper.selectById(row.getId()); + if (fresh != null && fresh.getShipStatus() != null && fresh.getShipStatus() == 1) { + row.setShipStatus(1); + return; + } + int maxAttempts = Math.max(1, goofishProperties.getShipRetryTimes() + 1); + String lastMsg = null; + for (int attempt = 1; attempt <= maxAttempts; attempt++) { + String resp; + try { + resp = callShipSerialized(ship); + } catch (Exception ex) { + lastMsg = ex.getMessage() == null ? ex.toString() : ex.getMessage(); + log.warn("闲管家发货调用异常 orderNo={} attempt={}/{} {}", row.getOrderNo(), attempt, maxAttempts, lastMsg); + if (attempt < maxAttempts && sleepBeforeShipRetry(attempt)) { + continue; + } + patchShipError(row, lastMsg); + appendShipLog(row, "发货异常:" + (row.getShipError() != null ? row.getShipError() : lastMsg)); + armShipCooldown(row.getId()); + return; + } + JSONObject r; + try { + r = StringUtils.isEmpty(resp) ? null : JSON.parseObject(resp); + } catch (Exception parseEx) { + r = null; + } + if (r != null && r.getIntValue("code") == 0) { + markShipSuccess(row, expressCode, waybill); + return; + } + lastMsg = shipFailureMessage(r, resp); + if (isAlreadyShippedMessage(lastMsg)) { + log.info("闲管家发货接口表示已发货,按成功处理 orderNo={} msg={}", row.getOrderNo(), lastMsg); + markShipSuccess(row, expressCode, waybill); + return; + } + if (attempt < maxAttempts && isRetryableShipFailure(r, resp, lastMsg)) { + log.warn("闲管家发货将重试 orderNo={} attempt={}/{} msg={}", row.getOrderNo(), attempt, maxAttempts, lastMsg); + if (sleepBeforeShipRetry(attempt)) { + continue; + } + } + patchShipError(row, lastMsg); + appendShipLog(row, "发货失败:" + (row.getShipError() != null ? row.getShipError() : lastMsg)); + if (isRetryableShipFailure(r, resp, lastMsg)) { + armShipCooldown(row.getId()); + } + return; + } + } finally { + try { + stringRedisTemplate.delete(orderLock); + } catch (Exception e) { + log.debug("释放闲管家发货锁失败 orderNo={} {}", row.getOrderNo(), e.getMessage()); + } + } + } + + /** + * 全局只放行一笔发货请求。调用结束后把锁缩短为最小间隔,下一笔必须等这段时间过去。 + */ + private String callShipSerialized(OrderShipRequest ship) { + long intervalMs = goofishProperties.getShipMinIntervalMs(); + long deadline = System.currentTimeMillis() + 90_000L; + boolean locked = false; + while (System.currentTimeMillis() < deadline) { + Boolean ok = stringRedisTemplate.opsForValue().setIfAbsent(REDIS_SHIP_GLOBAL_LOCK, "1", 45, TimeUnit.SECONDS); + if (Boolean.TRUE.equals(ok)) { + locked = true; + break; + } + if (!sleepQuiet(250L)) { + throw new IllegalStateException("发货排队被中断"); + } + } + if (!locked) { + throw new IllegalStateException("发货排队超时,避免并发触发操作频繁"); + } + try { + return ship.getResponseBody(); + } finally { + try { + if (intervalMs <= 0) { + stringRedisTemplate.delete(REDIS_SHIP_GLOBAL_LOCK); + } else { + long seconds = Math.max(1L, (intervalMs + 999L) / 1000L); + stringRedisTemplate.expire(REDIS_SHIP_GLOBAL_LOCK, seconds, TimeUnit.SECONDS); + } + } catch (Exception e) { + try { + stringRedisTemplate.delete(REDIS_SHIP_GLOBAL_LOCK); + } catch (Exception ignore) { + log.debug("释放闲管家发货全局锁失败 {}", e.getMessage()); + } + } + } + } + + private void markShipSuccess(ErpGoofishOrder row, String expressCode, String waybill) { + ErpGoofishOrder ok = new ErpGoofishOrder(); + ok.setId(row.getId()); + ok.setShipStatus(1); + ok.setShipTime(DateUtils.getNowDate()); + ok.setShipExpressCode(expressCode); + ok.setShipError(""); + ok.setUpdateTime(DateUtils.getNowDate()); + erpGoofishOrderMapper.update(ok); + row.setShipStatus(1); + row.setShipError(null); + appendShipLog(row, "发货成功,运单 " + waybill + + ",订单状态 " + GoofishStatusLabels.orderStatusHuman(row.getOrderStatus()) + + ",退款状态 " + GoofishStatusLabels.refundStatusHuman(row.getRefundStatus())); + } + + private void appendShipLog(ErpGoofishOrder row, String message) { + if (goofishOrderChangeLogger == null || row == null) { + return; + } + goofishOrderChangeLogger.append(row.getId(), row.getAppKey(), row.getOrderNo(), + GoofishOrderChangeLogger.TYPE_SHIP, "AUTO_SHIP", message); + } + + private static String shipFailureMessage(JSONObject r, String resp) { + if (r != null && StringUtils.isNotEmpty(r.getString("msg"))) { + return r.getString("msg"); + } + if (StringUtils.isEmpty(resp)) { + return "发货接口无响应"; + } + return "发货接口返回失败"; + } + + private static boolean isAlreadyShippedMessage(String msg) { + return StringUtils.isNotEmpty(msg) && (msg.contains("已发货") || msg.contains("重复发货")); + } + + private static boolean isRetryableShipFailure(JSONObject r, String resp, String msg) { + if (r == null || StringUtils.isEmpty(resp)) { + return true; + } + return isRetryableShipMessage(msg); + } + + private static boolean isRetryableShipMessage(String msg) { + if (StringUtils.isEmpty(msg)) { + return true; + } + String m = msg.toLowerCase(); + return m.contains("频繁") || m.contains("稍后") || m.contains("繁忙") + || m.contains("限流") || m.contains("超时") + || m.contains("too many") || m.contains("rate limit") || m.contains("timeout"); + } + + private boolean sleepBeforeShipRetry(int attempt) { + long base = Math.max(200L, goofishProperties.getShipRetryDelayMs()); + long delay = base * Math.max(1, attempt); + log.info("闲管家发货重试等待 {} ms", delay); + return sleepQuiet(delay); + } + + private void armShipCooldown(Long orderId) { + if (orderId == null) { + return; + } + try { + stringRedisTemplate.opsForValue().set(REDIS_SHIP_COOLDOWN_PREFIX + orderId, "1", + SHIP_FAIL_COOLDOWN_SECONDS, TimeUnit.SECONDS); + } catch (Exception e) { + log.debug("写入闲管家发货冷却失败 id={} {}", orderId, e.getMessage()); + } + } + + private boolean inShipCooldown(Long orderId) { + if (orderId == null) { + return false; + } + try { + return Boolean.TRUE.equals(stringRedisTemplate.hasKey(REDIS_SHIP_COOLDOWN_PREFIX + orderId)); + } catch (Exception e) { + return false; + } + } + + private static boolean sleepQuiet(long millis) { + if (millis <= 0) { + return true; + } + try { + Thread.sleep(millis); + return true; + } catch (InterruptedException ie) { + Thread.currentThread().interrupt(); + return false; + } + } + private void patchShipError(ErpGoofishOrder row, String err) { ErpGoofishOrder p = new ErpGoofishOrder(); p.setId(row.getId()); @@ -1197,30 +1394,43 @@ public class GoofishOrderPipeline { } public int syncWaybillAndTryShipBatch(int limit) { - List statuses = resolveAutoShipOrderStatuses(); - List rows = erpGoofishOrderMapper.selectPendingShip(statuses, limit); - if (rows == null) { + Boolean got = stringRedisTemplate.opsForValue().setIfAbsent(REDIS_SHIP_BATCH_LOCK, "1", 10, TimeUnit.MINUTES); + if (!Boolean.TRUE.equals(got)) { + log.info("闲管家发货扫描已在进行,跳过本轮"); return 0; } - int n = 0; - for (ErpGoofishOrder row : rows) { - ErpGoofishOrder full = erpGoofishOrderMapper.selectById(row.getId()); - if (full == null) { - continue; + try { + List statuses = resolveAutoShipOrderStatuses(); + List rows = erpGoofishOrderMapper.selectPendingShip(statuses, limit); + if (rows == null) { + return 0; } - tryLinkJdOrder(full); - syncWaybillFromRedis(full); - if ((full.getRefundStatus() == null || full.getRefundStatus() == 0) - && statuses.contains(full.getOrderStatus()) - && StringUtils.isEmpty(full.getLocalWaybillNo()) - && StringUtils.isEmpty(full.getDetailWaybillNo())) { - refreshDetail(full); + int n = 0; + for (ErpGoofishOrder row : rows) { + ErpGoofishOrder full = erpGoofishOrderMapper.selectById(row.getId()); + if (full == null) { + continue; + } + tryLinkJdOrder(full); syncWaybillFromRedis(full); + if ((full.getRefundStatus() == null || full.getRefundStatus() == 0) + && statuses.contains(full.getOrderStatus()) + && StringUtils.isEmpty(full.getLocalWaybillNo()) + && StringUtils.isEmpty(full.getDetailWaybillNo())) { + refreshDetail(full); + syncWaybillFromRedis(full); + } + tryAutoShip(full); + n++; + } + return n; + } finally { + try { + stringRedisTemplate.delete(REDIS_SHIP_BATCH_LOCK); + } catch (Exception e) { + log.debug("释放闲管家发货扫描锁失败 {}", e.getMessage()); } - tryAutoShip(full); - n++; } - return n; } }