上海郢昱
10 小时以前 0d9956362ef415d0b49318746e9a2ef955dcfebd
质检设备数采功能实现
已删除11个文件
已重命名6个文件
已添加6个文件
已修改4个文件
2461 ■■■■■ 文件已修改
README.md 45 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
pom.xml 12 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/chinaztt/mes/docx/config/SerialPortStartupRunner.java 37 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/chinaztt/mes/docx/constant/FieldMatchRuleConstants.java 30 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/chinaztt/mes/docx/controller/DocxController.java 38 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/chinaztt/mes/docx/dto/GetFileDto.java 57 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/chinaztt/mes/docx/handler/SerialPortListener.java 145 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/chinaztt/mes/docx/pojo/Chemical.java 22 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/chinaztt/mes/docx/pojo/TestBatch.java 19 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/chinaztt/mes/docx/service/DocxService.java 17 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/chinaztt/mes/docx/service/impl/DocxServiceImpl.java 195 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/chinaztt/mes/docx/util/TakeWords.java 554 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/chinaztt/mes/docx/util/XMLFileListener.java 231 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/hwtd/mes/collect/DataAcquisitionApplication.java 2 ●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/hwtd/mes/collect/config/TomcatConfig.java 2 ●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/hwtd/mes/collect/constant/CommonConstants.java 2 ●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/hwtd/mes/collect/controller/DataCollectionController.java 79 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/hwtd/mes/collect/dto/DatabaseDTO.java 75 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/hwtd/mes/collect/dto/SerialPortDTO.java 66 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/hwtd/mes/collect/handler/GlobalExceptionHandler.java 4 ●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/hwtd/mes/collect/handler/SerialPortListener.java 644 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/hwtd/mes/collect/service/DataCollectionService.java 14 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/hwtd/mes/collect/service/impl/DataCollectionServiceImpl.java 159 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/java/com/hwtd/mes/collect/util/Result.java 4 ●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/resources/META-INF/MANIFEST.MF 2 ●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/main/resources/application.yml 4 ●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/test/java/com/hwtd/mes/collect/DataAcquisitionApplicationTests.java 2 ●●● 补丁 | 查看 | 原始文档 | blame | 历史
README.md
@@ -1,4 +1,47 @@
## data-acquisition
设备数据采集器
亨旺特导MES数据采集器
### ä¸²å£é‡‡é›†ï¼ˆæ‰‹åŠ¨ç›‘å¬ï¼‰
串口不再随项目启动自动监听,改为手动监听:调用采集接口时开启监听,采集结束后调用关闭接口。
| æŽ¥å£ | è¯´æ˜Ž |
| --- | --- |
| `GET /lims/getFile?fileExtension=.serialPort` | é‡‡é›†æŽ¥å£ï¼šè°ƒç”¨æ—¶è‡ªåŠ¨å¼€å¯ç›‘å¬ï¼Œå¹¶è¿”å›žæœ¬æ¬¡é‡‡é›†åˆ°çš„æ•°æ®ï¼›æ²¡æœ‰æ•°æ®æ—¶è¿”å›žç©ºé›†åˆï¼Œç”±å‰ç«¯è½®è¯¢å–æ•° |
| `GET /lims/openSerialPort` | æ‰‹åŠ¨å¼€å¯ç›‘å¬ï¼Œå¹¶è¿”å›žå½“å‰æŽ¥æ”¶åˆ°çš„æ•°æ®åˆ—è¡¨ï¼ˆå¹‚ç­‰ï¼Œå‚æ•°ä¸€è‡´æ—¶å¤ç”¨å·²æ‰“å¼€çš„ä¸²å£ï¼‰ |
| `GET /lims/closeSerialPort` | æ‰‹åŠ¨å…³é—­ç›‘å¬ï¼ˆå¹‚ç­‰ï¼Œé‡‡é›†ç»“æŸåŽè°ƒç”¨ï¼Œé¿å…ä¸²å£ä¸€ç›´è¢«å ç”¨ï¼‰ |
| `GET /lims/serialPortStatus` | æŸ¥è¯¢ç›‘听状态:`listenName` / `listening` / `portOpen` / `dataSize` |
串口参数全部由接口传入,`/lims/getFile`(`fileExtension=.serialPort` æ—¶ï¼‰å’Œ `/lims/openSerialPort` éƒ½æ”¯æŒï¼Œ
不传的字段使用默认值(`listenName` å¿…传)。application.yml ä¸­çš„ `serialPort` é…ç½®å·²ä¸å†ç”Ÿæ•ˆï¼Œå¯ä»¥åˆ é™¤ã€‚
| å‚æ•° | é»˜è®¤å€¼ | è¯´æ˜Ž |
| --- | --- | --- |
| `listenName` | å¿…ä¼  | ä¸²å£åç§°ï¼Œå¦‚ `COM12`;不传会返回“串口名称不能为空” |
| `baudRate` | `9600` | æ³¢ç‰¹çއ |
| `dataBits` | `8` | æ•°æ®ä½ï¼Œ5~8 |
| `stopBits` | `1` | åœæ­¢ä½ï¼Œ1、2、3(1.5 ä½) |
| `parity` | `NONE` | æ ¡éªŒä½ï¼š`NONE`/`ODD`/`EVEN`/`MARK`/`SPACE`,也兼容 0~4 |
| `flowControl` | `NONE` | æµæŽ§ï¼š`NONE`/`RTS_CTS`/`XON_XOFF`,也兼容数字常量 |
| `readTimeout` | `1000` | è¯»å–超时(毫秒) |
| `endMark` | `*` | ä¸€å¸§æ•°æ®çš„结束标志,特殊字符用 URL ç¼–码,如换行 `%0A` |
| `maxCount` | `6` | æ•°æ®ç¼“存队列长度:队列满时插入最新一条、移除最旧一条 |
| `frameIdleMillis` | `200` | å¸§ç©ºé—²è¡¥å¸§æ—¶é—´ï¼ˆæ¯«ç§’):超过该时间没有新数据,就把缓冲区剩余数据当成一帧取出;`0` è¡¨ç¤ºå…³é—­ |
示例:`/lims/getFile?fileExtension=.serialPort&listenName=COM12&baudRate=19200&parity=EVEN&endMark=%0A&maxCount=6`
参数与当前生效参数一致时复用已经打开的串口(不会重复打开);参数变化时会自动按新参数重新打开串口;
参数非法时接口直接返回具体原因(如 `串口参数 dataBits éžæ³•:9,取值范围 5~8`)。
采集到的数据按队列缓存:插入最新的一条,超过 `maxCount` æ¡ï¼ˆé»˜è®¤ 6 æ¡ï¼‰æ—¶ç§»é™¤æœ€æ—§çš„一条,
因此接口每次最多返回 `maxCount` æ¡æœ€æ–°çš„æ•°æ®ï¼ˆæŒ‰æŽ¥æ”¶å…ˆåŽé¡ºåºï¼‰ï¼›æ•°æ®è¢«å–走后队列清空,重新开始缓存。
数据帧按 `endMark` åˆ‡åˆ†ã€‚如果设备是「每帧以某个字符开头」而不是结尾(例如发 `*0 1 0`、`*0 2 0`,
`*` ç”¨äºŽåˆ†éš”前后两帧),那么最后一帧后面没有新的 `*` æ¥è§¦å‘切分,只靠结束标志会漏掉最后一条;
此时由 `frameIdleMillis`(默认 200ms)兜底:串口空闲超过该时间后,把接收缓冲区里剩下的数据当作一帧取出。
如果设备每帧都以结束标志结尾,缓冲区每次切分后都是空的,空闲补帧不会做任何事。
并发说明:开启 / å…³é—­ä¸²å£ç”±åŒä¸€æŠŠé”ä¿æŠ¤ï¼Œå¤šçº¿ç¨‹å¹¶å‘调用采集接口也只会打开一次串口,不会重复注册监听器;
采集结果存放在同步集合中,多线程同时取数时同一条数据只会被一个线程取走,不会重复也不会丢失。
调用关闭接口后再次调用采集接口会重新开启监听,并清掉上一轮残留的半包数据和未取走的数据,保证每次采集的数据只属于本轮。
pom.xml
@@ -38,6 +38,16 @@
            <groupId>mysql</groupId>
            <artifactId>mysql-connector-java</artifactId>
        </dependency>
        <dependency>
            <groupId>org.postgresql</groupId>
            <artifactId>postgresql</artifactId>
            <version>42.7.7</version> <!-- å»ºè®®ä½¿ç”¨è¾ƒæ–°ç¨³å®šç‰ˆ -->
        </dependency>
        <dependency>
        <groupId>com.zaxxer</groupId>
        <artifactId>HikariCP</artifactId>
        <version>4.0.3</version>
    </dependency>
        <!--lombok-->
        <dependency>
            <groupId>org.projectlombok</groupId>
