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; /** * 串口监听(手动模式) *
* 与原来“项目启动就开始监听”的方式不同,现在改为手动监听: *
* 并发说明: *
* 帧切分说明:默认按结束标志 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
* 多线程同时调用时:第一个线程真正打开串口并注册监听器,其它线程直接复用,不会重复打开串口。
* 串口参数与当前生效参数不一致时会重新打开串口;串口被拔出导致异常断开时,下次调用会自动重新打开。
*
* @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
* 本轮使用的参数(结束标志、条数上限)直接固化在监听器里,避免解析时读到下一轮的参数。
*/
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
* 典型场景:设备每帧以某个字符“开头”而不是“结尾”(如 *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);
}
}
}