package com.hwtd.mes.collect.handler; import com.hwtd.mes.collect.dto.SerialPortDTO; import com.fazecast.jSerialComm.SerialPort; import com.fazecast.jSerialComm.SerialPortDataListener; import com.fazecast.jSerialComm.SerialPortEvent; import lombok.Data; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; import java.nio.charset.StandardCharsets; import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.Locale; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.locks.ReentrantLock; /** * 串口监听(手动模式) *

* 与原来“项目启动就开始监听”的方式不同,现在改为手动监听: *

* 串口参数(串口名、波特率、数据位、停止位、校验位、流控、读取超时、结束标志、单轮条数上限)全部由接口传入, * 不传时使用 {@link SerialPortDTO} 里的默认值;参数和当前生效的参数不一致时会自动按新参数重新打开串口, * 参数一致则复用已经打开的串口,不会重复打开。 *

* 并发说明: *

* 数据缓存说明:采集到的数据按队列缓存,插入最新的一条,超出 maxCount 条(默认 6 条)时移除最旧的一条, * 因此接口每次最多返回 maxCount 条最新的数据(按接收先后顺序);数据被取走后队列清空,重新开始缓存。 *

* 帧切分说明:默认按结束标志 endMark 切分数据帧;如果设备最后一帧不带结束标志(例如每帧以 * 开头而不是结尾), * 串口空闲 frameIdleMillis 毫秒(默认 200ms)后会把接收缓冲区里剩下的数据作为一帧取出,避免漏掉最后一条。 * * @author data-acquisition */ @Slf4j @Component public class SerialPortListener { /** 一直收不到结束标志时接收缓冲区的最大长度,超过则丢弃,防止内存无限增长 */ private static final int MAX_RECEIVE_BUFFER_LENGTH = 4096; /** * 串口开关锁:开启 / 关闭串口互斥。 * 多线程并发调用采集接口时,只有一个线程能真正打开串口,其它线程复用已经打开的串口。 */ private final ReentrantLock portLock = new ReentrantLock(); /** 串口对象,只在持有 {@link #portLock} 时创建和关闭;外部只做状态判断 */ private volatile SerialPort serialPort; /** 是否正在监听(volatile:串口回调线程与业务线程之间立即可见) */ private volatile boolean listening = false; /** 当前生效的串口参数,只在持有 {@link #portLock} 时变更 */ private volatile ActiveSetting activeSetting; /** 最后一次收到串口数据的时间,用于判断串口是否已经空闲(空闲补帧) */ private volatile long lastDataTime = 0L; /** 空闲补帧任务:设备最后一帧不带结束标志时,靠它把剩余数据取出来 */ private final ScheduledExecutorService frameFlushExecutor = Executors.newSingleThreadScheduledExecutor(runnable -> { Thread thread = new Thread(runnable, "serial-port-frame-flush"); thread.setDaemon(true); return thread; }); /** 接收缓冲区,由串口回调线程写入,用自身作为锁保护 */ private final StringBuilder receiveBuffer = new StringBuilder(); /** * 采集到的数据:串口回调线程写入,业务线程读取并清空,使用同步集合保证线程安全。 * 队列式缓存:从尾部插入最新的一条,超出 maxCount 时从头部移除最旧的一条,最多保留 maxCount 条数据。 */ private final List dataList = Collections.synchronizedList(new ArrayList<>()); /** * 开启串口监听(幂等操作) *

* 多线程同时调用时:第一个线程真正打开串口并注册监听器,其它线程直接复用,不会重复打开串口。 * 串口参数与当前生效参数不一致时会重新打开串口;串口被拔出导致异常断开时,下次调用会自动重新打开。 * * @param setting 接口传入的串口参数,为 null 时全部使用默认值 * @return true=监听已开启(包含本次调用之前就已开启的情况);false=开启失败(串口不存在、被占用,或串口功能未启用) * @throws IllegalArgumentException 串口参数非法(由 GlobalExceptionHandler 返回给前端) */ public boolean startListening(SerialPortDTO setting) { // 参数校验、补默认值放在加锁之前:参数不对直接抛异常,不占用串口 ActiveSetting requested = toActiveSetting(setting); portLock.lock(); try { // 已经在监听、串口正常且参数没有变化:直接复用(重复开启的唯一判定点) if (listening && isPortOpen() && requested.equals(activeSetting)) { log.debug("串口 {} 监听已开启且参数未变化,本次调用直接复用", requested.getSerialPortName()); return true; } if (activeSetting != null && !requested.equals(activeSetting)) { log.info("串口参数发生变化,按新参数重新打开串口:{} -> {}", activeSetting, requested); } // 首次开启、参数变化、上次打开失败、或串口异常断开时,先清理残留状态再重新打开 closePortInternal(); if (!openPort(requested)) { return false; } activeSetting = requested; log.info("串口 {} 监听已开启({})", requested.getSerialPortName(), requested.describe()); return true; } finally { portLock.unlock(); } } /** * 关闭串口监听(幂等操作) * * @return true=已关闭(包含本次调用之前就未开启的情况);false=关闭过程中出现异常 */ public boolean closeListening() { portLock.lock(); try { boolean success = closePortInternal(); if (success) { log.info("串口 {} 监听已关闭", getListenName()); } else { log.warn("串口 {} 监听关闭失败,请检查串口状态", getListenName()); } return success; } finally { portLock.unlock(); } } /** * 取走当前采集到的数据(取走后清空缓冲区),供采集接口调用。 * 加锁 + 快照 + 清空是一个原子操作,多线程并发取数时同一条数据只会被一个线程取走。 * * @return 本次取走的数据,没有数据时返回空集合 */ public List drainData() { synchronized (dataList) { if (dataList.isEmpty()) { return new ArrayList<>(); } List snapshot = new ArrayList<>(dataList); dataList.clear(); log.info("串口 {} 本次取走 {} 条数据", getListenName(), snapshot.size()); return snapshot; } } /** * @return 是否正在监听 */ public boolean isListening() { return listening && isPortOpen(); } /** * @return 串口是否已经打开 */ public boolean isPortOpen() { SerialPort port = serialPort; return port != null && port.isOpen(); } /** * @return 当前缓冲区中还未被取走的数据条数 */ public int dataSize() { synchronized (dataList) { return dataList.size(); } } /** * @return 当前生效的串口名称;还没开启过监听时返回配置的默认串口名称 */ public String getListenName() { ActiveSetting current = activeSetting; return current != null ? current.getSerialPortName() : ""; } /** * 项目停止时关闭串口 */ @PreDestroy public void destroy() { log.info("项目停止,关闭串口 {} 监听", getListenName()); closeListening(); frameFlushExecutor.shutdownNow(); } /** * 把接口传入的串口参数补齐默认值并校验,转换成当前会话使用的参数快照 * * @param setting 接口传入的参数,为 null 时全部使用默认值 * @return 归一化后的串口参数 * @throws IllegalArgumentException 参数非法 */ private ActiveSetting toActiveSetting(SerialPortDTO setting) { SerialPortDTO source = setting == null ? new SerialPortDTO() : setting; String portName = isBlank(source.getSerialPortName()) ? "" : source.getSerialPortName().trim(); if (isBlank(portName)) { throw new IllegalArgumentException("串口名称不能为空:请通过 listenName 参数传入"); } int baudRate = intOrDefault(source.getBaudRate(), SerialPortDTO.DEFAULT_BAUD_RATE); if (baudRate <= 0) { throw new IllegalArgumentException("串口参数 baudRate 非法:" + baudRate + ",必须大于 0"); } int dataBits = checkRange("dataBits", intOrDefault(source.getDataBits(), SerialPortDTO.DEFAULT_DATA_BITS), 5, 8); int stopBits = checkRange("stopBits", intOrDefault(source.getStopBits(), SerialPortDTO.DEFAULT_STOP_BITS), 1, 3); int readTimeout = intOrDefault(source.getReadTimeout(), SerialPortDTO.DEFAULT_READ_TIMEOUT); if (readTimeout < 0) { throw new IllegalArgumentException("串口参数 readTimeout 非法:" + readTimeout + ",不能小于 0"); } int maxCount = intOrDefault(source.getMaxCount(), SerialPortDTO.DEFAULT_MAX_COUNT); if (maxCount < 1) { throw new IllegalArgumentException("串口参数 maxCount 非法:" + maxCount + ",不能小于 1"); } int frameIdleMillis = intOrDefault(source.getFrameIdleMillis(), SerialPortDTO.DEFAULT_FRAME_IDLE_MILLIS); if (frameIdleMillis < 0) { throw new IllegalArgumentException("串口参数 frameIdleMillis 非法:" + frameIdleMillis + ",不能小于 0"); } String endMark = isBlank(source.getEndMark()) ? SerialPortDTO.DEFAULT_END_MARK : source.getEndMark(); String parity = isBlank(source.getParity()) ? SerialPortDTO.DEFAULT_PARITY : source.getParity(); String flowControl = isBlank(source.getFlowControl()) ? SerialPortDTO.DEFAULT_FLOW_CONTROL : source.getFlowControl(); return new ActiveSetting(portName, baudRate, dataBits, stopBits, parseParity(parity), parseFlowControl(flowControl), readTimeout, endMark, maxCount, frameIdleMillis); } /** * 打开串口并注册数据监听器(调用方必须持有 {@link #portLock}) * * @return 是否打开成功 */ private boolean openPort(ActiveSetting setting) { SerialPort port; try { port = SerialPort.getCommPort(setting.getSerialPortName()); port.setBaudRate(setting.getBaudRate()); // 波特率 port.setNumDataBits(setting.getDataBits()); // 数据位 port.setNumStopBits(setting.getStopBits()); // 停止位 port.setParity(setting.getParity()); // 校验位 port.setFlowControl(setting.getFlowControl()); // 流控 port.setComPortTimeouts(SerialPort.TIMEOUT_READ_SEMI_BLOCKING, setting.getReadTimeout(), 0); // 读取超时 if (!port.openPort()) { log.error("串口 {} 打开失败,请检查串口是否存在、是否被其他程序占用", setting.getSerialPortName()); return false; } } catch (Throwable t) { // 串口名称非法、参数不被驱动支持、驱动异常等情况统一返回失败,避免把异常抛给调用方(HTTP 线程) log.error("串口 {} 打开异常:{}", setting.getSerialPortName(), t.getMessage(), t); return false; } serialPort = port; // 新一轮采集开始:清掉上一轮未接收完整的半包和没被取走的数据,避免上一轮的数据混入本轮结果 synchronized (receiveBuffer) { receiveBuffer.setLength(0); } discardCollectedData(); // 先置为监听中再注册监听器,避免注册完成后第一批数据因状态未就绪被丢弃 listening = true; try { port.addDataListener(buildDataListener(port, setting)); } catch (Exception e) { log.error("串口 {} 注册数据监听器失败", setting.getSerialPortName(), e); // 注册失败时回滚,避免串口被占用却读不到数据 closePortInternal(); return false; } return true; } /** * 关闭串口并释放资源(调用方必须持有 {@link #portLock},本方法不会抛出异常) * * @return 是否正常关闭 */ private boolean closePortInternal() { SerialPort port = serialPort; serialPort = null; // 先置为未监听:串口回调线程不再处理后续数据 listening = false; boolean success = true; if (port != null) { try { port.removeDataListener(); // 移除监听,避免关闭过程中回调还在触发 if (port.isOpen() && !port.closePort()) { log.warn("串口 {} 关闭返回失败,请确认串口是否已被拔出", getListenName()); success = false; } } catch (Exception e) { log.error("关闭串口 {} 出现异常", getListenName(), e); success = false; } } // 丢弃未接收完整的半包数据,避免下次开启时拼出错误数据 synchronized (receiveBuffer) { receiveBuffer.setLength(0); } return success; } /** * 丢弃上一轮没有取走的数据(调用方必须持有 {@link #portLock}) */ private void discardCollectedData() { List abandoned = drainData(); if (!abandoned.isEmpty()) { log.warn("串口 {} 开启新一轮监听,丢弃上一轮未被取走的 {} 条数据:{}", getListenName(), abandoned.size(), abandoned); } } /** * 构建串口数据监听器,数据到达时在 jSerialComm 的监听线程中回调 *

* 本轮使用的参数(结束标志、条数上限)直接固化在监听器里,避免解析时读到下一轮的参数。 */ private SerialPortDataListener buildDataListener(final SerialPort port, final ActiveSetting setting) { return new SerialPortDataListener() { /** 指定监听的事件类型:数据可用时触发 */ @Override public int getListeningEvents() { return SerialPort.LISTENING_EVENT_DATA_AVAILABLE; } /** 数据可用时的处理逻辑 */ @Override public void serialEvent(SerialPortEvent event) { if (event.getEventType() != SerialPort.LISTENING_EVENT_DATA_AVAILABLE) { return; } // 关闭串口后可能还有一次回调在路上,这里再兜底判断一次 if (!listening || !setting.equals(activeSetting)) { return; } try { int available = port.bytesAvailable(); if (available <= 0) { return; } // 读取缓冲区中的所有字节 byte[] buffer = new byte[available]; int readBytes = port.readBytes(buffer, buffer.length); if (readBytes > 0) { // 把本次收到的零散数据转成字符串,追加到缓冲区 parseAndCollect(setting, new String(buffer, 0, readBytes, StandardCharsets.UTF_8)); } } catch (Exception e) { // 单次读取异常不影响后续监听 log.error("串口 {} 读取数据异常:{}", setting.getSerialPortName(), e.getMessage(), e); } } }; } /** * 解析接收到的数据:按结束标志切分,解析成数值后放入采集结果 * * @param setting 本轮采集使用的串口参数 * @param received 本次从串口读到的字符串 */ private void parseAndCollect(ActiveSetting setting, String received) { List completeValues = new ArrayList<>(); synchronized (receiveBuffer) { // 记录收到数据的时间,空闲补帧靠它判断串口是否已经没有新数据 lastDataTime = System.currentTimeMillis(); receiveBuffer.append(received); // 检查缓冲区是否包含“结束标志”,循环处理所有完整数据 int endIndex = receiveBuffer.indexOf(setting.getEndMark()); while (endIndex != -1) { // 提取“开头到结束标志”的完整数据 String completeData = normalizeFrame(receiveBuffer.substring(0, endIndex)); // 移除缓冲区中已处理的部分(保留剩余未完成的内容) receiveBuffer.delete(0, endIndex + setting.getEndMark().length()); if (!completeData.isEmpty()) { try { completeValues.add(completeData); } catch (NumberFormatException e) { // 干扰数据只忽略当前这一条,不影响后续数据解析 log.warn("串口 {} 收到无法解析的数据,已忽略:[{}]", setting.getSerialPortName(), completeData); } } endIndex = receiveBuffer.indexOf(setting.getEndMark()); } // 一直没有结束标志时缓冲区会持续增长,超过上限直接丢弃,防止内存泄漏 if (receiveBuffer.length() > MAX_RECEIVE_BUFFER_LENGTH) { log.warn("串口 {} 接收缓冲区超过 {} 个字符仍未收到结束标志,已丢弃缓存数据", setting.getSerialPortName(), MAX_RECEIVE_BUFFER_LENGTH); receiveBuffer.setLength(0); } } if (completeValues.isEmpty()) { return; } enqueue(setting, completeValues); } /** * 把数据放入缓存队列:从尾部插入最新的一条,超出 maxCount 时从头部移除最旧的一条 */ private void enqueue(ActiveSetting setting, List values) { synchronized (dataList) { for (String value : values) { dataList.add(value); while (dataList.size() > setting.getMaxCount()) { String removed = dataList.remove(0); log.debug("串口 {} 缓存已满 {} 条,移除最旧的一条数据 {}", setting.getSerialPortName(), setting.getMaxCount(), removed); } log.info("串口 {} 接收到数据-->{}", setting.getSerialPortName(), value); } } } /** * 启动空闲补帧任务:设备最后一帧不带结束标志时,靠它把剩余数据取出来 */ @PostConstruct void startFrameFlushTask() { frameFlushExecutor.scheduleWithFixedDelay(this::flushIdleFrame, 100, 100, TimeUnit.MILLISECONDS); } /** * 空闲补帧:串口超过 frameIdleMillis 没有新数据时,把接收缓冲区里剩下的数据当作一帧取出来。 *

* 典型场景:设备每帧以某个字符“开头”而不是“结尾”(如 *0 1 0), * 这时最后一帧不会有后续数据来触发切分,只靠结束标志永远取不出来。 * 如果设备每帧都带结束标志,缓冲区在每次切分后都是空的,这里不会做任何事。 */ private void flushIdleFrame() { try { ActiveSetting setting = activeSetting; if (!listening || setting == null || setting.getFrameIdleMillis() <= 0) { return; } String pending; synchronized (receiveBuffer) { if (receiveBuffer.length() == 0) { return; } long idle = System.currentTimeMillis() - lastDataTime; if (idle < setting.getFrameIdleMillis()) { return; } pending = normalizeFrame(receiveBuffer.toString()); receiveBuffer.setLength(0); log.info("串口 {} 已空闲 {}ms 没有收到新的结束标志,把缓冲区剩余数据 [{}] 作为一帧取出", setting.getSerialPortName(), idle, pending); } if (!pending.isEmpty()) { enqueue(setting, Collections.singletonList(pending)); } } catch (Exception e) { // 补帧异常不能影响后续补帧 log.error("串口空闲补帧异常:{}", e.getMessage(), e); } } /** * 去掉数据里的换行和首尾空白 */ private static String normalizeFrame(String value) { return value.replace("\n", "").replace("\r", "").trim(); } /** * 解析校验位参数:支持 NONE/ODD/EVEN/MARK/SPACE(不区分大小写),也兼容 0~4 的数字 */ private static int parseParity(String parity) { String value = parity.trim().toUpperCase(Locale.ROOT); switch (value) { case "NONE": case "N": case "NO": return SerialPort.NO_PARITY; case "ODD": case "O": return SerialPort.ODD_PARITY; case "EVEN": case "E": return SerialPort.EVEN_PARITY; case "MARK": case "M": return SerialPort.MARK_PARITY; case "SPACE": case "S": return SerialPort.SPACE_PARITY; default: int number = parseInt(value, -1); if (number >= SerialPort.NO_PARITY && number <= SerialPort.SPACE_PARITY) { return number; } throw new IllegalArgumentException("串口参数 parity 非法:" + parity + ",可选 NONE/ODD/EVEN/MARK/SPACE,或 0~4 的数字"); } } /** * 解析流控参数:支持 NONE/RTS_CTS/XON_XOFF(不区分大小写),也兼容 jSerialComm 的数字常量 */ private static int parseFlowControl(String flowControl) { String value = flowControl.trim().toUpperCase(Locale.ROOT).replace('-', '_'); switch (value) { case "NONE": case "OFF": case "DISABLED": return SerialPort.FLOW_CONTROL_DISABLED; case "RTS_CTS": case "RTSCTS": case "HARDWARE": return SerialPort.FLOW_CONTROL_RTS_ENABLED | SerialPort.FLOW_CONTROL_CTS_ENABLED; case "XON_XOFF": case "XONXOFF": case "SOFTWARE": return SerialPort.FLOW_CONTROL_XONXOFF_IN_ENABLED | SerialPort.FLOW_CONTROL_XONXOFF_OUT_ENABLED; default: int number = parseInt(value, -1); if (number >= 0) { return number; } throw new IllegalArgumentException("串口参数 flowControl 非法:" + flowControl + ",可选 NONE/RTS_CTS/XON_XOFF,或数字常量"); } } /** * 校验取值范围 */ private static int checkRange(String name, int value, int min, int max) { if (value < min || value > max) { throw new IllegalArgumentException("串口参数 " + name + " 非法:" + value + ",取值范围 " + min + "~" + max); } return value; } /** * 接口没传(null)时使用默认值 */ private static int intOrDefault(Integer value, int defaultValue) { return value == null ? defaultValue : value; } /** * 解析数字,解析不了返回默认值 */ private static int parseInt(String value, int defaultValue) { try { return Integer.parseInt(value); } catch (NumberFormatException e) { return defaultValue; } } private static boolean isBlank(String value) { return value == null || value.trim().isEmpty(); } /** * 当前生效的串口参数(接口参数补默认值、校验、归一化之后的快照,用于判断参数是否发生变化) */ @Data private static final class ActiveSetting { /** 串口名称 */ private final String serialPortName; /** 波特率 */ private final int baudRate; /** 数据位 */ private final int dataBits; /** 停止位 */ private final int stopBits; /** 校验位(jSerialComm 常量) */ private final int parity; /** 流控(jSerialComm 常量) */ private final int flowControl; /** 读取超时(毫秒) */ private final int readTimeout; /** 一帧数据的结束标志 */ private final String endMark; /** 单轮采集最多保留的数据条数 */ private final int maxCount; /** 帧空闲补帧时间(毫秒),0 表示关闭 */ private final int frameIdleMillis; /** * 日志用的参数描述 */ String describe() { return "波特率 " + baudRate + ",数据位 " + dataBits + ",停止位 " + stopBits + ",校验位 " + parityName() + ",流控 " + flowControlName() + ",读取超时 " + readTimeout + "ms" + ",结束标志 [" + endMark + "]" + ",单轮最多 " + maxCount + " 条" + ",空闲补帧 " + (frameIdleMillis > 0 ? frameIdleMillis + "ms" : "关闭"); } private String parityName() { switch (parity) { case SerialPort.ODD_PARITY: return "奇校验"; case SerialPort.EVEN_PARITY: return "偶校验"; case SerialPort.MARK_PARITY: return "MARK"; case SerialPort.SPACE_PARITY: return "SPACE"; default: return "无校验"; } } private String flowControlName() { if (flowControl == SerialPort.FLOW_CONTROL_DISABLED) { return "无"; } if (flowControl == (SerialPort.FLOW_CONTROL_RTS_ENABLED | SerialPort.FLOW_CONTROL_CTS_ENABLED)) { return "RTS/CTS"; } if (flowControl == (SerialPort.FLOW_CONTROL_XONXOFF_IN_ENABLED | SerialPort.FLOW_CONTROL_XONXOFF_OUT_ENABLED)) { return "XON/XOFF"; } return String.valueOf(flowControl); } } }