1
This commit is contained in:
@@ -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<String> 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<String> hotDisks = new ArrayList<>();
|
||||||
|
double maxDisk = -1D;
|
||||||
|
String maxDiskDir = "";
|
||||||
|
List<SysFile> 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<String> 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();
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -219,6 +219,16 @@ jarvis:
|
|||||||
order-delay-ms: 250
|
order-delay-ms: 250
|
||||||
# 0=不限制;例如 40 可控制单轮最长耗时(余下下轮再扫)
|
# 0=不限制;例如 40 可控制单轮最长耗时(余下下轮再扫)
|
||||||
max-orders-per-round: 0
|
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:
|
fetch-comments:
|
||||||
base-url: http://192.168.8.60:5008
|
base-url: http://192.168.8.60:5008
|
||||||
@@ -322,5 +332,11 @@ jarvis:
|
|||||||
# true=仅拉 auto-ship-order-statuses(省调用,其它状态依赖推送);false=时间窗内全状态(推荐,与本地对齐)
|
# true=仅拉 auto-ship-order-statuses(省调用,其它状态依赖推送);false=时间窗内全状态(推荐,与本地对齐)
|
||||||
pull-list-only-auto-ship-statuses: false
|
pull-list-only-auto-ship-statuses: false
|
||||||
auto-ship-batch-size: 20
|
auto-ship-batch-size: 20
|
||||||
|
# 两次发货接口最小间隔(毫秒),避免连续/并发发货被平台返回「操作频繁」
|
||||||
|
ship-min-interval-ms: 2000
|
||||||
|
# 「操作频繁」、超时等可恢复错误的额外重试次数(不含首次)
|
||||||
|
ship-retry-times: 3
|
||||||
|
# 重试基础等待(毫秒),第 n 次等待 n × 该值
|
||||||
|
ship-retry-delay-ms: 2000
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -216,6 +216,16 @@ jarvis:
|
|||||||
cron: "0 */5 * * * ?"
|
cron: "0 */5 * * * ?"
|
||||||
order-delay-ms: 250
|
order-delay-ms: 250
|
||||||
max-orders-per-round: 0
|
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:
|
fetch-comments:
|
||||||
base-url: http://192.168.8.60:5008
|
base-url: http://192.168.8.60:5008
|
||||||
|
|||||||
@@ -71,4 +71,19 @@ public class JarvisGoofishProperties {
|
|||||||
|
|
||||||
/** 与 defaultShipExpressCode 配套的展示名称 */
|
/** 与 defaultShipExpressCode 配套的展示名称 */
|
||||||
private String defaultShipExpressName = "日日顺";
|
private String defaultShipExpressName = "日日顺";
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 两次发货接口调用的最小间隔(毫秒)。开放平台对连续发货会返回「操作频繁」。
|
||||||
|
*/
|
||||||
|
private long shipMinIntervalMs = 2000L;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 发货遇到「操作频繁」、超时等可恢复错误时的额外重试次数(不含首次)。
|
||||||
|
*/
|
||||||
|
private int shipRetryTimes = 3;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 发货重试基础等待(毫秒)。第 n 次重试等待 n × 该值。
|
||||||
|
*/
|
||||||
|
private long shipRetryDelayMs = 2000L;
|
||||||
}
|
}
|
||||||
|
|||||||
+256
-46
@@ -29,6 +29,7 @@ import java.util.ArrayList;
|
|||||||
import java.util.Date;
|
import java.util.Date;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.Objects;
|
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:";
|
private static final String REDIS_WAYBILL_KEY_PREFIX = "logistics:waybill:order:";
|
||||||
/** F-闲鱼等京东单:运单可按第三方单号=闲鱼单号写入 Redis,见 LogisticsServiceImpl 镜像写入 */
|
/** F-闲鱼等京东单:运单可按第三方单号=闲鱼单号写入 Redis,见 LogisticsServiceImpl 镜像写入 */
|
||||||
public static final String REDIS_WAYBILL_GOOFISH_ORDER_PREFIX = "logistics:waybill:goofish:";
|
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
|
@Resource
|
||||||
private ErpGoofishOrderMapper erpGoofishOrderMapper;
|
private ErpGoofishOrderMapper erpGoofishOrderMapper;
|
||||||
@@ -565,34 +575,7 @@ public class GoofishOrderPipeline {
|
|||||||
ship.setWaybillNo(waybill.trim());
|
ship.setWaybillNo(waybill.trim());
|
||||||
ship.setExpressCode(expressCode);
|
ship.setExpressCode(expressCode);
|
||||||
ship.setExpressName(expressName);
|
ship.setExpressName(expressName);
|
||||||
String resp = ship.getResponseBody();
|
submitShip(row, ship, expressCode, waybill.trim(), manualRetry);
|
||||||
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));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
} catch (Exception ex) {
|
} catch (Exception ex) {
|
||||||
patchShipError(row, ex.getMessage());
|
patchShipError(row, ex.getMessage());
|
||||||
if (goofishOrderChangeLogger != null) {
|
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) {
|
private void patchShipError(ErpGoofishOrder row, String err) {
|
||||||
ErpGoofishOrder p = new ErpGoofishOrder();
|
ErpGoofishOrder p = new ErpGoofishOrder();
|
||||||
p.setId(row.getId());
|
p.setId(row.getId());
|
||||||
@@ -1197,30 +1394,43 @@ public class GoofishOrderPipeline {
|
|||||||
}
|
}
|
||||||
|
|
||||||
public int syncWaybillAndTryShipBatch(int limit) {
|
public int syncWaybillAndTryShipBatch(int limit) {
|
||||||
List<Integer> statuses = resolveAutoShipOrderStatuses();
|
Boolean got = stringRedisTemplate.opsForValue().setIfAbsent(REDIS_SHIP_BATCH_LOCK, "1", 10, TimeUnit.MINUTES);
|
||||||
List<ErpGoofishOrder> rows = erpGoofishOrderMapper.selectPendingShip(statuses, limit);
|
if (!Boolean.TRUE.equals(got)) {
|
||||||
if (rows == null) {
|
log.info("闲管家发货扫描已在进行,跳过本轮");
|
||||||
return 0;
|
return 0;
|
||||||
}
|
}
|
||||||
int n = 0;
|
try {
|
||||||
for (ErpGoofishOrder row : rows) {
|
List<Integer> statuses = resolveAutoShipOrderStatuses();
|
||||||
ErpGoofishOrder full = erpGoofishOrderMapper.selectById(row.getId());
|
List<ErpGoofishOrder> rows = erpGoofishOrderMapper.selectPendingShip(statuses, limit);
|
||||||
if (full == null) {
|
if (rows == null) {
|
||||||
continue;
|
return 0;
|
||||||
}
|
}
|
||||||
tryLinkJdOrder(full);
|
int n = 0;
|
||||||
syncWaybillFromRedis(full);
|
for (ErpGoofishOrder row : rows) {
|
||||||
if ((full.getRefundStatus() == null || full.getRefundStatus() == 0)
|
ErpGoofishOrder full = erpGoofishOrderMapper.selectById(row.getId());
|
||||||
&& statuses.contains(full.getOrderStatus())
|
if (full == null) {
|
||||||
&& StringUtils.isEmpty(full.getLocalWaybillNo())
|
continue;
|
||||||
&& StringUtils.isEmpty(full.getDetailWaybillNo())) {
|
}
|
||||||
refreshDetail(full);
|
tryLinkJdOrder(full);
|
||||||
syncWaybillFromRedis(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;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user