package com.gaotao.modules.base.task;
import com.gaotao.modules.base.entity.PrintTask;
import com.gaotao.modules.base.service.PrintTaskService;
import com.google.gson.Gson;
import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import java.io.OutputStream;
import java.net.InetSocketAddress;
import java.net.Socket;
import java.util.Date;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.ThreadFactory;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.nio.charset.StandardCharsets;
/**
* 打印任务定时调度器
*
*
功能说明:
*
* - 每2秒扫描一次打印任务队列
* - 按站点(site)拆分任务,不同站点并发处理
* - 同一站点内按FIFO顺序(created_date)取出PENDING状态的任务
* - 执行打印任务并更新状态
* - 失败任务记录错误信息
*
*
* 并发控制:
*
* - 使用数据库乐观锁(markAsProcessing时检查状态)
* - 确保同一任务不会被重复执行
*
*/
@Slf4j
@Component
public class PrintTaskScheduler {
@Autowired
private PrintTaskService printTaskService;
private static final Gson gson = new Gson();
/**
* 打印调度总开关(未配置时默认开启)
*/
@Value("${dashboard.print.enabled:false}")
private boolean dashboardPushEnabled;
/**
* 每轮调度最多扫描的站点数量(未配置时默认20)
*/
@Value("${print.scheduler.site.max-sites-per-round:20}")
private int maxSitesPerRound;
/**
* 每个站点每轮最多处理的任务数量(未配置时默认10)
*/
@Value("${print.scheduler.site.batch-size:10}")
private int siteBatchSize;
/**
* 站点并发线程池大小(未配置时默认8)
*/
@Value("${print.scheduler.site.thread-pool-size:8}")
private int siteThreadPoolSize;
/**
* 打印机连接超时(毫秒,未配置时默认3000)
*/
@Value("${print.scheduler.printer-connect-timeout-ms:3000}")
private int printerConnectTimeoutMs;
/**
* 打印机读写超时(毫秒,未配置时默认10000)
*/
@Value("${print.scheduler.printer-read-timeout-ms:10000}")
private int printerReadTimeoutMs;
/** Maximum age of a PROCESSING task before it is considered abandoned. */
@Value("${print.scheduler.stale-processing-timeout-minutes:30}")
private int staleProcessingTimeoutMinutes;
/** How often abandoned PROCESSING tasks are recovered. */
@Value("${print.scheduler.recovery-interval-ms:60000}")
private long recoveryIntervalMs;
private static final AtomicInteger THREAD_INDEX = new AtomicInteger(1);
private static final AtomicLong LAST_IDLE_LOG_AT = new AtomicLong(0);
private final Map siteProcessingFlags = new ConcurrentHashMap<>();
private ExecutorService siteTaskExecutor;
@PostConstruct
public void initSiteExecutor() {
if (maxSitesPerRound <= 0) {
log.warn("Invalid print.scheduler.site.max-sites-per-round={}, using 20", maxSitesPerRound);
maxSitesPerRound = 20;
}
if (siteBatchSize <= 0) {
log.warn("Invalid print.scheduler.site.batch-size={}, using 10", siteBatchSize);
siteBatchSize = 10;
}
if (staleProcessingTimeoutMinutes <= 0) {
log.warn("Invalid print.scheduler.stale-processing-timeout-minutes={}, using 30",
staleProcessingTimeoutMinutes);
staleProcessingTimeoutMinutes = 30;
}
if (recoveryIntervalMs < 1000) {
log.warn("Invalid print.scheduler.recovery-interval-ms={}, using 60000", recoveryIntervalMs);
recoveryIntervalMs = 60000;
}
int actualThreadPoolSize = Math.max(1, siteThreadPoolSize);
if (actualThreadPoolSize != siteThreadPoolSize) {
log.warn("打印站点线程池配置无效,使用默认最小值1,原配置={}", siteThreadPoolSize);
}
ThreadFactory threadFactory = runnable -> {
Thread thread = new Thread(runnable);
thread.setName("print-site-worker-" + THREAD_INDEX.getAndIncrement());
return thread;
};
siteTaskExecutor = Executors.newFixedThreadPool(actualThreadPoolSize, threadFactory);
if (dashboardPushEnabled) {
log.info("Print task scheduler config: enabled={}, maxSites={}, batchSize={}, siteWorkers={}, " +
"staleTimeoutMinutes={}, recoveryIntervalMs={}",
dashboardPushEnabled, maxSitesPerRound, siteBatchSize, actualThreadPoolSize,
staleProcessingTimeoutMinutes, recoveryIntervalMs);
}
log.info("打印任务站点线程池初始化完成,线程数={}", actualThreadPoolSize);
}
@PreDestroy
public void shutdownSiteExecutor() {
if (siteTaskExecutor != null) {
siteTaskExecutor.shutdown();
log.info("打印任务站点线程池已关闭");
}
}
// 每2秒执行一次
@Scheduled(fixedDelayString = "${print.scheduler.poll-interval-ms:2000}",
scheduler = "printQueueTaskScheduler")
public void processPrintTasks() {
// 检查定时任务开关
if (!dashboardPushEnabled) {
return;
} else {
log.debug("Print task scheduler tick: enabled={}, thread={}",
dashboardPushEnabled, Thread.currentThread().getName());
}
if (siteTaskExecutor == null || siteTaskExecutor.isShutdown()) {
log.warn("打印任务站点线程池不可用,跳过本轮调度");
return;
}
try {
// 1. 获取有待执行任务的站点列表
List pendingSites = printTaskService.getPendingSites(maxSitesPerRound);
if (pendingSites.isEmpty()) {
long now = System.currentTimeMillis();
long lastLog = LAST_IDLE_LOG_AT.get();
if (now - lastLog >= 30000 && LAST_IDLE_LOG_AT.compareAndSet(lastLog, now)) {
log.info("Print task scheduler is running; no PENDING sites found");
}
return; // 没有待处理任务
}
log.info("【打印队列】发现 {} 个待执行站点", pendingSites.size());
// 2. 按站点并发处理,不同站点互不影响
for (String site : pendingSites) {
if (site == null || site.trim().isEmpty()) {
continue;
}
AtomicBoolean siteProcessingFlag = siteProcessingFlags
.computeIfAbsent(site, key -> new AtomicBoolean(false));
// 同一个site只允许一个worker在执行,避免站内并发
if (!siteProcessingFlag.compareAndSet(false, true)) {
log.debug("【打印队列】site={} 正在执行中,跳过本轮", site);
continue;
}
try {
siteTaskExecutor.submit(() -> processSiteTasks(site, siteProcessingFlag));
} catch (Exception submitEx) {
siteProcessingFlag.set(false);
log.error("【打印队列】提交site任务失败 site={}, error={}",
site, submitEx.getMessage(), submitEx);
}
}
} catch (Exception e) {
log.error("打印任务调度器执行异常: {}", e.getMessage(), e);
}
}
// Recover tasks left in PROCESSING after a server crash or forced restart.
@Scheduled(fixedDelayString = "${print.scheduler.recovery-interval-ms:60000}",
initialDelayString = "${print.scheduler.recovery-initial-delay-ms:5000}",
scheduler = "printQueueTaskScheduler")
public void recoverStaleProcessingTasks() {
if (!dashboardPushEnabled) {
return;
}
try {
long timeoutMillis = staleProcessingTimeoutMinutes * 60_000L;
Date before = new Date(System.currentTimeMillis() - timeoutMillis);
int recovered = printTaskService.resetStaleProcessingTasks(before);
if (recovered > 0) {
log.warn("Recovered {} stale PROCESSING print tasks older than {} minutes",
recovered, staleProcessingTimeoutMinutes);
}
} catch (Exception e) {
log.error("Failed to recover stale PROCESSING print tasks", e);
}
}
/**
* 处理指定站点的打印任务(站内串行)
*/
private void processSiteTasks(String site, AtomicBoolean siteProcessingFlag) {
try {
List pendingTasks = printTaskService.getPendingTasksBySite(site, siteBatchSize);
if (pendingTasks.isEmpty()) {
return;
}
log.info("【打印队列】site={} 发现 {} 个待执行任务", site, pendingTasks.size());
for (PrintTask task : pendingTasks) {
processTask(task);
}
} catch (Exception e) {
log.error("【打印队列】site任务执行异常 site={}, error={}", site, e.getMessage(), e);
} finally {
siteProcessingFlag.set(false);
}
}
/**
* 处理单个打印任务
*/
private void processTask(PrintTask task) {
Long taskId = task.getId();
String site = task.getSite();
String printerIp = task.getPrinterIp();
try {
// 1. 标记为执行中(乐观锁,防止并发重复执行)
boolean marked = printTaskService.markAsProcessing(taskId);
if (!marked) {
log.debug("任务 {} 已被其他线程处理,跳过", taskId);
return;
}
log.info("【打印队列】开始执行任务 taskId={}, site={}, 打印机={}, 份数={}",
taskId, site, printerIp, task.getCopies());
// 2. 解析标签数据
Map labelData = null;
if (task.getLabelData() != null && !task.getLabelData().trim().isEmpty()) {
labelData = gson.fromJson(task.getLabelData(), Map.class);
}
// 3. 执行打印
boolean success = executePrint(
printerIp,
task.getZplCode(),
task.getCopies(),
task.getRfidFlag(),
labelData
);
// 4. 更新任务状态
if (success) {
printTaskService.markAsSuccess(taskId);
log.info("【打印队列】✓ 任务完成 taskId={}, site={}, 打印机={}", taskId, site, printerIp);
} else {
printTaskService.markAsFailed(taskId, "打印失败,未返回成功标识");
log.error("【打印队列】✗ 任务失败 taskId={}, site={}, 打印机={}", taskId, site, printerIp);
}
} catch (Exception e) {
log.error("【打印队列】✗ 任务执行异常 taskId={}, site={}, 错误: {}", taskId, site, e.getMessage(), e);
printTaskService.markAsFailed(taskId, "执行异常: " + e.getMessage());
}
}
/**
* 执行打印(发送ZPL到打印机)
* 注意:ZPL代码已在入队时组装好(包括RFID指令),这里只需要直接发送
*
* @param printerIp 打印机IP
* @param zplCode ZPL代码(已包含RFID指令)
* @param copies 打印份数
* @param rfidFlag RFID标识(用于日志)
* @param labelData 标签数据(未使用,保留以备后续扩展)
* @return 是否成功
*/
private boolean executePrint(String printerIp, String zplCode, Integer copies,
String rfidFlag, Map labelData) {
Socket socket = null;
try {
// 1. 连接打印机
socket = new Socket();
socket.connect(new InetSocketAddress(printerIp, 9100), printerConnectTimeoutMs);
socket.setSoTimeout(printerReadTimeoutMs);
OutputStream os = socket.getOutputStream();
int actualCopies = copies != null && copies > 0 ? copies : 1;
// 2. 直接发送ZPL代码(已在入队时组装好,包括RFID指令)
log.debug("开始发送ZPL到打印机 {}, 份数: {}, RFID: {}", printerIp, actualCopies, rfidFlag);
for (int i = 0; i < actualCopies; i++) {
os.write(zplCode.getBytes(StandardCharsets.UTF_8));
os.flush();
// RFID标签需要等待处理时间
if ("Y".equals(rfidFlag)) {
Thread.sleep(1000); // RFID标签等待1秒
log.debug("RFID标签打印完成 {} / {}", i + 1, actualCopies);
}
}
log.debug("✓ ZPL已发送到打印机 {}, 份数: {}", printerIp, actualCopies);
return true;
} catch (Exception e) {
log.error("发送ZPL到打印机失败: printerIp={}, error={}", printerIp, e.getMessage(), e);
return false;
} finally {
if (socket != null) {
try {
socket.close();
} catch (Exception e) {
log.warn("关闭socket失败: {}", e.getMessage());
}
}
}
}
/**
* 归档打印任务历史数据
*
* 功能说明:
*
* - 每天凌晨5点自动执行
* - 将已完成的打印任务备份到 print_task_history 表
* - 清理10天之前的历史数据
*
*
* 执行时间:每天凌晨5点(固定)
* 归档规则:归档10天前的已完成任务(固定)
*/
@Scheduled(cron = "0 0 5 * * ?", scheduler = "printQueueTaskScheduler")
public void archivePrintTaskHistory() {
// 检查定时任务开关
if (!dashboardPushEnabled) {
return;
}
log.info("=== 开始归档打印任务历史数据 ===");
try {
// 调用Service层执行归档操作
int archivedCount = printTaskService.archivePrintTaskHistory();
log.info("=== 打印任务历史数据归档完成:归档任务数量={} ===", archivedCount);
} catch (Exception e) {
log.error("=== 打印任务历史数据归档失败 ===", e);
}
}
}