| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182 |
- package com.uas.eis.serviceImpl;
- import com.uas.eis.core.config.SpObserver;
- import com.uas.eis.core.support.HeaderRowHandler;
- import com.uas.eis.core.support.SaxDataBatchHandler;
- import com.uas.eis.dao.BaseDao;
- import com.uas.eis.entity.DataChipTestLog;
- import com.uas.eis.utils.FTPUtil;
- import org.apache.commons.net.ftp.FTPClient;
- import org.apache.poi.openxml4j.opc.OPCPackage;
- import org.apache.poi.xssf.eventusermodel.ReadOnlySharedStringsTable;
- import org.apache.poi.xssf.eventusermodel.XSSFReader;
- import org.apache.poi.xssf.model.SharedStrings;
- import org.apache.poi.xssf.model.StylesTable;
- import org.slf4j.Logger;
- import org.slf4j.LoggerFactory;
- import org.springframework.beans.factory.annotation.Autowired;
- import org.springframework.scheduling.annotation.Async;
- import org.springframework.stereotype.Service;
- import org.xml.sax.InputSource;
- import org.xml.sax.XMLReader;
- import org.xml.sax.helpers.XMLReaderFactory;
- import java.io.BufferedInputStream;
- import java.io.IOException;
- import java.io.InputStream;
- import java.util.ArrayList;
- import java.util.List;
- import java.util.concurrent.CountDownLatch;
- @Service
- @Async("taskExecutor")
- public class FileParse2ServiceImpl {
- private static final Logger logger = LoggerFactory.getLogger(FileParse2ServiceImpl.class);
- // 常量配置
- private static final int BATCH_ROW_SIZE = 1000;
- private static final String DB_LINK = "@SHENAI_EDCDATA";
- private static final int LABEL_ROW_INDEX = 13;
- private static final int DATA_START_INDEX = 23;
- // 外部注入
- @Autowired
- private BaseDao baseDao;
- public void EDCDataDeal(DataChipTestLog dataChipTestLog, CountDownLatch countDownLatch) {
- String waferId = dataChipTestLog.getWafer_id();
- String lotId = dataChipTestLog.getLot_id();
- Integer dbId = dataChipTestLog.getDbid_();
- String tableName = "DATA$" + lotId.replace("-", "_");
- FTPClient ftpClient = null;
- InputStream ftpInputStream = null;
- OPCPackage opcPackage = null;
- try {
- logger.info("WaferId:{} Data文件解析开始", waferId);
- String filePath = extractFilePath(dataChipTestLog.getData_path());
- ftpClient = FTPUtil.getFTPClient();
- boolean fileExists = FTPUtil.exists(ftpClient, filePath);
- logger.info("文件路径:{}, FTP文件存在:{}", filePath, fileExists);
- if (ftpClient == null) {
- logger.warn("FTP客户端为空,跳过 waferId={}", waferId);
- countDownLatch.countDown();
- return;
- }
- if(!fileExists) {
- logger.warn("文件不存在,跳过 waferId={}", waferId);
- String errSql = String.format("UPDATE DATACENTER$MESFILE SET DEALSTATE_=-100 , DEALTIME_=SYSDATE WHERE DBID_=%s", dbId);
- baseDao.execute(errSql);
- countDownLatch.countDown();
- return;
- }
- // 打开FTP文件流 & OPC压缩包
- ftpInputStream = new BufferedInputStream(FTPUtil.readFileAsStream(ftpClient, filePath));
- opcPackage = OPCPackage.open(ftpInputStream);
- XSSFReader xssfReader = new XSSFReader(opcPackage);
- // 开启只读共享字符串表,降低内存占用
- SharedStrings sharedStrings = xssfReader.getSharedStringsTable();
- StylesTable stylesTable = xssfReader.getStylesTable();
- XMLReader xmlReader = XMLReaderFactory.createXMLReader();
- // 1. 读取表头行,生成ITEM1,ITEM2...动态列
- List<String> headerList = new ArrayList<>();
- readHeaderRow(xssfReader, xmlReader, sharedStrings, stylesTable, LABEL_ROW_INDEX, headerList);
- List<String> itemColumnList = new ArrayList<>();
- for (int i = 0; i < headerList.size()+1; i++) {
- itemColumnList.add("ITEM" + (i + 1));
- }
- String dynamicCols = String.join(",", itemColumnList);
- // 2. 先删除当前wafer历史数据
- String deleteSql = String.format("DELETE FROM %s%s WHERE WAFER_ID='%s'", tableName, DB_LINK, waferId);
- baseDao.execute(deleteSql);
- // 3. SAX流式读取数据分批入库,传入JdbcTemplate
- SaxDataBatchHandler dataHandler = new SaxDataBatchHandler(
- baseDao.getJdbcTemplate(),
- waferId,
- tableName,
- DB_LINK,
- dynamicCols,
- headerList,
- DATA_START_INDEX,
- BATCH_ROW_SIZE,
- sharedStrings
- );
- dataHandler.parseSheet(xssfReader, xmlReader, stylesTable);
- // 4. 后置业务SQL:更新文件状态、归档、日志
- StringBuilder postSqlBlock = new StringBuilder(1024);
- postSqlBlock.append("BEGIN ");
- postSqlBlock.append("UPDATE DATACENTER$MESFILE SET DEALSTATE_=1,DEALTIME_=SYSDATE WHERE DBID_=").append(dbId).append("; ");
- postSqlBlock.append("INSERT INTO chip$file (WAFER_ID,LOT_ID,PASS_QTY,FAIL_QTY,PASS_YFIELD,EQUIP_ID,STATUS,COUNTER_PATH,DATA_PATH,TEST_LOT,CREATETIME_,DEALSTATE_,DEALTIME_,DBID_) ");
- postSqlBlock.append("SELECT WAFER_ID,LOT_ID,PASS_QTY,FAIL_QTY,PASS_YFIELD,EQUIP_ID,STATUS,TRIM(COUNTER_PATH),TRIM(DATA_PATH),TEST_LOT,CREATETIME_,DEALSTATE_,DEALTIME_,DBID_ ");
- postSqlBlock.append("FROM DATACENTER$MESFILE WHERE DBID_=").append(dbId);
- postSqlBlock.append(" AND WAFER_ID NOT IN (SELECT WAFER_ID FROM chip$file WHERE WAFER_ID='").append(waferId).append("'); ");
- postSqlBlock.append("INSERT INTO chip$data$log (WAFER_ID,ID_) SELECT '").append(waferId).append("',CHIP$DATA$LOG_SEQ.NEXTVAL FROM DUAL; END;");
- baseDao.execute(postSqlBlock.toString());
- logger.info("WaferId:{} 文件处理完成", waferId);
- } catch (Exception e) {
- logger.error("芯片号:{} 解析异常{}", waferId, e.getMessage());
- // 异常标记失败状态
- String errSql = String.format("UPDATE DATACENTER$MESFILE SET DEALSTATE_=-99 , DEALTIME_=SYSDATE WHERE DBID_=%s", dbId);
- baseDao.execute(errSql);
- } finally {
- // 严格顺序释放全部资源,杜绝内存/句柄泄漏
- try {
- if (opcPackage != null) {
- opcPackage.close();
- }
- } catch (Exception ex) {
- logger.error("关闭OPCPackage失败", ex);
- }
- try {
- if (ftpInputStream != null) {
- ftpInputStream.close();
- }
- } catch (IOException ex) {
- logger.error("关闭FTP文件输入流失败", ex);
- }
- try {
- if (ftpClient != null) {
- ftpClient.completePendingCommand();
- ftpClient.logout();
- ftpClient.disconnect();
- }
- } catch (IOException ex) {
- logger.error("释放FTP连接失败", ex);
- }
- // 线程计数器放行
- countDownLatch.countDown();
- }
- }
- /**
- * 读取指定行作为表头
- */
- private void readHeaderRow(XSSFReader reader, XMLReader xmlReader, SharedStrings sst,
- StylesTable styles, int targetRow, List<String> outHeader) throws Exception {
- InputSource sheetSource = new InputSource(reader.getSheetsData().next());
- HeaderRowHandler headerHandler = new HeaderRowHandler(targetRow, sst, outHeader);
- xmlReader.setContentHandler(headerHandler);
- xmlReader.parse(sheetSource);
- }
- private String extractFilePath(String ftpUrl) {
- if (ftpUrl == null || ftpUrl.isEmpty()) {
- return null;
- }
- // 查找 ftp:// 的位置
- int protocolIndex = ftpUrl.indexOf("ftp://");
- if (protocolIndex == -1) {
- // 不是FTP URL,直接返回原路径
- return ftpUrl;
- }
- int pathStartIndex = ftpUrl.indexOf("/", protocolIndex + 6);
- if (pathStartIndex == -1) {
- // 没有路径部分,返回根目录
- return "/";
- }
- return "/home/mes"+ftpUrl.substring(pathStartIndex).trim();
- }
- }
|