HcFtpToFileSystemQuartz.java 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379
  1. package com.java110.job.smo;
  2. import com.java110.common.constant.RuleDomain;
  3. import com.java110.common.util.DateUtil;
  4. import com.java110.common.util.StringUtil;
  5. import com.java110.job.dao.IHcFtpFileDAO;
  6. import com.java110.job.model.FtpTaskLog;
  7. import com.java110.job.model.FtpTaskLogDetail;
  8. import org.apache.commons.lang.StringUtils;
  9. import org.slf4j.Logger;
  10. import org.slf4j.LoggerFactory;
  11. import org.springframework.beans.factory.annotation.Autowired;
  12. import java.text.DateFormat;
  13. import java.text.SimpleDateFormat;
  14. import java.util.*;
  15. /**
  16. *
  17. * @author
  18. *
  19. */
  20. public abstract class HcFtpToFileSystemQuartz{
  21. protected static final Logger logger = LoggerFactory.getLogger(HcFtpToFileSystemQuartz.class);
  22. @Autowired
  23. private IHcFtpFileDAO iHcFtpFileDAO;
  24. @Autowired
  25. private IHcFtpFileSMO iHcFtpFileSMO;
  26. /*private IPrvncDumpSMO prvncDumpSMO;*/
  27. // 运行状态,R:正在执行 T:等待运行 TD1:文件下载失败 TD2:文件内容保存失败 TU1:数据文件生成失败 TU2:数据文件上传失败
  28. private static final String TASK_STATE_R = "R";// 正在运行
  29. private static final String TASK_STATE_T = "T";// 等待运行
  30. private static final String TASK_STATE_E1 = "E1";// 执行事前过程失败
  31. private static final String TASK_STATE_E2 = "E2";// 处理数据失败
  32. private static final String TASK_STATE_E3 = "E3";// 执行事后过程失败
  33. public void initTask() {
  34. // 将所有的任务状态改为等待运行状态
  35. Map paramIn = new HashMap();
  36. paramIn.put("oldRunState", "R");
  37. paramIn.put("runState", "T");
  38. int updateFtpItemRunStateFlag = iHcFtpFileDAO.updateFtpItemRunState(paramIn);
  39. if (updateFtpItemRunStateFlag < 1) {
  40. logger.error("--【PrvncFtpToFileSystemQuartz.initTask】,没有需要更新的内容(没有下载一半后停止应用的情况)", paramIn);
  41. }
  42. }
  43. /**
  44. * 启动任务
  45. *
  46. * @param ftpItemConfigInfo
  47. */
  48. public void startFtpTask(Map ftpItemConfigInfo) throws Exception {
  49. // 这么做是为了,单线程调用,防止多线程导致数据重复处理
  50. if (!ftpItemConfigInfo.containsKey("RUN_STATE") || "R".equals(ftpItemConfigInfo.get("RUN_STATE"))) {
  51. return;
  52. }
  53. long taskId = Long.parseLong(ftpItemConfigInfo.get("TASKID").toString());
  54. if (logger.isDebugEnabled()) {
  55. logger.debug("---【PrvncFtpToFileSystemQuartz.startFtpTask】:任务【" + taskId + "】开始运行!", taskId);
  56. }
  57. // 保存任务执行主要日志信息
  58. //获取LOGID 默认生成规则为tadkid去掉年月日之前的值+66
  59. String id = ftpItemConfigInfo.get("TASKID").toString();
  60. id = id.substring(10,id.length());
  61. long logid = Long.parseLong (id+"22");
  62. ftpItemConfigInfo.put("logid",logid);
  63. long taskLogID = insertTaskInfo(ftpItemConfigInfo);
  64. ftpItemConfigInfo.put("logid", taskLogID);
  65. ftpItemConfigInfo.put("taskid", taskId);
  66. ftpItemConfigInfo.put("threadrunstate", TASK_STATE_R);
  67. ftpItemConfigInfo.put("tnum", 1);
  68. // 修改任务状态为正在执行状态
  69. updateTaskState(taskId, TASK_STATE_R);
  70. // 方法调用是否成功,S成功(默认),E表示失败(在方法中失败时,需要修改)
  71. ftpItemConfigInfo.put("PRE_METHOD_FLAG", "S");
  72. try {
  73. // 1.0空方法,让子类去实现
  74. prepare(ftpItemConfigInfo);
  75. // 2.0调用事前过程
  76. if (ftpItemConfigInfo.containsKey("PREFLAG") && "0".equals(ftpItemConfigInfo.get("PREFLAG"))) {
  77. callPreFunction(ftpItemConfigInfo);
  78. }
  79. if (ftpItemConfigInfo.containsKey("PRE_METHOD_FLAG") && "E".equals(ftpItemConfigInfo.get("PRE_METHOD_FLAG"))) {
  80. // 此时调用事前过程失败,直接返回 查询标识为E,更新日志
  81. udpateTaskLog(ftpItemConfigInfo);
  82. updateTaskState(taskId, TASK_STATE_E1);
  83. return;
  84. }
  85. // 3.0核心业务处理逻辑,需要子类去实现
  86. process(ftpItemConfigInfo);
  87. if (ftpItemConfigInfo.containsKey("PRE_METHOD_FLAG") && "E".equals(ftpItemConfigInfo.get("PRE_METHOD_FLAG"))) {
  88. // 程序处理失败,直接返回
  89. ftpItemConfigInfo.put("threadrunstate", TASK_STATE_E2);
  90. updateTaskState(taskId, TASK_STATE_E2);
  91. udpateTaskLog(ftpItemConfigInfo);
  92. saveTaskLogDetail(ftpItemConfigInfo);// 保存detail
  93. return;
  94. }
  95. // 记录详细日志
  96. ftpItemConfigInfo.put("threadrunstate", "T");
  97. saveTaskLogDetail(ftpItemConfigInfo);
  98. // 4.0调用事后过程
  99. if (ftpItemConfigInfo.containsKey("AFTERFLAG") && "0".equals(ftpItemConfigInfo.get("AFTERFLAG"))) {
  100. callAfterFunction(ftpItemConfigInfo);
  101. }
  102. if (ftpItemConfigInfo.containsKey("PRE_METHOD_FLAG") && "E".equals(ftpItemConfigInfo.get("PRE_METHOD_FLAG"))) {
  103. // 此时调用事前过程失败,直接返回 查询标识为E,更新日志
  104. udpateTaskLog(ftpItemConfigInfo);
  105. updateTaskState(taskId, TASK_STATE_E3);
  106. return;
  107. }
  108. // 5.0空方法,让子类去实现
  109. post(ftpItemConfigInfo);
  110. } catch (Exception ex) {
  111. ftpItemConfigInfo.put("threadrunstate", TASK_STATE_E2);
  112. udpateTaskLog(ftpItemConfigInfo);
  113. ftpItemConfigInfo.put("threadrunstate", TASK_STATE_E2);
  114. ftpItemConfigInfo.put("remark", ex);
  115. saveTaskLogDetail(ftpItemConfigInfo);
  116. updateTaskState(taskId, TASK_STATE_E2);
  117. // 接续向外抛出去
  118. logger.error("处理出现问题:", ex);
  119. return;
  120. }
  121. // 修改任务状态为执行完毕状态
  122. updateTaskState(taskId, TASK_STATE_T);
  123. ftpItemConfigInfo.put("threadrunstate", TASK_STATE_T);
  124. udpateTaskLog(ftpItemConfigInfo);
  125. // 发送任务运行结果通知短信给相关人员 **暂时不调用短信
  126. if (!TASK_STATE_T.equals(ftpItemConfigInfo.get("RUN_STATE").toString())) {
  127. /*sendErrLogPhoneMsg(ftpItemConfigInfo, taskLogID);*/
  128. }
  129. }
  130. // 如果有事前存过需要调用,则先调用存过
  131. public void callPreFunction(Map taskInfo) {
  132. if (taskInfo.containsKey("PREFUNCTION") && taskInfo.get("PREFUNCTION") != null && !"".equals(taskInfo.get("PREFUNCTION"))) {
  133. try {
  134. iHcFtpFileSMO.saveDbFunction(taskInfo.get("PREFUNCTION").toString());
  135. taskInfo.put("threadrunstate", "T");
  136. taskInfo.put("remark", "调用事前存过结束");
  137. saveTaskLogDetail(taskInfo);
  138. } catch (Exception ex) {
  139. logger.error("调用事前存过失败:", ex);
  140. taskInfo.put("threadrunstate", "E1");
  141. taskInfo.put("remark", "调用事前存过失败" + ex);
  142. taskInfo.put("PRE_METHOD_FLAG", "E");
  143. saveTaskLogDetail(taskInfo);
  144. }
  145. }
  146. }
  147. /**
  148. * 主要业务处理(上传下载),让子类去实现
  149. *
  150. * @param ftpItemConfigInfo
  151. */
  152. protected abstract void process(Map ftpItemConfigInfo) throws Exception;
  153. // 如果有事后存过需要调用,则调用存过
  154. public void callAfterFunction(Map taskInfo) {
  155. if (taskInfo.containsKey("AFTERFUNCTION") && taskInfo.get("AFTERFUNCTION") != null && !"".equals(taskInfo.get("AFTERFUNCTION"))) {
  156. try {
  157. taskInfo.put("functionname", taskInfo.get("AFTERFUNCTION"));
  158. // taskInfo 参数param需要在process方法中需要自己写入
  159. iHcFtpFileSMO.saveDbFunctionWithParam(taskInfo);
  160. taskInfo.put("threadrunstate", "T");
  161. taskInfo.put("remark", "调用事后存过结束");
  162. saveTaskLogDetail(taskInfo);
  163. } catch (Exception ex) {
  164. ex.printStackTrace();
  165. taskInfo.put("threadrunstate", "E3");
  166. taskInfo.put("remark", "调用事后存过失败" + ex);
  167. taskInfo.put("PRE_METHOD_FLAG", "E");
  168. saveTaskLogDetail(taskInfo);
  169. }
  170. }
  171. }
  172. /**
  173. * 修改任务状态
  174. *
  175. * @return
  176. */
  177. private void updateTaskState(long taskId, String state) {
  178. Map info = new HashMap();
  179. info.put("taskId", taskId);
  180. info.put("runState", state);
  181. int updateFtpItemFlag = iHcFtpFileDAO.updateFtpItemByTaskId(info);
  182. // 这里只是后台提示,不进行日志保存
  183. if (updateFtpItemFlag < 1) {
  184. logger.error("---【PrvncFtpToFileSystemQuartz.updateTaskState】修改任务【" + taskId + "】的状态失败", info);
  185. }
  186. }
  187. /**
  188. * 修改任务执行日志的状态
  189. */
  190. private void udpateTaskLog(Map taskInfo) {
  191. FtpTaskLog loginfo = new FtpTaskLog();
  192. loginfo.setLogid(Long.valueOf(taskInfo.get("logid").toString()));
  193. loginfo.setState(taskInfo.get("threadrunstate").toString());
  194. iHcFtpFileSMO.updateTaskRunLog(loginfo);
  195. }
  196. /**
  197. * 保存任务执行的详细日志
  198. */
  199. protected void saveTaskLogDetail(Map taskInfo) {
  200. FtpTaskLogDetail logdetail = new FtpTaskLogDetail();
  201. logdetail.setId(Long.valueOf(taskInfo.get("logid").toString()+"66"));
  202. logdetail.setLogid(Long.valueOf(taskInfo.get("logid").toString()));
  203. logdetail.setTaskid(Long.valueOf(taskInfo.get("taskid").toString()));
  204. logdetail.setState((String) taskInfo.get("threadrunstate"));
  205. logdetail.setTnum(Integer.valueOf(taskInfo.get("tnum").toString()));
  206. if (taskInfo.get("begin") != null) {
  207. logdetail.setBegin(Long.valueOf(taskInfo.get("begin").toString()));
  208. }
  209. if (taskInfo.get("end") != null) {
  210. logdetail.setEnd(Long.valueOf(taskInfo.get("end").toString()));
  211. }
  212. if (taskInfo.get("havedown") != null) {
  213. logdetail.setHavedown(Long.valueOf(taskInfo.get("havedown").toString()));
  214. }
  215. logdetail.setRemark(taskInfo.get("remark") == null ? "" : (taskInfo.get("remark").toString().trim().length() > 2000 ? taskInfo.get("remark").toString().trim().substring(0,
  216. 1600) : taskInfo.get("remark").toString().trim()));
  217. logdetail.setData(taskInfo.get("data") == null ? "" : taskInfo.get("data").toString());
  218. logdetail.setServerfilename(taskInfo.get("serverfilename") == null ? "" : taskInfo.get("serverfilename").toString());
  219. logdetail.setLocalfilename(taskInfo.get("localfilename") == null ? "" : taskInfo.get("localfilename").toString());
  220. int logdetailid = iHcFtpFileSMO.saveTaskRunDetailLog(logdetail);
  221. taskInfo.put("logdetailid", logdetailid);
  222. }
  223. // /**
  224. // * 修改任务执行的详细日志的状态
  225. // */
  226. // private void updateTaskLogDetail(Map taskInfo){
  227. // FtpTaskLogDetail logdetail=new FtpTaskLogDetail();//
  228. // logdetail.setId(Long.valueOf(taskInfo.get("logdetailid").toString()));
  229. // logdetail.setState(taskInfo.get("threadrunstate").toString());
  230. // logdetail.setRemark((String)taskInfo.get("remark"));
  231. // logdetail.setData((String)taskInfo.get("data"));
  232. // if(taskInfo.get("downedlength")!=null)
  233. // logdetail.setHavedown(Long.valueOf(taskInfo.get("downedlength").toString()));
  234. // prvncFtpFileSMO.updateTaskRunDetailLog(logdetail);
  235. // }
  236. /**
  237. * 生成任务执行日志
  238. */
  239. private long insertTaskInfo(Map taskInfo) {
  240. FtpTaskLog loginfo = new FtpTaskLog();
  241. loginfo.setTaskid(Long.valueOf(taskInfo.get("TASKID").toString()));
  242. loginfo.setState("R");
  243. loginfo.setServerfilename("");// taskInfo.get("serverfilename").toString()
  244. loginfo.setLocalfilename("");// taskInfo.get("localfilename").toString()
  245. loginfo.setUord(taskInfo.get("U_OR_D").toString());
  246. return iHcFtpFileSMO.saveTaskRunLog(loginfo);
  247. }
  248. /**
  249. * 如果任务运行有异常,则发送警告短信给配置的手机号码
  250. */
  251. private void sendErrLogPhoneMsg(Map taskInfo, long taskLogID) {
  252. Map msginfo = new HashMap();
  253. String phone = (String) taskInfo.get("errphone");
  254. if (phone != null && !"".equals(phone)) {
  255. String[] phonelist = phone.split(",");
  256. for (int i = 0; i < phonelist.length; i++) {
  257. msginfo.put("taskid", taskInfo.get("taskid"));
  258. msginfo.put("phone", phonelist[i]);
  259. msginfo.put("msg", "通用FTP数据文件传接任务:" + (String) taskInfo.get("taskname") + "运行提示");
  260. DateFormat df = new SimpleDateFormat("yyyy-mm-dd HH:mm:ss");
  261. String detail = "任务已于" + df.format(new Date()) + "运行完毕。运行过程中出现异常,详情请登录系统查看!";
  262. msginfo.put("detail", detail);
  263. /*prvncDumpSMO.saveTaskErrInfoPhoneMsg(msginfo);*/
  264. }
  265. }
  266. }
  267. /**
  268. * 处理文件名,校验文件名是中是否存在****(4个),表示通配符,如果不存在就是确定唯一文件名
  269. * 文件名支持日期型的如CRM_########001.txt 程序处理后是 CRM_20170105001.txt 文件名支持sql 语句生成的
  270. * 文件名支持通配符的如863_****.txt 程序下载所有以863_开头的文件 863_****001.txt
  271. * 以863_开头,以001结尾,****001.txt 以001结尾的
  272. *
  273. * @param fileName
  274. * @return
  275. */
  276. protected List<String> dealFileName(String fileName) {
  277. // TODO Auto-generated method stub
  278. List<String> results = new ArrayList<String>();
  279. String result = "";
  280. // 文件中使用的日期
  281. if (StringUtils.contains(fileName, RuleDomain.REPLAY_TYPE_F)) {
  282. result = StringUtil.replace(fileName, RuleDomain.REPLAY_TYPE_F, DateUtil.getFormatTimeString(new Date(), "yyyyMMddHHmm"));
  283. } else if (StringUtils.contains(fileName, RuleDomain.REPLAY_TYPE_E)) {
  284. result = StringUtil.replace(fileName, RuleDomain.REPLAY_TYPE_E, DateUtil.getFormatTimeString(new Date(), "yyyyMMddHH"));
  285. } else if (StringUtils.contains(fileName, RuleDomain.REPLAY_TYPE_A)) {
  286. result = StringUtil.replace(fileName, RuleDomain.REPLAY_TYPE_A, DateUtil.getFormatTimeString(new Date(), "yyyyMMdd"));
  287. } else if (StringUtils.contains(fileName == null ? "" : fileName.toLowerCase(), RuleDomain.REPLAY_TYPE_SQL)) {
  288. // 后期改造,文件名如果配置的是sql的话,以sql查询文件名
  289. List<String> fileNames = this.getPrvncFtpFileDAO().execConfigSql(fileName);
  290. // if (fileNames != null && fileNames.size() > 0) {
  291. // result = fileNames.get(0);
  292. // }
  293. return fileNames;
  294. } else {
  295. result = fileName;
  296. }
  297. results.add(result);
  298. return results;
  299. }
  300. /**
  301. * 空方法,如果在事前过程处理前,还需要做一定的处理,需要子类重写这个方法,实现业务逻辑
  302. *
  303. * @param ftpItemConfigInfo
  304. */
  305. protected void prepare(Map ftpItemConfigInfo) {
  306. }
  307. /**
  308. * 空方法,如果在事后过程处理完后,还需要做一定的处理,需要子类重写这个方法,实现业务逻辑
  309. *
  310. * @param ftpItemConfigInfo
  311. */
  312. protected void post(Map ftpItemConfigInfo) {
  313. }
  314. public IHcFtpFileDAO getPrvncFtpFileDAO() {
  315. return iHcFtpFileDAO;
  316. }
  317. public void setPrvncFtpFileDAO(IHcFtpFileDAO prvncFtpFileDAO) {
  318. this.iHcFtpFileDAO = prvncFtpFileDAO;
  319. }
  320. public IHcFtpFileSMO getPrvncFtpFileSMO() {
  321. return iHcFtpFileSMO;
  322. }
  323. public void setPrvncFtpFileSMO(IHcFtpFileSMO prvncFtpFileSMO) {
  324. this.iHcFtpFileSMO = prvncFtpFileSMO;
  325. }
  326. /*public IPrvncDumpSMO getPrvncDumpSMO() {
  327. return prvncDumpSMO;
  328. }
  329. public void setPrvncDumpSMO(IPrvncDumpSMO prvncDumpSMO) {
  330. this.prvncDumpSMO = prvncDumpSMO;
  331. }
  332. */
  333. }