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 headerList = new ArrayList<>(); readHeaderRow(xssfReader, xmlReader, sharedStrings, stylesTable, LABEL_ROW_INDEX, headerList); List 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 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(); } }