@@ -148,7 +158,7 @@
                <artifactId>spring-boot-maven-plugin</artifactId>
                <version>${spring-boot.version}</version>
                <configuration>
                    <mainClass>com.chinaztt.mes.docx.DataAcquisitionApplication</mainClass>
                    <mainClass>com.hwtd.mes.collect.DataAcquisitionApplication</mainClass>
                    <skip>false</skip>
                </configuration>
                <executions>
src/main/java/com/chinaztt/mes/docx/config/SerialPortStartupRunner.java
ÎļþÒÑɾ³ý
src/main/java/com/chinaztt/mes/docx/constant/FieldMatchRuleConstants.java
ÎļþÒÑɾ³ý
src/main/java/com/chinaztt/mes/docx/controller/DocxController.java
ÎļþÒÑɾ³ý
src/main/java/com/chinaztt/mes/docx/dto/GetFileDto.java
ÎļþÒÑɾ³ý
src/main/java/com/chinaztt/mes/docx/handler/SerialPortListener.java
ÎļþÒÑɾ³ý
src/main/java/com/chinaztt/mes/docx/pojo/Chemical.java
ÎļþÒÑɾ³ý
src/main/java/com/chinaztt/mes/docx/pojo/TestBatch.java
ÎļþÒÑɾ³ý
src/main/java/com/chinaztt/mes/docx/service/DocxService.java
ÎļþÒÑɾ³ý
src/main/java/com/chinaztt/mes/docx/service/impl/DocxServiceImpl.java
ÎļþÒÑɾ³ý
src/main/java/com/chinaztt/mes/docx/util/TakeWords.java
ÎļþÒÑɾ³ý
src/main/java/com/chinaztt/mes/docx/util/XMLFileListener.java
ÎļþÒÑɾ³ý
src/main/java/com/hwtd/mes/collect/DataAcquisitionApplication.java
ÎļþÃû´Ó src/main/java/com/chinaztt/mes/docx/DataAcquisitionApplication.java ÐÞ¸Ä
@@ -1,4 +1,4 @@
package com.chinaztt.mes.docx;
package com.hwtd.mes.collect;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
src/main/java/com/hwtd/mes/collect/config/TomcatConfig.java
ÎļþÃû´Ó src/main/java/com/chinaztt/mes/docx/config/TomcatConfig.java ÐÞ¸Ä
@@ -1,4 +1,4 @@
package com.chinaztt.mes.docx.config;
package com.hwtd.mes.collect.config;
import org.apache.catalina.connector.Connector;
import org.springframework.boot.web.embedded.tomcat.TomcatServletWebServerFactory;
src/main/java/com/hwtd/mes/collect/constant/CommonConstants.java
ÎļþÃû´Ó src/main/java/com/chinaztt/mes/docx/constant/CommonConstants.java ÐÞ¸Ä
@@ -1,4 +1,4 @@
package com.chinaztt.mes.docx.constant;
package com.hwtd.mes.collect.constant;
public interface CommonConstants {
src/main/java/com/hwtd/mes/collect/controller/DataCollectionController.java
¶Ô±ÈÐÂÎļþ
@@ -0,0 +1,79 @@
package com.hwtd.mes.collect.controller;
import com.hwtd.mes.collect.dto.DatabaseDTO;
import com.hwtd.mes.collect.dto.SerialPortDTO;
import com.hwtd.mes.collect.handler.SerialPortListener;
import com.hwtd.mes.collect.service.DataCollectionService;
import com.hwtd.mes.collect.util.Result;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import java.util.LinkedHashMap;
import java.util.Map;
@RestController
@RequestMapping("/collection")
public class DataCollectionController {
    @Autowired
    private DataCollectionService dataCollectionService;
    @Autowired
    private SerialPortListener serialPortListener;
    /**
     * æ‰‹åŠ¨å¼€å¯ä¸²å£ç›‘å¬ï¼Œå¹¶è¿”å›žå½“å‰æŽ¥æ”¶åˆ°çš„æ•°æ®åˆ—è¡¨ï¼ˆé˜Ÿåˆ—å¼ç¼“å­˜ï¼Œæœ€å¤š maxCount æ¡ï¼Œæ²¡æœ‰æ•°æ®æ—¶è¿”回空集合)。
     * <p>
     * ä¸²å£å‚数通过请求参数传入(不传的字段使用默认值,listenName å¿…传):
     * listenName、baudRate、dataBits、stopBits、parity、flowControl、readTimeout、endMark、maxCount。
     * å¹‚等接口:参数一致时直接复用已打开的串口、不会重复打开,参数变化时按新参数重新打开;
     * å¤šçº¿ç¨‹å¹¶å‘调用同样只会打开一次。参数非法时由 GlobalExceptionHandler è¿”回具体原因。
     */
    @GetMapping("/openSerialPort")
    public Result<?> openSerialPort(SerialPortDTO serialPortDTO) {
        if (!serialPortListener.startListening(serialPortDTO)) {
            String listenName = serialPortDTO == null || serialPortDTO.getSerialPortName() == null
                    ? "" : serialPortDTO.getSerialPortName();
            return Result.failed("串口 " + listenName + " ç›‘听开启失败,请检查串口是否存在、是否被其他程序占用!");
        }
        return Result.ok(serialPortListener.drainData());
    }
    /**
     * æ‰‹åŠ¨å…³é—­ä¸²å£ç›‘å¬ã€‚
     * å¹‚等接口:未开启时调用也会正常返回,不会报错。
     */
    @GetMapping("/closeSerialPort")
    public Result<?> closeSerialPort() {
        if (serialPortListener.closeListening()) {
            return Result.ok("串口 " + serialPortListener.getListenName() + " ç›‘听已关闭");
        }
        return Result.failed("串口 " + serialPortListener.getListenName() + " ç›‘听关闭失败,请检查串口状态!");
    }
    /**
     * æŸ¥è¯¢ä¸²å£ç›‘听状态,便于前端判断是否需要调用开启 / å…³é—­æŽ¥å£
     */
    @GetMapping("/serialPortStatus")
    public Result<?> serialPortStatus() {
        Map<String, Object> status = new LinkedHashMap<>();
        status.put("listenName", serialPortListener.getListenName());
        status.put("listening", serialPortListener.isListening());
        status.put("portOpen", serialPortListener.isPortOpen());
        status.put("dataSize", serialPortListener.dataSize());
        return Result.ok(status);
    }
    @GetMapping("/getAccessData")
    public Result<?> getAccessData(DatabaseDTO databaseDTO) {
        return Result.ok(dataCollectionService.getAccessData(databaseDTO));
    }
    @GetMapping("/getPostgreSqlData")
    public Result<?> getPostgreSqlData(DatabaseDTO databaseDTO) {
        return Result.ok(dataCollectionService.getPostgreSqlData(databaseDTO));
    }
}
src/main/java/com/hwtd/mes/collect/dto/DatabaseDTO.java
¶Ô±ÈÐÂÎļþ
@@ -0,0 +1,75 @@
package com.hwtd.mes.collect.dto;
import lombok.Data;
/**
 *
 * æ•°æ®åº“采集参数
 * @author 27233
 * @date 2026-09-14 9:44
 */
@Data
public class DatabaseDTO {
    /**
     * æ•°æ®åº“地址
     */
    private String filePath;
    /**
     * ip地址
     */
    private String ipAddress;
    /**
     *
     * ç«¯å£å·
     */
    private Integer serverPort;
    /**
     * ä¸»é”®å­—段
     */
    private String mainColumn;
    /**
     * æ•°æ®åº“名
     */
    private String databaseName;
    /**
     * æ•°æ®è¡¨å
     */
    private String tableName;
    /**
     * ç”¨æˆ·å
     */
    private String userName;
    /**
     * å¯†ç 
     */
    private String password;
    /**
     * ç›˜å·
     */
    private String batchCode;
    /**
     * é‡‡é›†ç‚¹ä½å­—段
     */
    private String pointColumns;
    /**
     * æŽ’序字段名称
     */
    private String orderColumn;
    /**
     * æŽ’序规则
     */
    private String orderRule;
}
src/main/java/com/hwtd/mes/collect/dto/SerialPortDTO.java
¶Ô±ÈÐÂÎļþ
@@ -0,0 +1,66 @@
package com.hwtd.mes.collect.dto;
import lombok.Data;
/**
 * ä¸²å£å‚数,由接口传入(不传的字段使用默认值)。
 * <p>
 * é‡‡é›†æŽ¥å£ {@code /lims/getFile?fileExtension=.serialPort} å’Œ {@code /lims/openSerialPort} éƒ½æ”¯æŒè¿™äº›å‚数,
 * ä¾‹å¦‚:{@code /lims/openSerialPort?listenName=COM7&baudRate=19200&dataBits=8&stopBits=1&parity=NONE&endMark=*&maxCount=6}
 */
@Data
public class SerialPortDTO {
    /** é»˜è®¤æ³¢ç‰¹çއ */
    public static final int DEFAULT_BAUD_RATE = 9600;
    /** é»˜è®¤æ•°æ®ä½ */
    public static final int DEFAULT_DATA_BITS = 8;
    /** é»˜è®¤åœæ­¢ä½ */
    public static final int DEFAULT_STOP_BITS = 1;
    /** é»˜è®¤æ ¡éªŒä½ */
    public static final String DEFAULT_PARITY = "NONE";
    /** é»˜è®¤æµæŽ§ */
    public static final String DEFAULT_FLOW_CONTROL = "NONE";
    /** é»˜è®¤è¯»å–超时(毫秒) */
    public static final int DEFAULT_READ_TIMEOUT = 1000;
    /** é»˜è®¤ç»“束标志 */
    public static final String DEFAULT_END_MARK = "*";
    /** é»˜è®¤å•轮采集条数上限 */
    public static final int DEFAULT_MAX_COUNT = 1;
    /** é»˜è®¤å¸§ç©ºé—²è¡¥å¸§æ—¶é—´ï¼ˆæ¯«ç§’) */
    public static final int DEFAULT_FRAME_IDLE_MILLIS = 200;
    /** ä¸²å£åç§°ï¼Œå¦‚ COM12;必传,不传会返回“串口名称不能为空” */
    private String serialPortName;
    /** æ³¢ç‰¹çŽ‡ï¼Œé»˜è®¤ 9600 */
    private Integer baudRate = DEFAULT_BAUD_RATE;
    /** æ•°æ®ä½ 5~8,默认 8 */
    private Integer dataBits = DEFAULT_DATA_BITS;
    /** åœæ­¢ä½ 1、2、3(1.5 ä½),默认 1 */
    private Integer stopBits = DEFAULT_STOP_BITS;
    /** æ ¡éªŒä½ï¼šNONE/ODD/EVEN/MARK/SPACE(不区分大小写),也兼容 0~4 çš„æ•°å­—,默认 NONE */
    private String parity = DEFAULT_PARITY;
    /** æµæŽ§ï¼šNONE/RTS_CTS/XON_XOFF(不区分大小写),也兼容数字常量,默认 NONE */
    private String flowControl = DEFAULT_FLOW_CONTROL;
    /** è¯»å–超时(毫秒),默认 1000 */
    private Integer readTimeout = DEFAULT_READ_TIMEOUT;
    /** ä¸€å¸§æ•°æ®çš„结束标志,默认 *;特殊字符可用 URL ç¼–码传入,如换行 %0A */
    private String endMark = DEFAULT_END_MARK;
    /** æ•°æ®ç¼“存队列长度,默认 6;队列满时插入最新一条、移除最旧一条,接口最多返回 maxCount æ¡æ•°æ® */
    private Integer maxCount = DEFAULT_MAX_COUNT;
    /**
     * å¸§ç©ºé—²è¡¥å¸§æ—¶é—´ï¼ˆæ¯«ç§’),默认 200:
     * è®¾å¤‡å‘完一帧后超过该时间没有新数据,就把接收缓冲区里剩下的数据当作一帧取出来,
     * é¿å…è®¾å¤‡æœ€åŽä¸€å¸§æ²¡æœ‰ç»“束标志时被漏掉;设为 0 è¡¨ç¤ºå…³é—­ç©ºé—²è¡¥å¸§ï¼ŒåªæŒ‰ç»“束标志切帧。
     */
    private Integer frameIdleMillis = DEFAULT_FRAME_IDLE_MILLIS;
}
src/main/java/com/hwtd/mes/collect/handler/GlobalExceptionHandler.java
ÎļþÃû´Ó src/main/java/com/chinaztt/mes/docx/handler/GlobalExceptionHandler.java ÐÞ¸Ä
@@ -1,6 +1,6 @@
package com.chinaztt.mes.docx.handler;
package com.hwtd.mes.collect.handler;
import com.chinaztt.mes.docx.util.Result;
import com.hwtd.mes.collect.util.Result;
import org.springframework.web.bind.annotation.ExceptionHandler;
import org.springframework.web.bind.annotation.RestControllerAdvice;
src/main/java/com/hwtd/mes/collect/handler/SerialPortListener.java
¶Ô±ÈÐÂÎļþ
@@ -0,0 +1,644 @@
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;
/**
 * ä¸²å£ç›‘听(手动模式)
 * <p>
 * ä¸ŽåŽŸæ¥â€œé¡¹ç›®å¯åŠ¨å°±å¼€å§‹ç›‘å¬â€çš„æ–¹å¼ä¸åŒï¼ŒçŽ°åœ¨æ”¹ä¸ºæ‰‹åŠ¨ç›‘å¬ï¼š
 * <ul>
 *     <li>调用采集接口(/lims/getFile ä¸” fileExtension=.serialPort)时开启监听,串口没有数据时接口返回空集合,由前端轮询取数;</li>
 *     <li>也可以调用 /lims/openSerialPort æ‰‹åŠ¨å¼€å¯ç›‘å¬ï¼›</li>
 *     <li>采集结束后调用 /lims/closeSerialPort æ‰‹åŠ¨å…³é—­ç›‘å¬ï¼ˆå¹‚ç­‰ï¼Œæœªå¼€å¯æ—¶è°ƒç”¨ä¹Ÿä¸ä¼šæŠ¥é”™ï¼‰ï¼›</li>
 *     <li>项目停止时自动关闭监听。</li>
 * </ul>
 * ä¸²å£å‚数(串口名、波特率、数据位、停止位、校验位、流控、读取超时、结束标志、单轮条数上限)全部由接口传入,
 * ä¸ä¼ æ—¶ä½¿ç”¨ {@link SerialPortDTO} é‡Œçš„默认值;参数和当前生效的参数不一致时会自动按新参数重新打开串口,
 * å‚数一致则复用已经打开的串口,不会重复打开。
 * <p>
 * å¹¶å‘说明:
 * <ul>
 *     <li>开启/关闭串口统一由 {@link #portLock} ä¿è¯äº’斥,重复调用只会真正打开一次串口,
 *         å¤šçº¿ç¨‹å¹¶å‘调用采集接口也不会重复打开、重复注册监听器;</li>
 *     <li>采集到的数据放在同步集合中,多线程同时取数时每条数据只会被一个线程取走,不会重复、不会丢失;</li>
 *     <li>监听状态使用 volatile æ ‡è®°ï¼Œå…³é—­ä¸²å£åŽå›žè°ƒçº¿ç¨‹ä¸ä¼šå†å¤„理数据。</li>
 * </ul>
 * æ•°æ®ç¼“存说明:采集到的数据按队列缓存,插入最新的一条,超出 maxCount æ¡ï¼ˆé»˜è®¤ 6 æ¡ï¼‰æ—¶ç§»é™¤æœ€æ—§çš„一条,
 * å› æ­¤æŽ¥å£æ¯æ¬¡æœ€å¤šè¿”回 maxCount æ¡æœ€æ–°çš„æ•°æ®ï¼ˆæŒ‰æŽ¥æ”¶å…ˆåŽé¡ºåºï¼‰ï¼›æ•°æ®è¢«å–走后队列清空,重新开始缓存。
 * <p>
 * å¸§åˆ‡åˆ†è¯´æ˜Žï¼šé»˜è®¤æŒ‰ç»“束标志 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<String> dataList = Collections.synchronizedList(new ArrayList<>());
    /**
     * å¼€å¯ä¸²å£ç›‘听(幂等操作)
     * <p>
     * å¤šçº¿ç¨‹åŒæ—¶è°ƒç”¨æ—¶ï¼šç¬¬ä¸€ä¸ªçº¿ç¨‹çœŸæ­£æ‰“开串口并注册监听器,其它线程直接复用,不会重复打开串口。
     * ä¸²å£å‚数与当前生效参数不一致时会重新打开串口;串口被拔出导致异常断开时,下次调用会自动重新打开。
     *
     * @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<String> drainData() {
        synchronized (dataList) {
            if (dataList.isEmpty()) {
                return new ArrayList<>();
            }
            List<String> 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<String> abandoned = drainData();
        if (!abandoned.isEmpty()) {
            log.warn("串口 {} å¼€å¯æ–°ä¸€è½®ç›‘听,丢弃上一轮未被取走的 {} æ¡æ•°æ®ï¼š{}", getListenName(), abandoned.size(), abandoned);
        }
    }
    /**
     * æž„建串口数据监听器,数据到达时在 jSerialComm çš„监听线程中回调
     * <p>
     * æœ¬è½®ä½¿ç”¨çš„参数(结束标志、条数上限)直接固化在监听器里,避免解析时读到下一轮的参数。
     */
    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<String> 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<String> 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 æ²¡æœ‰æ–°æ•°æ®æ—¶ï¼ŒæŠŠæŽ¥æ”¶ç¼“冲区里剩下的数据当作一帧取出来。
     * <p>
     * å…¸åž‹åœºæ™¯ï¼šè®¾å¤‡æ¯å¸§ä»¥æŸä¸ªå­—符“开头”而不是“结尾”(如 *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);
        }
    }
}
src/main/java/com/hwtd/mes/collect/service/DataCollectionService.java
¶Ô±ÈÐÂÎļþ
@@ -0,0 +1,14 @@
package com.hwtd.mes.collect.service;
import com.hwtd.mes.collect.dto.DatabaseDTO;
import java.util.List;
import java.util.Map;
public interface DataCollectionService {
    List<Map<String,Object>> getAccessData(DatabaseDTO databaseDTO);
    List<Map<String,Object>> getPostgreSqlData(DatabaseDTO databaseDTO);
}
src/main/java/com/hwtd/mes/collect/service/impl/DataCollectionServiceImpl.java
¶Ô±ÈÐÂÎļþ
@@ -0,0 +1,159 @@
package com.hwtd.mes.collect.service.impl;
import com.hwtd.mes.collect.dto.DatabaseDTO;
import com.hwtd.mes.collect.service.DataCollectionService;
import com.zaxxer.hikari.HikariConfig;
import com.zaxxer.hikari.HikariDataSource;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.springframework.stereotype.Service;
import java.sql.*;
import java.util.*;
@Service
@Slf4j
public class DataCollectionServiceImpl implements DataCollectionService {
    /**
     * å¤„理mdb数据库排除字段类型
     */
    private final static List<String> MDB_EXCLUDE_TYPES = Arrays.asList("java.sql.Blob");
    @Override
    public List<Map<String, Object>> getAccessData(DatabaseDTO databaseDTO) {
        List<Map<String,Object>> list = new ArrayList<>();
        try{
            Properties prop = new Properties();
            //设置编码
            prop.put("charSet", "UTF-8");
            prop.put("user",  StringUtils.isNotBlank(databaseDTO.getUserName())?databaseDTO.getUserName():"");
            prop.put("password", StringUtils.isNotBlank(databaseDTO.getPassword())?databaseDTO.getPassword():"");
            //数据地址
            String dbUrl = "jdbc:ucanaccess://" + databaseDTO.getFilePath();
            //引入驱动
            Class.forName("net.ucanaccess.jdbc.UcanaccessDriver").newInstance();
            Connection conn = null;
            PreparedStatement preparedStatement = null;
            ResultSet rs = null;
            //连接数据库资源
            conn = DriverManager.getConnection(dbUrl, prop);
            try {
                //遍历获取多张表数据
                String s = "SELECT "+databaseDTO.getPointColumns()+" FROM " + databaseDTO.getTableName() + " WHERE 1=1";
                if(StringUtils.isNotBlank(databaseDTO.getMainColumn()) && StringUtils.isNotBlank(databaseDTO.getBatchCode())){
                    s+=" and " + databaseDTO.getMainColumn() + " = '" + databaseDTO.getBatchCode()+ "'";
                }
                if(StringUtils.isNotBlank(databaseDTO.getOrderColumn())){
                    String orderRule = StringUtils.isNotBlank(databaseDTO.getOrderRule())?databaseDTO.getOrderRule():"ASC";
                    s+=" ORDER BY " + databaseDTO.getOrderColumn() + " " + orderRule;
                }
                preparedStatement = conn.prepareStatement(s);
                rs = preparedStatement.executeQuery();
                ResultSetMetaData data = rs.getMetaData();
                while (rs.next()) {
                    Map<String, Object> map = new HashMap<>();
                    for (int i = 1; i <= data.getColumnCount(); i++) {
                        //列名
                        String columnName = data.getColumnName(i);
                        //列字段类型
                        String columnClassName = data.getColumnClassName(i);
                        Object columnValue = null;
                        if(!MDB_EXCLUDE_TYPES.contains(columnClassName)){
                            columnValue = rs.getObject(i);
                        }
                        map.put(columnName, columnValue);
                    }
                    list.add(map);
                }
            } catch (Exception e) {
                e.printStackTrace();
            } finally {
                closeA1l(conn, preparedStatement, rs);
            }
            return list;
        }catch (Exception e){
            throw new RuntimeException("Access数据库采集异常:"+e.getMessage());
        }
    }
    @Override
    public List<Map<String, Object>> getPostgreSqlData(DatabaseDTO databaseDTO) {
        List<Map<String, Object>> dataList = new ArrayList<>();
        try{
            String dbName = databaseDTO.getDatabaseName();
            String user = databaseDTO.getUserName();
            String password = databaseDTO.getPassword();
            // ä»Ž GetFileDto èŽ·å–æ•°æ®è¡¨åï¼Œå¯¹åº”ã€æ•°æ®åº“è¡¨åã€‘å­—æ®µ
            String table = databaseDTO.getTableName();
            // æ£€æŸ¥æ•°æ®åº“名和表名是否为空
            if (dbName == null || dbName.isEmpty() || table == null || table.isEmpty()) {
                throw new RuntimeException("数据库名或表名不能为空");
            }
            // æ•°æ®åº“连接信息
            String url = String.format("jdbc:postgresql://%s:%s/%s",databaseDTO.getIpAddress(),databaseDTO.getServerPort(),dbName);
            Connection connection = null;
            PreparedStatement preparedStatement = null;
            ResultSet resultSet = null;
            HikariConfig config = new HikariConfig();
            config.setJdbcUrl(url);
            config.setUsername(user);
            config.setPassword(password);
            config.setMaximumPoolSize(10);
            config.setConnectionTimeout(30000);
            HikariDataSource ds = new HikariDataSource(config);
            try {
                // å»ºç«‹è¿žæŽ¥
                connection = ds.getConnection();
                // æž„建基础 SQL
                String sql = "SELECT "+databaseDTO.getPointColumns()+" FROM "+table+" WHERE 1=1";
                if(StringUtils.isNotBlank(databaseDTO.getMainColumn()) && StringUtils.isNotBlank(databaseDTO.getBatchCode())){
                    sql+=" AND (" + databaseDTO.getMainColumn() + " = TRIM('" + databaseDTO.getBatchCode()+ "')";
                }
                if(StringUtils.isNotBlank(databaseDTO.getOrderColumn())){
                    String orderRule = StringUtils.isNotBlank(databaseDTO.getOrderRule())?databaseDTO.getOrderRule():"ASC";
                    sql+=" ORDER BY " + databaseDTO.getOrderColumn() + " " + orderRule;
                }
                // åˆ›å»º PreparedStatement å¯¹è±¡æ‰§è¡Œ SQL
                preparedStatement = connection.prepareStatement(sql);
                resultSet = preparedStatement.executeQuery();
                ResultSetMetaData metaData = resultSet.getMetaData();
                int columnCount = metaData.getColumnCount();
                // éåŽ†ç»“æžœé›†èŽ·å–æ•°æ®
                while (resultSet.next()) {
                    Map<String, Object> rowData = new HashMap<>();
                    for (int i = 1; i <= columnCount; i++) {
                        String columnName = metaData.getColumnName(i);
                        rowData.put(columnName, resultSet.getObject(i));
                    }
                    dataList.add(rowData);
                }
            } catch (Exception e) {
                e.printStackTrace();
            } finally {
                closeA1l(connection, preparedStatement, resultSet);
            }
            return dataList;
        }catch (Exception e){
            throw new RuntimeException("PostgreSql数据库采集异常:"+e.getMessage());
        }
    }
    private static void closeA1l(Connection conn, PreparedStatement preparedStatement, ResultSet rs) {
        try {
            if (null != rs) {
                rs.close();
            }
            if (null != preparedStatement) {
                preparedStatement.close();
            }
            if (null != conn) {
                conn.close();
            }
        } catch (Exception ignore) {
        }
    }
}
src/main/java/com/hwtd/mes/collect/util/Result.java
ÎļþÃû´Ó src/main/java/com/chinaztt/mes/docx/util/Result.java ÐÞ¸Ä
@@ -3,9 +3,9 @@
// (powered by FernFlower decompiler)
//
package com.chinaztt.mes.docx.util;
package com.hwtd.mes.collect.util;
import com.chinaztt.mes.docx.constant.CommonConstants;
import com.hwtd.mes.collect.constant.CommonConstants;
import java.io.Serializable;
src/main/resources/META-INF/MANIFEST.MF
@@ -1,3 +1,3 @@
Manifest-Version: 1.0
Main-Class: com.chinaztt.mes.docx.DataAcquisitionApplication
Main-Class: com.hwtd.mes.collect.DataAcquisitionApplication
src/main/resources/application.yml
@@ -3,7 +3,3 @@
logging:
  file-location: D:\lims-acquistion\logs
serialPort:
  enable: false # æ˜¯å¦å¼€å¯ä¸²å£ç›‘听
  listenName: COM4 # ç›‘听的串口名称
src/test/java/com/hwtd/mes/collect/DataAcquisitionApplicationTests.java
ÎļþÃû´Ó src/test/java/com/chinaztt/mes/docx/DataAcquisitionApplicationTests.java ÐÞ¸Ä
@@ -1,4 +1,4 @@
package com.chinaztt.mes.docx;
package com.hwtd.mes.collect;
import org.apache.commons.lang3.ObjectUtils;
import org.junit.jupiter.api.Test;