DownloadFileFromFtpToTable.java 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441
  1. package com.java110.job.smo;
  2. import com.java110.job.util.FTPClientTemplate;
  3. import org.apache.commons.net.ftp.FTP;
  4. import org.apache.commons.net.ftp.FTPClient;
  5. import org.apache.commons.net.ftp.FTPReply;
  6. import java.io.File;
  7. import java.io.RandomAccessFile;
  8. import java.util.ArrayList;
  9. import java.util.HashMap;
  10. import java.util.List;
  11. import java.util.Map;
  12. import java.util.concurrent.Executors;
  13. import java.util.concurrent.Future;
  14. import java.util.concurrent.ThreadPoolExecutor;
  15. /**
  16. * 从Ftp文件系统下载文件内容存到对应配置的表中
  17. *
  18. *
  19. * @author wuxw7 2016-01-04
  20. *
  21. */
  22. public class DownloadFileFromFtpToTable extends PrvncFtpToFileSystemQuartz {
  23. private static final String ITEM_SPEC_CD_10011 = "10011";// FTP地址
  24. private static final String ITEM_SPEC_CD_10012 = "10012";// FTP端口号
  25. private static final String ITEM_SPEC_CD_10013 = "10013";// FTP账号
  26. private static final String ITEM_SPEC_CD_10014 = "10014";// FTP密码
  27. private static final String ITEM_SPEC_CD_10015 = "10015";// FTP路径
  28. private static final String ITEM_SPEC_CD_10016 = "10016";// 本地路径
  29. private static final String ITEM_SPEC_CD_10007 = "10007";// 文件头
  30. private static final String ITEM_SPEC_CD_10008 = "10008";// 分隔符
  31. private static final String ITEM_SPEC_CD_10009 = "10009";// 总记录数
  32. private static final String ITEM_SPEC_CD_10010 = "10010";// 处理脚本
  33. /**
  34. * 如果运行失败时,需要在ftpItemConfigInfo 中PRE_METHOD_FLAG 值改成E,remark 备注错误原因
  35. */
  36. @Override
  37. protected void process(Map ftpItemConfigInfo) throws Exception {
  38. // TODO Auto-generated method stub
  39. String taskId = ftpItemConfigInfo.get("TASKID").toString();
  40. FTPClientTemplate ftpClientTemplate = null;
  41. // 1.0 读取配置,包括ftp 服务器信息,和tfs相关配置,根据taskId 关联信息
  42. if (!ftpItemConfigInfo.containsKey("FTP_ITEM_ATTRS") || ftpItemConfigInfo.get("FTP_ITEM_ATTRS") == null) {
  43. ftpItemConfigInfo.put("PRE_METHOD_FLAG", "E");
  44. ftpItemConfigInfo.put("remark", "当前ftp任务【" + taskId + "】没有配置属性,在从Ftp文件系统下载文件存到TFS文件系统模板中必须配置属性");
  45. return;
  46. }
  47. List<Map> ftpItemAttrs = (List<Map>) ftpItemConfigInfo.get("FTP_ITEM_ATTRS");
  48. // FTP_ITEM_ATTRS
  49. String ftpIp = null;
  50. int ftpPort = 21;
  51. String ftpUsername = null;
  52. String ftpPassword = null;
  53. String ftpPath = null;
  54. String titleflag = null;
  55. String sign = null;
  56. String linecountflag = null;
  57. String dbsql = null; // 处理脚本
  58. String localPath = "";// 本地文件保存路径
  59. int tnum = ftpItemConfigInfo.get("TNUM") == null ? 1 : Integer.parseInt(ftpItemConfigInfo.get("TNUM").toString());
  60. // 如果没有配置默认获取指定目录下的所有文件
  61. String fileName = ftpItemConfigInfo.get("FILE_NAME") == null ? "" : ftpItemConfigInfo.get("FILE_NAME").toString();
  62. // 文件名如果为空,直接返回再不处理
  63. if ("".equals(fileName)) {
  64. ftpItemConfigInfo.put("PRE_METHOD_FLAG", "E");
  65. ftpItemConfigInfo.put("remark", "当前ftp任务【" + taskId + "】没有配置文件名,在从Ftp文件系统下载文件存到TFS文件系统模板中必须配置文件名");
  66. return;
  67. }
  68. // 获取属性表中的数据
  69. for (Map ftpItemAttr : ftpItemAttrs) {
  70. if (ftpItemAttr.containsKey("ITEM_SPEC_ID") && ITEM_SPEC_CD_10011.equals(ftpItemAttr.get("ITEM_SPEC_ID").toString())) {
  71. ftpIp = ftpItemAttr.get("VALUE") == null ? "" : ftpItemAttr.get("VALUE").toString();
  72. } else if (ftpItemAttr.containsKey("ITEM_SPEC_ID") && ITEM_SPEC_CD_10012.equals(ftpItemAttr.get("ITEM_SPEC_ID").toString())) {
  73. ftpPort = ftpItemAttr.get("VALUE") == null ? 21 : Integer.parseInt(ftpItemAttr.get("VALUE").toString());
  74. } else if (ftpItemAttr.containsKey("ITEM_SPEC_ID") && ITEM_SPEC_CD_10013.equals(ftpItemAttr.get("ITEM_SPEC_ID").toString())) {
  75. ftpUsername = ftpItemAttr.get("VALUE") == null ? "" : ftpItemAttr.get("VALUE").toString();
  76. } else if (ftpItemAttr.containsKey("ITEM_SPEC_ID") && ITEM_SPEC_CD_10014.equals(ftpItemAttr.get("ITEM_SPEC_ID").toString())) {
  77. ftpPassword = ftpItemAttr.get("VALUE") == null ? "" : ftpItemAttr.get("VALUE").toString();
  78. } else if (ftpItemAttr.containsKey("ITEM_SPEC_ID") && ITEM_SPEC_CD_10015.equals(ftpItemAttr.get("ITEM_SPEC_ID").toString())) {
  79. ftpPath = ftpItemAttr.get("VALUE") == null ? "" : ftpItemAttr.get("VALUE").toString();
  80. } else if (ftpItemAttr.containsKey("ITEM_SPEC_ID") && ITEM_SPEC_CD_10007.equals(ftpItemAttr.get("ITEM_SPEC_ID").toString())) {
  81. titleflag = ftpItemAttr.get("VALUE") == null ? "" : ftpItemAttr.get("VALUE").toString();
  82. } else if (ftpItemAttr.containsKey("ITEM_SPEC_ID") && ITEM_SPEC_CD_10008.equals(ftpItemAttr.get("ITEM_SPEC_ID").toString())) {
  83. sign = ftpItemAttr.get("VALUE") == null ? "" : ftpItemAttr.get("VALUE").toString();
  84. } else if (ftpItemAttr.containsKey("ITEM_SPEC_ID") && ITEM_SPEC_CD_10009.equals(ftpItemAttr.get("ITEM_SPEC_ID").toString())) {
  85. linecountflag = ftpItemAttr.get("VALUE") == null ? "" : ftpItemAttr.get("VALUE").toString();
  86. } else if (ftpItemAttr.containsKey("ITEM_SPEC_ID") && ITEM_SPEC_CD_10010.equals(ftpItemAttr.get("ITEM_SPEC_ID").toString())) {
  87. dbsql = ftpItemAttr.get("VALUE") == null ? "" : ftpItemAttr.get("VALUE").toString();
  88. } else if (ftpItemAttr.containsKey("ITEM_SPEC_ID") && ITEM_SPEC_CD_10016.equals(ftpItemAttr.get("ITEM_SPEC_ID").toString())) {
  89. localPath = ftpItemAttr.get("VALUE") == null ? "" : ftpItemAttr.get("VALUE").toString();
  90. }
  91. }
  92. // 将 contentSpiltChar totalCount dealSql fileTop属性回写到
  93. ftpItemConfigInfo.put("sign", sign);
  94. ftpItemConfigInfo.put("linecountflag", linecountflag);
  95. ftpItemConfigInfo.put("dbsql", dbsql);
  96. ftpItemConfigInfo.put("titleflag", titleflag);
  97. ftpItemConfigInfo.put("localPath", localPath);
  98. ftpItemConfigInfo.put("port", ftpPort);
  99. ftpItemConfigInfo.put("username", ftpUsername);
  100. ftpItemConfigInfo.put("pwd", ftpPassword);
  101. ftpItemConfigInfo.put("ip", ftpIp);
  102. // 初始化FTPClientTemplate
  103. ftpClientTemplate = new FTPClientTemplate(ftpIp, ftpPort, ftpUsername, ftpPassword, ftpPath, null, 0, tnum, 100);
  104. if (!ftpPath.endsWith("/")) {
  105. ftpPath += "/";
  106. }
  107. // 处理文件名
  108. List<String> fileNames = dealFileName(fileName);
  109. // **** 则不支持select出来文件名是list 的情况,如果是list 默认处理第一个
  110. // 校验文件名是中是否存在****(4个),表示通配符,如果不存在就是确定唯一文件名
  111. fileName = fileNames.get(0);
  112. if (fileName != null && fileName.contains("****")) {
  113. String[] reFileNames = ftpClientTemplate.listNames(ftpPath + fileName, true);
  114. List<String> fNames = new ArrayList<String>();
  115. for (int reFileNamesIndex = 0; reFileNamesIndex < reFileNames.length; reFileNamesIndex++) {
  116. if (reFileNames[reFileNamesIndex] != null && !"".equals(reFileNames[reFileNamesIndex])) {
  117. fNames.add(reFileNames[reFileNamesIndex].replace(ftpPath, ""));
  118. }
  119. }
  120. // 多文件处理
  121. anyFilesDownload(ftpItemConfigInfo, taskId, ftpClientTemplate, ftpPath, fNames);
  122. return;
  123. }
  124. // 单文件支持文件名为list,如果是list,则处理list所有文件
  125. // 单文件处理
  126. anyFilesDownload(ftpItemConfigInfo, taskId, ftpClientTemplate, ftpPath, fileNames);
  127. }
  128. /**
  129. * 文件名是中是存在****,多文件下载处理
  130. *
  131. * @param ftpItemConfigInfo
  132. * @param taskId
  133. * @param ftpClientTemplate
  134. * @param ftpPath
  135. * @param
  136. * @param
  137. * @param
  138. * @throws Exception
  139. */
  140. private void anyFilesDownload(Map ftpItemConfigInfo, String taskId, FTPClientTemplate ftpClientTemplate, String ftpPath, List<String> fileNames) throws Exception {
  141. // 这种需要列出所有文件,根据名称匹配
  142. String param = "";
  143. // 获取FTP上的文件
  144. List<Map> needDownloadFiles = null;
  145. String downLoadFailFileNames = "";// 下载失败的fileName
  146. String downLoadSuccessFileNames = "";// 下载失败的fileName
  147. for (String fileName : fileNames) {
  148. if (fileName.length() > 0) {
  149. // 查询数据库,那写文件还没有下载
  150. Map logInfo = new HashMap();
  151. logInfo.put("fileNames", fileName);
  152. logInfo.put("taskId", taskId);
  153. needDownloadFiles = this.getPrvncFtpFileDAO().queryFileNamesWithOutFtpLog(logInfo);
  154. }
  155. if (needDownloadFiles == null || needDownloadFiles.size() < 1) {
  156. continue;
  157. }
  158. ftpItemConfigInfo.put("newFileName", fileName);
  159. // 保存文件至table
  160. Map resultInfo = downLoadFileToTable(ftpPath + fileName, ftpClientTemplate, ftpItemConfigInfo);
  161. Map remoteFileInfo = new HashMap();
  162. if (resultInfo.containsKey("SAVE_FILE_FLAG") && "S".equals(resultInfo.get("SAVE_FILE_FLAG"))) {
  163. param += (fileName + "@@");
  164. // 将下载成功的文件名需要保存至表中,防止以后重复下载
  165. remoteFileInfo.put("taskId", taskId);
  166. remoteFileInfo.put("fileName", fileName);
  167. saveDownLoadSuccessFile(remoteFileInfo);
  168. // 记录成功时的文件名
  169. downLoadSuccessFileNames += (fileName + ",");
  170. } else {
  171. // 记录下载失败的文件名
  172. downLoadFailFileNames += (fileName + ",");
  173. }
  174. // }
  175. }
  176. // 做这个校验主要为了,如果一个都没有成功,就不去调事后过程
  177. if ("".equals(param)) {
  178. ftpItemConfigInfo.put("PRE_METHOD_FLAG", "E");
  179. ftpItemConfigInfo.put("remark", "当前ftp任务【" + taskId + "】没有可下载的文件,或下载文件时失败");
  180. return;
  181. }
  182. ftpItemConfigInfo.put("PRE_METHOD_FLAG", "S");
  183. ftpItemConfigInfo.put("remark", "当前ftp任务【" + taskId + "】下载成功{" + downLoadSuccessFileNames + "},失败的{" + downLoadFailFileNames + "}");
  184. ftpItemConfigInfo.put("param", param);
  185. return;
  186. }
  187. /**
  188. * 下载并保存文件至表中。
  189. *
  190. * @param remoteFileInfo
  191. */
  192. private void saveDownLoadSuccessFile(Map remoteFileInfo) {
  193. // TODO Auto-generated method stub
  194. int addDownloadFlag = this.getPrvncFtpFileDAO().addDownloadFileName(remoteFileInfo);
  195. if (addDownloadFlag < 1) {
  196. logger.error("---【DownloadFileFromFtpToTFS.saveDownLoadSuccessFile】保存下载文件名失败", remoteFileInfo);
  197. }
  198. }
  199. /**
  200. * 下载文件并且存至tfs文件系统
  201. *
  202. * @param remoteFileNameTmp
  203. * @param
  204. * @return
  205. */
  206. private Map downLoadFileToTable(String remoteFileNameTmp, FTPClientTemplate ftpClientTemplate, Map ftpItemConfigInfo) {
  207. Map resultInfo = new HashMap();
  208. String tfsReturnFileName = null;
  209. long block = 10 * 1024;// 默认
  210. if (ftpClientTemplate == null) {
  211. resultInfo.put("SAVE_FILE_FLAG", "E");
  212. return resultInfo;
  213. }
  214. String localPathName = ftpItemConfigInfo.get("localPath").toString().endsWith("/") ? ftpItemConfigInfo.get("localPath").toString()
  215. + ftpItemConfigInfo.get("newFileName").toString() : ftpItemConfigInfo.get("localPath").toString() + "/" + ftpItemConfigInfo.get("newFileName").toString();
  216. ftpItemConfigInfo.put("localfilename", localPathName);// 本地带路径的文件名回写,后面读文件时使用
  217. try {
  218. File file = new File(localPathName);
  219. RandomAccessFile accessFile = new RandomAccessFile(file, "rwd");// 建立随机访问
  220. FTPClient ftpClient = new FTPClient();
  221. ftpClient.connect(ftpClientTemplate.getHost(), ftpClientTemplate.getPort());
  222. if (FTPReply.isPositiveCompletion(ftpClient.getReplyCode())) {
  223. if (!ftpClient.login(ftpClientTemplate.getUsername(), ftpClientTemplate.getPassword())) {
  224. resultInfo.put("SAVE_FILE_FLAG", "E");
  225. resultInfo.put("remark", "登录失败,用户名【" + ftpClientTemplate.getUsername() + "】密码【" + ftpClientTemplate.getPassword() + "】");
  226. return resultInfo;
  227. }
  228. }
  229. ftpClient.setControlEncoding("UTF-8");
  230. ftpClient.setFileType(FTP.BINARY_FILE_TYPE); // 二进制
  231. ftpClient.enterLocalPassiveMode(); // 被动模式
  232. ftpClient.sendCommand("PASV");
  233. ftpClient.sendCommand("SIZE " + remoteFileNameTmp + "\r\n");
  234. String replystr = ftpClient.getReplyString();
  235. String[] replystrL = replystr.split(" ");
  236. long filelen = 0;
  237. if (Integer.valueOf(replystrL[0]) == 213) {
  238. filelen = Long.valueOf(replystrL[1].trim());
  239. } else {
  240. resultInfo.put("SAVE_FILE_FLAG", "E");
  241. resultInfo.put("remark", "无法获取要下载的文件的大小!");
  242. return resultInfo;
  243. }
  244. accessFile.setLength(filelen);
  245. accessFile.close();
  246. ftpClient.disconnect();
  247. int tnum = Integer.valueOf(ftpItemConfigInfo.get("TNUM").toString());
  248. block = (filelen + tnum - 1) / tnum;// 每个线程下载的快大小
  249. ThreadPoolExecutor cachedThreadPool = (ThreadPoolExecutor) Executors.newCachedThreadPool();
  250. List<Future<Map>> threadR = new ArrayList<Future<Map>>();
  251. for (int i = 0; i < tnum; i++) {
  252. logger.debug("发起线程:" + i);
  253. // 保存线程日志
  254. ftpItemConfigInfo.put("threadrunstate", "R");
  255. ftpItemConfigInfo.put("remark", "开始下载文件");
  256. ftpItemConfigInfo.put("data", "文件名:" + remoteFileNameTmp);
  257. long start = i * block;
  258. long end = (i + 1) * block - 1;
  259. ftpItemConfigInfo.put("begin", start);
  260. ftpItemConfigInfo.put("end", end);
  261. saveTaskLogDetail(ftpItemConfigInfo);
  262. Map para = new HashMap();
  263. para.putAll(ftpItemConfigInfo);
  264. para.put("serverfilename", remoteFileNameTmp);
  265. para.put("filelength", filelen);
  266. para.put("tnum", i + 0);
  267. para.put("threadDownSize", block);
  268. para.put("transferflag", FTPClientTemplate.TransferType.download);
  269. FTPClientTemplate dumpThread = new FTPClientTemplate(para);
  270. Future<Map> runresult = cachedThreadPool.submit(dumpThread);
  271. threadR.add(runresult);
  272. }
  273. do {
  274. // 等待下载完成
  275. Thread.sleep(1000);
  276. } while (cachedThreadPool.getCompletedTaskCount() < threadR.size());
  277. saveDownFileData(ftpItemConfigInfo);
  278. // 下载已经完成,多线程保存数据至表中
  279. } catch (Exception e) {
  280. // TODO Auto-generated catch block
  281. logger.error("保存文件失败:", e);
  282. resultInfo.put("SAVE_FILE_FLAG", "E");
  283. resultInfo.put("remark", "保存文件失败:" + e);
  284. return resultInfo;
  285. }
  286. resultInfo.put("SAVE_FILE_FLAG", "S");
  287. return resultInfo;
  288. }
  289. @Override
  290. protected void post(Map ftpItemConfigInfo) {
  291. // TODO Auto-generated method stub
  292. if (ftpItemConfigInfo != null && ftpItemConfigInfo.containsValue("AFTERFLAG") && "0".equals(ftpItemConfigInfo.get("AFTERFLAG"))
  293. && ftpItemConfigInfo.containsKey("AFTERFUNCTION") && ftpItemConfigInfo.get("AFTERFUNCTION") != null && !"".equals(ftpItemConfigInfo.get("AFTERFUNCTION"))) {
  294. // 这个时候确定已经进入了事后过程
  295. if (ftpItemConfigInfo.containsKey("retVal") && !"0000".equals(ftpItemConfigInfo.get("retVal"))) {
  296. ftpItemConfigInfo.put("threadrunstate", "E3");
  297. ftpItemConfigInfo.put("remark", ftpItemConfigInfo.get("retVal"));
  298. saveTaskLogDetail(ftpItemConfigInfo);
  299. }
  300. }
  301. }
  302. /**
  303. * 保存下载的文件里的内容到配置的数据表中,多线程同时保存
  304. */
  305. public void saveDownFileData(Map taskInfo) {
  306. // 先分配每个线程处理的起止位置
  307. List contthr = contSaveThreadContInfo(taskInfo);
  308. // 启动多个线程,保存文件内容
  309. if (contthr != null) {
  310. try {
  311. int tnum = Integer.valueOf(taskInfo.get("TNUM").toString());
  312. ThreadPoolExecutor cachedThreadPool = (ThreadPoolExecutor) Executors.newCachedThreadPool();
  313. List<Future<Map>> threadR = new ArrayList<Future<Map>>();
  314. for (int i = 0; i < tnum; i++) {
  315. logger.debug("发起线程:" + i);
  316. Map tmp = (Map) contthr.get(i);
  317. Map para = new HashMap();
  318. para.putAll(taskInfo);
  319. para.putAll(tmp);
  320. para.put("transferflag", FTPClientTemplate.TransferType.savedata);
  321. FTPClientTemplate dumpThread = new FTPClientTemplate(para);
  322. Future<Map> runresult = cachedThreadPool.submit(dumpThread);
  323. threadR.add(runresult);
  324. }
  325. do {
  326. // 等待保存数据
  327. Thread.sleep(1000);
  328. } while (cachedThreadPool.getCompletedTaskCount() < threadR.size());
  329. taskInfo.put("SAVE_FILE_FLAG", "S");
  330. logger.debug("文件内容保存完毕!");
  331. } catch (Exception ex) {
  332. logger.error("保存文件失败", ex);
  333. taskInfo.put("SAVE_FILE_FLAG", "E");
  334. taskInfo.put("remark", "保存文件失败" + ex);
  335. } catch (Throwable ex) {
  336. logger.error("保存文件失败", ex);
  337. taskInfo.put("SAVE_FILE_FLAG", "E");
  338. taskInfo.put("remark", "保存文件失败" + ex);
  339. }
  340. }
  341. }
  342. /**
  343. * 分配每个线程处理的起止位置
  344. */
  345. public static List contSaveThreadContInfo(Map taskInfo) {
  346. List contlist = new ArrayList();
  347. try {
  348. RandomAccessFile raf = new RandomAccessFile(taskInfo.get("localfilename").toString(), "r");
  349. int tnum = Integer.valueOf(taskInfo.get("TNUM").toString());
  350. long filelen = raf.length();
  351. if (filelen != 0L) {
  352. long block = (filelen + tnum - 1) / tnum;// 每个线程下载的快大小
  353. long begin = 0;
  354. // 修复在保存数据时的线程数太大时,无法去除titleflag 行的问题处理
  355. if ("0".equals(taskInfo.get("linecountflag"))) {
  356. String temp = raf.readLine();
  357. begin += temp.length();
  358. }
  359. if ("0".equals(taskInfo.get("titleflag"))) {
  360. String temp = raf.readLine();
  361. begin += temp.length();
  362. }
  363. for (int i = 0; i < tnum; i++) {
  364. if (i == tnum - 1) {
  365. Map tmp = new HashMap();
  366. tmp.put("begin", begin);
  367. tmp.put("end", filelen - 1);
  368. contlist.add(tmp);
  369. break;
  370. }
  371. long end = (i + 1) * block - 1;
  372. // 处理如果有总数行和文件头行时,线程数大于数据行数,映入的问题
  373. if (end < begin) {
  374. begin = 0;
  375. end = begin;
  376. Map tmp = new HashMap();
  377. tmp.put("begin", begin);
  378. tmp.put("end", end);
  379. contlist.add(tmp);
  380. begin = end + 1;
  381. continue;
  382. }
  383. while (end > 0) {
  384. raf.seek(end);
  385. if (raf.readByte() == '\n') {
  386. Map tmp = new HashMap();
  387. tmp.put("begin", begin);
  388. tmp.put("end", end);
  389. contlist.add(tmp);
  390. begin = end + 1;
  391. break;
  392. }
  393. end--;
  394. }
  395. }
  396. }
  397. raf.close();
  398. } catch (Exception ex) {
  399. ex.printStackTrace();
  400. return null;
  401. }
  402. return contlist;
  403. }
  404. }