FileParse2ServiceImpl.java 8.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182
  1. package com.uas.eis.serviceImpl;
  2. import com.uas.eis.core.config.SpObserver;
  3. import com.uas.eis.core.support.HeaderRowHandler;
  4. import com.uas.eis.core.support.SaxDataBatchHandler;
  5. import com.uas.eis.dao.BaseDao;
  6. import com.uas.eis.entity.DataChipTestLog;
  7. import com.uas.eis.utils.FTPUtil;
  8. import org.apache.commons.net.ftp.FTPClient;
  9. import org.apache.poi.openxml4j.opc.OPCPackage;
  10. import org.apache.poi.xssf.eventusermodel.ReadOnlySharedStringsTable;
  11. import org.apache.poi.xssf.eventusermodel.XSSFReader;
  12. import org.apache.poi.xssf.model.SharedStrings;
  13. import org.apache.poi.xssf.model.StylesTable;
  14. import org.slf4j.Logger;
  15. import org.slf4j.LoggerFactory;
  16. import org.springframework.beans.factory.annotation.Autowired;
  17. import org.springframework.scheduling.annotation.Async;
  18. import org.springframework.stereotype.Service;
  19. import org.xml.sax.InputSource;
  20. import org.xml.sax.XMLReader;
  21. import org.xml.sax.helpers.XMLReaderFactory;
  22. import java.io.BufferedInputStream;
  23. import java.io.IOException;
  24. import java.io.InputStream;
  25. import java.util.ArrayList;
  26. import java.util.List;
  27. import java.util.concurrent.CountDownLatch;
  28. @Service
  29. @Async("taskExecutor")
  30. public class FileParse2ServiceImpl {
  31. private static final Logger logger = LoggerFactory.getLogger(FileParse2ServiceImpl.class);
  32. // 常量配置
  33. private static final int BATCH_ROW_SIZE = 1000;
  34. private static final String DB_LINK = "@SHENAI_EDCDATA";
  35. private static final int LABEL_ROW_INDEX = 13;
  36. private static final int DATA_START_INDEX = 23;
  37. // 外部注入
  38. @Autowired
  39. private BaseDao baseDao;
  40. public void EDCDataDeal(DataChipTestLog dataChipTestLog, CountDownLatch countDownLatch) {
  41. String waferId = dataChipTestLog.getWafer_id();
  42. String lotId = dataChipTestLog.getLot_id();
  43. Integer dbId = dataChipTestLog.getDbid_();
  44. String tableName = "DATA$" + lotId.replace("-", "_");
  45. FTPClient ftpClient = null;
  46. InputStream ftpInputStream = null;
  47. OPCPackage opcPackage = null;
  48. try {
  49. logger.info("WaferId:{} Data文件解析开始", waferId);
  50. String filePath = extractFilePath(dataChipTestLog.getData_path());
  51. ftpClient = FTPUtil.getFTPClient();
  52. boolean fileExists = FTPUtil.exists(ftpClient, filePath);
  53. logger.info("文件路径:{}, FTP文件存在:{}", filePath, fileExists);
  54. if (ftpClient == null) {
  55. logger.warn("FTP客户端为空,跳过 waferId={}", waferId);
  56. countDownLatch.countDown();
  57. return;
  58. }
  59. if(!fileExists) {
  60. logger.warn("文件不存在,跳过 waferId={}", waferId);
  61. String errSql = String.format("UPDATE DATACENTER$MESFILE SET DEALSTATE_=-100 , DEALTIME_=SYSDATE WHERE DBID_=%s", dbId);
  62. baseDao.execute(errSql);
  63. countDownLatch.countDown();
  64. return;
  65. }
  66. // 打开FTP文件流 & OPC压缩包
  67. ftpInputStream = new BufferedInputStream(FTPUtil.readFileAsStream(ftpClient, filePath));
  68. opcPackage = OPCPackage.open(ftpInputStream);
  69. XSSFReader xssfReader = new XSSFReader(opcPackage);
  70. // 开启只读共享字符串表,降低内存占用
  71. SharedStrings sharedStrings = xssfReader.getSharedStringsTable();
  72. StylesTable stylesTable = xssfReader.getStylesTable();
  73. XMLReader xmlReader = XMLReaderFactory.createXMLReader();
  74. // 1. 读取表头行,生成ITEM1,ITEM2...动态列
  75. List<String> headerList = new ArrayList<>();
  76. readHeaderRow(xssfReader, xmlReader, sharedStrings, stylesTable, LABEL_ROW_INDEX, headerList);
  77. List<String> itemColumnList = new ArrayList<>();
  78. for (int i = 0; i < headerList.size()+1; i++) {
  79. itemColumnList.add("ITEM" + (i + 1));
  80. }
  81. String dynamicCols = String.join(",", itemColumnList);
  82. // 2. 先删除当前wafer历史数据
  83. String deleteSql = String.format("DELETE FROM %s%s WHERE WAFER_ID='%s'", tableName, DB_LINK, waferId);
  84. baseDao.execute(deleteSql);
  85. // 3. SAX流式读取数据分批入库,传入JdbcTemplate
  86. SaxDataBatchHandler dataHandler = new SaxDataBatchHandler(
  87. baseDao.getJdbcTemplate(),
  88. waferId,
  89. tableName,
  90. DB_LINK,
  91. dynamicCols,
  92. headerList,
  93. DATA_START_INDEX,
  94. BATCH_ROW_SIZE,
  95. sharedStrings
  96. );
  97. dataHandler.parseSheet(xssfReader, xmlReader, stylesTable);
  98. // 4. 后置业务SQL:更新文件状态、归档、日志
  99. StringBuilder postSqlBlock = new StringBuilder(1024);
  100. postSqlBlock.append("BEGIN ");
  101. postSqlBlock.append("UPDATE DATACENTER$MESFILE SET DEALSTATE_=1,DEALTIME_=SYSDATE WHERE DBID_=").append(dbId).append("; ");
  102. 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_) ");
  103. 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_ ");
  104. postSqlBlock.append("FROM DATACENTER$MESFILE WHERE DBID_=").append(dbId);
  105. postSqlBlock.append(" AND WAFER_ID NOT IN (SELECT WAFER_ID FROM chip$file WHERE WAFER_ID='").append(waferId).append("'); ");
  106. postSqlBlock.append("INSERT INTO chip$data$log (WAFER_ID,ID_) SELECT '").append(waferId).append("',CHIP$DATA$LOG_SEQ.NEXTVAL FROM DUAL; END;");
  107. baseDao.execute(postSqlBlock.toString());
  108. logger.info("WaferId:{} 文件处理完成", waferId);
  109. } catch (Exception e) {
  110. logger.error("芯片号:{} 解析异常{}", waferId, e.getMessage());
  111. // 异常标记失败状态
  112. String errSql = String.format("UPDATE DATACENTER$MESFILE SET DEALSTATE_=-99 , DEALTIME_=SYSDATE WHERE DBID_=%s", dbId);
  113. baseDao.execute(errSql);
  114. } finally {
  115. // 严格顺序释放全部资源,杜绝内存/句柄泄漏
  116. try {
  117. if (opcPackage != null) {
  118. opcPackage.close();
  119. }
  120. } catch (Exception ex) {
  121. logger.error("关闭OPCPackage失败", ex);
  122. }
  123. try {
  124. if (ftpInputStream != null) {
  125. ftpInputStream.close();
  126. }
  127. } catch (IOException ex) {
  128. logger.error("关闭FTP文件输入流失败", ex);
  129. }
  130. try {
  131. if (ftpClient != null) {
  132. ftpClient.completePendingCommand();
  133. ftpClient.logout();
  134. ftpClient.disconnect();
  135. }
  136. } catch (IOException ex) {
  137. logger.error("释放FTP连接失败", ex);
  138. }
  139. // 线程计数器放行
  140. countDownLatch.countDown();
  141. }
  142. }
  143. /**
  144. * 读取指定行作为表头
  145. */
  146. private void readHeaderRow(XSSFReader reader, XMLReader xmlReader, SharedStrings sst,
  147. StylesTable styles, int targetRow, List<String> outHeader) throws Exception {
  148. InputSource sheetSource = new InputSource(reader.getSheetsData().next());
  149. HeaderRowHandler headerHandler = new HeaderRowHandler(targetRow, sst, outHeader);
  150. xmlReader.setContentHandler(headerHandler);
  151. xmlReader.parse(sheetSource);
  152. }
  153. private String extractFilePath(String ftpUrl) {
  154. if (ftpUrl == null || ftpUrl.isEmpty()) {
  155. return null;
  156. }
  157. // 查找 ftp:// 的位置
  158. int protocolIndex = ftpUrl.indexOf("ftp://");
  159. if (protocolIndex == -1) {
  160. // 不是FTP URL,直接返回原路径
  161. return ftpUrl;
  162. }
  163. int pathStartIndex = ftpUrl.indexOf("/", protocolIndex + 6);
  164. if (pathStartIndex == -1) {
  165. // 没有路径部分,返回根目录
  166. return "/";
  167. }
  168. return "/home/mes"+ftpUrl.substring(pathStartIndex).trim();
  169. }
  170. }