HcFtpToFileSystemConfigAction.java 26 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748
  1. package com.java110.job.Api;
  2. import com.alibaba.fastjson.JSONArray;
  3. import com.alibaba.fastjson.JSONObject;
  4. import com.java110.common.util.SpringBeanInvoker;
  5. import com.java110.core.factory.GenerateCodeFactory;
  6. import com.java110.job.common.CustomizedPropertyPlaceholderConfigurer;
  7. import com.java110.job.dao.IHcFtpFileDAO;
  8. import com.java110.job.smo.DownloadFileFromFtpToTable;
  9. import com.java110.job.task.HcFtpToFileSystemJob;
  10. import org.apache.commons.validator.GenericValidator;
  11. import org.apache.commons.validator.util.ValidatorUtils;
  12. import org.quartz.*;
  13. import org.slf4j.Logger;
  14. import org.slf4j.LoggerFactory;
  15. import org.springframework.beans.factory.annotation.Autowired;
  16. import org.springframework.stereotype.Service;
  17. import javax.servlet.http.HttpServletRequest;
  18. import java.text.DateFormat;
  19. import java.text.SimpleDateFormat;
  20. import java.util.*;
  21. /**
  22. * 将ftp上的文件保存到支持的文件系统
  23. *
  24. * @author wuxw7 add by 20170103
  25. * shiyj update by 2019.08.29
  26. *
  27. */
  28. @Service
  29. public class HcFtpToFileSystemConfigAction {
  30. private static final Logger logger = LoggerFactory.getLogger(HcFtpToFileSystemConfigAction.class);
  31. private static final String defaultCronExpression = "0 * * * * ?";// 每分钟执行一次
  32. private static final String prefixJobName = "HcFtpToSystem_"; // job
  33. // 名称前缀,防止和其他的job名称产生冲突
  34. private static final String RUNFLAG_START = "1";
  35. private static final String RUNFLAG_STOP = "0";
  36. @Autowired
  37. private IHcFtpFileDAO iHcFtpFileDAO;
  38. @Autowired
  39. private Scheduler scheduler;
  40. // 每页数据条数
  41. private static int pageSize = 20;
  42. public JSONObject resultMsg;
  43. /**
  44. *
  45. */
  46. private static final long serialVersionUID = 1L;
  47. private DateFormat df = new SimpleDateFormat("yyyy-MM-dd");
  48. /**
  49. * 查询配置的需要下载任务的列表
  50. *
  51. * @return
  52. */
  53. public JSONObject queryFtpItems(HttpServletRequest request) {
  54. String curPage = request.getParameter("curPage");
  55. curPage = curPage == null || "".equals(curPage) ? "1" : curPage;
  56. // 查询在用状态时的下载任务列表
  57. Map info = new HashMap();
  58. info.put("curPage", (Integer.parseInt(curPage)-1)*pageSize);
  59. info.put("pageSize", pageSize*Integer.parseInt(curPage));
  60. Map resultInfo = iHcFtpFileDAO.queryFtpItems(info);
  61. // 获取总数据数
  62. int dataCount = resultInfo.get("ITEMSCOUNT") == null ? 0 : Integer.parseInt(resultInfo.get("ITEMSCOUNT").toString());
  63. // 计算页数
  64. if (dataCount % pageSize == 0) {
  65. dataCount /= pageSize;
  66. } else {
  67. dataCount = dataCount / pageSize + 1;
  68. }
  69. // 获取数据
  70. List<Map> ftpItems = resultInfo.get("DATA") == null ? null : (List) resultInfo.get("DATA");
  71. // {"total":10,"rows":[{},{}]}
  72. JSONObject data = new JSONObject();
  73. data.put("total", dataCount);
  74. data.put("currentPage", curPage);
  75. if (ftpItems != null && ftpItems.size() > 0) {
  76. JSONArray rows = new JSONArray();
  77. for (int itemIndex = 0; itemIndex < ftpItems.size(); itemIndex++) {
  78. // 处理时间显示和界面显示传输类型
  79. Map ftpItemMap = ftpItems.get(itemIndex);
  80. ftpItemMap.put("U_OR_D_NAME", ftpItemMap.get("U_OR_D"));// 暂且写死,最终还是读取配置
  81. ftpItemMap.put("CREATE_DATE", df.format(ftpItemMap.get("CREATE_DATE")));// 暂且写死,最终还是读取配置
  82. rows.add(JSONObject.parseObject(JSONObject.toJSONString(ftpItems.get(itemIndex))));
  83. }
  84. data.put("rows", rows);
  85. resultMsg = data;
  86. return data;
  87. }
  88. data.put("rows", "[]");
  89. resultMsg = data;
  90. return data;
  91. }
  92. /**
  93. * 增加Ftp配置
  94. *
  95. * @return
  96. */
  97. public JSONObject addFtpItem(HttpServletRequest request) {
  98. // 请求参数
  99. String ftpItemJson = request.getParameter("ftpItemJson");
  100. JSONObject ftpItemJsonObj = null;
  101. try {
  102. // 校验格式是否正确
  103. ftpItemJsonObj = JSONObject.parseObject(ftpItemJson);
  104. } catch (Exception e) {
  105. logger.error("传入参数格式不正确:" + ftpItemJson, e);
  106. resultMsg = createResultMsg("1999", "传入参数格式不正确:" + ftpItemJson, "");
  107. return resultMsg;
  108. }
  109. // 将ftpItemJson装为Map保存操作
  110. Map paramIn = JSONObject.parseObject(ftpItemJsonObj.getJSONObject("taskInfo").toJSONString(), Map.class);
  111. // 数据规范性校验
  112. Object dealClassObj = null;
  113. // 在prvncCrm.properties 文件中获取对应处理类
  114. if ("DT".equals(paramIn.get("uOrD").toString())) {
  115. dealClassObj = DownloadFileFromFtpToTable.class;
  116. }else{
  117. resultMsg = this.createResultMsg("1999", "对应模板不存在,请联系管理员", "");
  118. return resultMsg;
  119. }
  120. // Object dealClassObj = "provInner.DownloadFileFromFtpToTFS";
  121. if (dealClassObj == null) {
  122. resultMsg = this.createResultMsg("1999", "对应模板不存在,请联系管理员", "");
  123. return resultMsg;
  124. }
  125. String dealClass = dealClassObj.toString();
  126. String taskId = GenerateCodeFactory.getGeneratorId(GenerateCodeFactory.CODE_PREFIX_HCJOBId);
  127. // 保存数据
  128. paramIn.put("taskId", taskId);
  129. paramIn.put("dealClass", dealClass);
  130. int addFtpItemFlag = iHcFtpFileDAO.addFtpItem(paramIn);
  131. if (addFtpItemFlag > 0) {
  132. // #taskId#,#itemSpecId#,#value#
  133. // 保存属性信息
  134. JSONArray taskAttrs = ftpItemJsonObj.getJSONArray("taskAttrs");
  135. List<Map> taskAttrsList = new ArrayList<Map>();
  136. for (int taskAttrIndex = 0; taskAttrIndex < taskAttrs.size(); taskAttrIndex++) {
  137. JSONObject taskAttr = taskAttrs.getJSONObject(taskAttrIndex);
  138. Map taskAttrMap = new HashMap();
  139. taskAttrMap.put("taskId", taskId);
  140. taskAttrMap.put("itemSpecId", taskAttr.get("itemSpecId"));
  141. taskAttrMap.put("value", taskAttr.get("value"));
  142. taskAttrsList.add(taskAttrMap);
  143. }
  144. int addFtpItemAttrFlag = iHcFtpFileDAO.addFtpItemAttrs(taskAttrsList);
  145. if (addFtpItemAttrFlag > 0) {
  146. resultMsg = this.createResultMsg("0000", "成功", ftpItemJson);
  147. } else {
  148. resultMsg = this.createResultMsg("1999", "保存属性失败", "");
  149. }
  150. return resultMsg;
  151. }
  152. resultMsg = this.createResultMsg("1999", "保存数据失败", "");
  153. return resultMsg;
  154. }
  155. /**
  156. * 编辑Ftp配置,编辑时将会修改所有的ftpItem 信息,所以传递是所有字段需要传递全
  157. *
  158. * @return
  159. */
  160. public String editFtpItem(HttpServletRequest request) {
  161. // 请求参数为{"taskId":"12","taskName":"经办人照片同步处理","ftpUserName":"weblogic",.....}
  162. String ftpItemJson = request.getParameter("ftpItemJson");
  163. JSONObject ftpItemJsonObj = null;
  164. try {
  165. // 校验格式是否正确
  166. ftpItemJsonObj = JSONObject.parseObject(ftpItemJson);
  167. } catch (Exception e) {
  168. logger.error("传入参数格式不正确:" + ftpItemJson, e);
  169. resultMsg = createResultMsg("1999", "传入参数格式不正确:" + ftpItemJson, "");
  170. return "editFtpItem";
  171. }
  172. // 将ftpItemJson装为Map保存操作
  173. Map paramIn = JSONObject.parseObject(ftpItemJsonObj.getJSONObject("taskInfo").toJSONString(), Map.class);
  174. // 在prvncCrm.properties 文件中获取对应处理类
  175. Object dealClassObj = CustomizedPropertyPlaceholderConfigurer.getContextProperty("task.deal.class." + paramIn.get("uOrD"));
  176. // Object dealClassObj = "provInner.DownloadFileFromFtpToTFS";
  177. if (dealClassObj == null) {
  178. resultMsg = this.createResultMsg("1999", "对应模板不存在,请联系管理员", "");
  179. return "editFtpItem";
  180. }
  181. String dealClass = dealClassObj.toString();
  182. paramIn.put("dealClass", dealClass);
  183. // 根据taskId 查询记录是否存在,如果不存在直接返回失败
  184. Map ftpItem = iHcFtpFileDAO.queryFtpItemByTaskId(paramIn);
  185. // 判断是否有对应的数据
  186. if (ftpItem != null && ftpItem.containsKey("TASKID")) {
  187. // 更新数据
  188. int updateFtpItemFlag = iHcFtpFileDAO.updateFtpItemByTaskId(paramIn);
  189. if (updateFtpItemFlag > 0) {
  190. // 首先先删除
  191. iHcFtpFileDAO.deleteFtpItemAttrsbyTaskId(paramIn);
  192. // 保存属性信息
  193. JSONArray taskAttrs = ftpItemJsonObj.getJSONArray("taskAttrs");
  194. List<Map> taskAttrsList = new ArrayList<Map>();
  195. for (int taskAttrIndex = 0; taskAttrIndex < taskAttrs.size(); taskAttrIndex++) {
  196. JSONObject taskAttr = taskAttrs.getJSONObject(taskAttrIndex);
  197. Map taskAttrMap = new HashMap();
  198. taskAttrMap.put("taskId", paramIn.get("taskId"));
  199. taskAttrMap.put("itemSpecId", taskAttr.get("itemSpecId"));
  200. taskAttrMap.put("value", taskAttr.get("value"));
  201. taskAttrsList.add(taskAttrMap);
  202. }
  203. int addFtpItemAttrFlag = iHcFtpFileDAO.addFtpItemAttrs(taskAttrsList);
  204. if (addFtpItemAttrFlag > 0) {
  205. resultMsg = this.createResultMsg("0000", "成功", ftpItemJson);
  206. } else {
  207. resultMsg = this.createResultMsg("1999", "更新属性失败", "");
  208. }
  209. return "editFtpItem";
  210. }
  211. resultMsg = this.createResultMsg("1999", "修改的数据不存在或修改失败", "");
  212. return "editFtpItem";
  213. }
  214. resultMsg = this.createResultMsg("1999", "未找到对应的数据更新失败【" + paramIn.get("taskId") + "】", "");
  215. return "editFtpItem";
  216. }
  217. /**
  218. * 删除ftp配置
  219. *
  220. * @return
  221. */
  222. public String deleteFtpItem(HttpServletRequest request) {
  223. // 请求参数为{"tasks":[{"taskId":1},{"taskId":2}],"state":"DELETE"}
  224. String ftpItemJson = request.getParameter("ftpItemJson");
  225. if (logger.isDebugEnabled()) {
  226. logger.debug("---【PrvncFtpToFileSystemConfigAction.deleteFtpItem】入参为:" + ftpItemJson, ftpItemJson);
  227. }
  228. JSONObject paramIn = null;
  229. try {
  230. // 校验格式是否正确
  231. paramIn = JSONObject.parseObject(ftpItemJson);
  232. } catch (Exception e) {
  233. logger.error("传入参数格式不正确:" + ftpItemJson, e);
  234. resultMsg = createResultMsg("1999", "传入参数格式不正确:" + ftpItemJson + e, "");
  235. return "deleteFtpItem";
  236. }
  237. // 传入报文不为空
  238. if (paramIn == null || !paramIn.containsKey("tasks") || !paramIn.containsKey("state")) {
  239. resultMsg = createResultMsg("1999", "传入参数格式不正确(必须包含tasks 和 state节点):" + ftpItemJson, "");
  240. return "deleteFtpItem";
  241. }
  242. // 校验当前是否为启动侦听
  243. if (!"DELETE".equals(paramIn.get("state"))) {
  244. resultMsg = createResultMsg("1999", "传入参数格式不正确(state的值必须是DELETE):" + ftpItemJson, "");
  245. return "deleteFtpItem";
  246. }
  247. // 查询需要操作的任务
  248. JSONArray taskInfos = paramIn.getJSONArray("tasks");
  249. String taskIds = "";
  250. for (int taskIndex = 0; taskIndex < taskInfos.size(); taskIndex++) {
  251. taskIds += (taskInfos.getJSONObject(taskIndex).getString("taskId") + ",");
  252. }
  253. if (taskIds.length() > 0) {
  254. taskIds = taskIds.substring(0, taskIds.length() - 1);
  255. }
  256. // 将ftpItemJson装为Map保存操作
  257. Map paramInfo = new HashMap();
  258. paramInfo.put("taskIds", taskIds.split(","));
  259. // 更新数据
  260. int updateFtpItemFlag = iHcFtpFileDAO.deleteFtpItemByTaskId(paramInfo);
  261. if (updateFtpItemFlag > 0) {
  262. resultMsg = this.createResultMsg("0000", "成功", ftpItemJson);
  263. return "deleteFtpItem";
  264. }
  265. resultMsg = this.createResultMsg("1999", "删除数据已经不存在,或删除失败", "");
  266. return "deleteFtpItem";
  267. }
  268. /**
  269. * 根据taskId 获取 ftp配置信息
  270. *
  271. * @return
  272. */
  273. public String queryFtpItemByTaskId(HttpServletRequest request) {
  274. // 请求参数为{"taskId":"12"}
  275. String ftpItemJson = request.getParameter("ftpItemJson");
  276. if (logger.isDebugEnabled()) {
  277. logger.debug("---【PrvncFtpToFileSystemConfigAction.queryFtpItemByTaskId】入参为:" + ftpItemJson, ftpItemJson);
  278. }
  279. try {
  280. // 校验格式是否正确
  281. JSONObject.parseObject(ftpItemJson);
  282. } catch (Exception e) {
  283. logger.error("传入参数格式不正确:" + ftpItemJson, e);
  284. resultMsg = createResultMsg("1999", "传入参数格式不正确:" + ftpItemJson, "");
  285. return "queryFtpItemByTaskId";
  286. }
  287. // 将ftpItemJson装为Map保存操作
  288. Map paramIn = JSONObject.parseObject(ftpItemJson, Map.class);
  289. // 根据taskId 查询记录是否存在,如果不存在直接返回失败
  290. Map ftpItem = iHcFtpFileDAO.queryFtpItemByTaskId(paramIn);
  291. // 判断是否有对应的数据
  292. if (ftpItem != null && ftpItem.containsKey("TASKID")) {
  293. // 更新数据
  294. this.createResultMsg("0000", "成功", JSONObject.toJSONString(ftpItem));
  295. return "queryFtpItemByTaskId";
  296. }
  297. resultMsg = this.createResultMsg("1999", "删除数据已经不存在,或删除失败", "");
  298. return "queryFtpItemByTaskId";
  299. }
  300. /**
  301. * 查询任务模板 ftpItemJson:{'uOrD':'U'}
  302. *
  303. * @return
  304. */
  305. public JSONObject questTaskTample(HttpServletRequest request) {
  306. // 请求参数为{"taskId":"12"}
  307. String ftpItemJson = request.getParameter("ftpItemJson");
  308. if (logger.isDebugEnabled()) {
  309. logger.debug("---【PrvncFtpToFileSystemConfigAction.queryFtpItemByTaskId】入参为:" + ftpItemJson, ftpItemJson);
  310. }
  311. JSONObject paramIn = null;
  312. try {
  313. // 校验格式是否正确
  314. paramIn = JSONObject.parseObject(ftpItemJson);
  315. } catch (Exception e) {
  316. logger.error("传入参数格式不正确:" + ftpItemJson, e);
  317. resultMsg = createResultMsg("1999", "传入参数格式不正确:" + ftpItemJson, "");
  318. return resultMsg;
  319. }
  320. String tample = paramIn.getString("uOrD");
  321. Map info = new HashMap();
  322. info.put("domain", tample);
  323. List<Map> itemSpecs = iHcFtpFileDAO.queryItemSpec(info);
  324. String taskItems = JSONObject.toJSONString(itemSpecs);
  325. resultMsg = this.createResultMsg("0000", "成功", "{\"U_OR_D\":\"" + tample + "\",\"TASK_ITEMS\":" + taskItems + "}");
  326. return resultMsg;
  327. }
  328. /**
  329. * 根据TaskId 获取任务属性
  330. *
  331. * @return
  332. */
  333. public JSONObject queryTaskAttrs(HttpServletRequest request) {
  334. // 请求参数为{"taskId":"12"}
  335. String ftpItemJson = request.getParameter("ftpItemJson");
  336. if (logger.isDebugEnabled()) {
  337. logger.debug("---【PrvncFtpToFileSystemConfigAction.queryTaskAttrs】入参为:" + ftpItemJson, ftpItemJson);
  338. }
  339. JSONObject paramIn = null;
  340. try {
  341. // 校验格式是否正确
  342. paramIn = JSONObject.parseObject(ftpItemJson);
  343. } catch (Exception e) {
  344. logger.error("传入参数格式不正确:" + ftpItemJson, e);
  345. resultMsg = createResultMsg("1999", "传入参数格式不正确:" + ftpItemJson, "");
  346. return resultMsg;
  347. }
  348. long taskId = paramIn.getLong("taskId");
  349. Map info = new HashMap();
  350. info.put("taskId", taskId);
  351. List<Map> itemAttrs = iHcFtpFileDAO.queryFtpItemAttrsByTaskId(info);
  352. String itemsAttrs = JSONObject.toJSONString(itemAttrs);
  353. resultMsg = this.createResultMsg("0000", "成功", "{\"TASK_ATTRS\":" + itemsAttrs + "}");
  354. return resultMsg;
  355. }
  356. /**
  357. * 启动侦听(多节点启动)
  358. *
  359. * @return
  360. */
  361. public String startJob(HttpServletRequest request) {
  362. // 请求参数为{"tasks":[{"taskId":1},{"taskId":2}],"state":"START"}
  363. String ftpItemJson = request.getParameter("ftpItemJson");
  364. if (logger.isDebugEnabled()) {
  365. logger.debug("---【PrvncFtpToFileSystemConfigAction.startJob】入参为:" + ftpItemJson, ftpItemJson);
  366. }
  367. JSONObject paramIn = null;
  368. try {
  369. // 校验格式是否正确
  370. paramIn = JSONObject.parseObject(ftpItemJson);
  371. } catch (Exception e) {
  372. logger.error("传入参数格式不正确:" + ftpItemJson, e);
  373. resultMsg = createResultMsg("1999", "传入参数格式不正确:" + ftpItemJson, "");
  374. return "startJob";
  375. }
  376. // 传入报文不为空
  377. if (paramIn == null || !paramIn.containsKey("tasks") || !paramIn.containsKey("state")) {
  378. resultMsg = createResultMsg("1999", "传入参数格式不正确(必须包含tasks 和 state节点):" + ftpItemJson, "");
  379. return "startJob";
  380. }
  381. // 校验当前是否为启动侦听
  382. if (!"START".equals(paramIn.get("state"))) {
  383. resultMsg = createResultMsg("1999", "传入参数格式不正确(state的值必须是START):" + ftpItemJson, "");
  384. return "startJob";
  385. }
  386. // 查询需要操作的任务
  387. JSONArray taskInfos = paramIn.getJSONArray("tasks");
  388. String taskIds = "";
  389. for (int taskIndex = 0; taskIndex < taskInfos.size(); taskIndex++) {
  390. taskIds += (taskInfos.getJSONObject(taskIndex).getString("taskId") + ",");
  391. }
  392. if (taskIds.length() > 0) {
  393. taskIds = taskIds.substring(0, taskIds.length() - 1);
  394. }
  395. Map info = new HashMap();
  396. info.put("taskIds", taskIds.split(","));
  397. List<Map> doFtpItems = iHcFtpFileDAO.queryFtpItemsByTaskIds(info);
  398. // 获取Spring调度器
  399. Scheduler scheduler = (Scheduler) SpringBeanInvoker.getBean("schedulerFactoryBean");
  400. int linstenCount = 0;
  401. int updateTaskStateFailCount = 0;
  402. try {
  403. for (int doIndex = 0; doIndex < doFtpItems.size(); doIndex++) {
  404. Map doFtpItem = doFtpItems.get(doIndex);
  405. // 侦听运行状态
  406. String runState = doFtpItem.get("RUN_STATE") == null ? "1" : doFtpItem.get("RUN_STATE").toString();// 如果为空,为了操作正常,默认为1,也就是启动状态,在后面再关闭一次
  407. // 获取taskId
  408. String taskId = doFtpItem.get("TASKID").toString();// 这个就不用三目判断了,因为如果为空,则直接抛出异常,正常情况下不会出现空的情况
  409. // 获取定时时间
  410. String cronExpression = doFtpItem.get("TASKCRON") == null ? defaultCronExpression : doFtpItem.get("TASKCRON").toString();// 如果没有配置则,每一分运行一次
  411. // 设置触发时间点
  412. CronScheduleBuilder cronScheduleBuilder =CronScheduleBuilder.cronSchedule(cronExpression);
  413. String jobName = prefixJobName + taskId;
  414. String triggerName = prefixJobName + taskId;
  415. //设置任务名称
  416. JobKey jobKey = new JobKey(jobName);
  417. JobDetail jobDetail = scheduler.getJobDetail(jobKey); // 说明这个没有启动,则需要重新启动,如果启动着不做处理
  418. if (jobDetail == null) {
  419. // 任务名称
  420. String taskCfgName = (String) doFtpItem.get("TASKNAME");
  421. //构建job信息
  422. JobDetail warnJob = JobBuilder.newJob(HcFtpToFileSystemJob.class).withIdentity(jobName, HcFtpToFileSystemJob.JOB_GROUP_NAME).withDescription("任务启动").build();
  423. warnJob.getJobDataMap().put(HcFtpToFileSystemJob.JOB_DATA_CONFIG_NAME, taskCfgName);
  424. warnJob.getJobDataMap().put(HcFtpToFileSystemJob.JOB_DATA_TASK_ID, taskId);
  425. // 触发时间点
  426. CronTrigger warnTrigger = TriggerBuilder.newTrigger().withIdentity(triggerName, triggerName).withSchedule(cronScheduleBuilder).build();
  427. // 错过执行后,立即执行
  428. //warnTrigger(CronTrigger.MISFIRE_INSTRUCTION_FIRE_ONCE_NOW);
  429. //交由Scheduler安排触发
  430. scheduler.scheduleJob(warnJob, warnTrigger);
  431. // 修改数据状态,将任务数据状态改为运行状态
  432. Map updateTaskInfo = new HashMap();
  433. updateTaskInfo.put("taskId", taskId);
  434. updateTaskInfo.put("runFlag", RUNFLAG_START);
  435. // 这里更新状态没有成功的,只是在后台打印日志,再前台不进行展示
  436. int updateTaskStateFlag = iHcFtpFileDAO.updateFtpItemByTaskId(updateTaskInfo);
  437. if (updateTaskStateFlag < 1) {
  438. logger.error("---侦听【" + taskId + "】启动成功,但是更新任务状态失败,请关注!!!", info);
  439. updateTaskStateFailCount++;
  440. }
  441. }
  442. linstenCount++;
  443. }
  444. String resultMsgStr = "侦听启动成功,启动目标侦听个数为【" + doFtpItems.size() + "】,成功启动个数为【" + linstenCount + "】";
  445. if (updateTaskStateFailCount > 0) {
  446. resultMsgStr += (",有【" + updateTaskStateFailCount + "】任务启动成功,更新数据失败");
  447. }
  448. resultMsg = this.createResultMsg("0000", resultMsgStr, taskIds);
  449. } catch (Exception e) {
  450. // TODO Auto-generated catch block
  451. logger.error("调度器启动出错:" + ftpItemJson, e);
  452. resultMsg = createResultMsg("1999", "调度器启动出错:" + e, "");
  453. return "startJob";
  454. }
  455. if (logger.isDebugEnabled()) {
  456. logger.debug("---【PrvncFtpToFileSystemConfigAction.startJob】出参为:" + resultMsg, resultMsg);
  457. }
  458. return "startJob";
  459. }
  460. /**
  461. * 停止侦听
  462. *
  463. * @return
  464. */
  465. public String stopJob(HttpServletRequest request) {
  466. // 请求参数为{"tasks":[{"taskId":1},{"taskId":2}],"state":"STOP"}
  467. String ftpItemJson = request.getParameter("ftpItemJson");
  468. if (logger.isDebugEnabled()) {
  469. logger.debug("---【PrvncFtpToFileSystemConfigAction.stopJob】入参为:" + ftpItemJson, ftpItemJson);
  470. }
  471. JSONObject paramIn = null;
  472. try {
  473. // 校验格式是否正确
  474. paramIn = JSONObject.parseObject(ftpItemJson);
  475. } catch (Exception e) {
  476. logger.error("传入参数格式不正确:" + ftpItemJson, e);
  477. resultMsg = createResultMsg("1999", "传入参数格式不正确:" + ftpItemJson, "");
  478. return "stopJob";
  479. }
  480. // 传入报文不为空
  481. if (paramIn == null || !paramIn.containsKey("tasks") || !paramIn.containsKey("state")) {
  482. resultMsg = createResultMsg("1999", "传入参数格式不正确(必须包含tasks 和 state节点):" + ftpItemJson, "");
  483. return "stopJob";
  484. }
  485. // 校验当前是否为启动侦听
  486. if (!"STOP".equals(paramIn.get("state"))) {
  487. resultMsg = createResultMsg("1999", "传入参数格式不正确(state的值必须是START):" + ftpItemJson, "");
  488. return "stopJob";
  489. }
  490. // 查询需要操作的任务
  491. JSONArray taskInfos = paramIn.getJSONArray("tasks");
  492. String taskIds = "";
  493. for (int taskIndex = 0; taskIndex < taskInfos.size(); taskIndex++) {
  494. taskIds += (taskInfos.getJSONObject(taskIndex).getString("taskId") + ",");
  495. }
  496. if (taskIds.length() > 0) {
  497. taskIds = taskIds.substring(0, taskIds.length() - 1);
  498. }
  499. Map info = new HashMap();
  500. info.put("taskIds", taskIds.split(","));
  501. List<Map> doFtpItems = iHcFtpFileDAO.queryFtpItemsByTaskIds(info);
  502. // 获取Spring调度器
  503. Scheduler scheduler = (Scheduler) SpringBeanInvoker.getBean("schedulerFactoryBean");
  504. int linstenCount = 0;
  505. int updateTaskStateFailCount = 0;
  506. try {
  507. for (Map doFtpItem : doFtpItems) {
  508. // 获取taskId
  509. String taskId = doFtpItem.get("TASKID").toString();// 这个就不用三目判断了,因为如果为空,则直接抛出异常,正常情况下不会出现空的情况
  510. // 获取定时时间
  511. String cronExpression = doFtpItem.get("TASKCRON") == null ? defaultCronExpression : doFtpItem.get("TASKCRON").toString();// 如果没有配置则,每一分运行一次
  512. String jobName = prefixJobName + taskId;
  513. String triggerName = prefixJobName + taskId;
  514. TriggerKey triggerKey = TriggerKey.triggerKey(jobName, HcFtpToFileSystemJob.JOB_GROUP_NAME);
  515. // 停止触发器
  516. scheduler.pauseTrigger(triggerKey);
  517. // 移除触发器
  518. scheduler.unscheduleJob(triggerKey);
  519. JobKey jobKey = new JobKey(jobName, HcFtpToFileSystemJob.JOB_GROUP_NAME);
  520. // 删除任务
  521. scheduler.deleteJob(jobKey);
  522. // 修改数据状态,将任务数据状态改为运行状态
  523. Map updateTaskInfo = new HashMap();
  524. updateTaskInfo.put("taskId", taskId);
  525. updateTaskInfo.put("runFlag", RUNFLAG_STOP);
  526. // 这里更新状态没有成功的,只是在后台打印日志,再前台不进行展示
  527. int updateTaskStateFlag = iHcFtpFileDAO.updateFtpItemByTaskId(updateTaskInfo);
  528. if (updateTaskStateFlag < 1) {
  529. logger.error("---侦听【" + taskId + "】停止成功,但是更新任务状态失败,请关注!!!", info);
  530. updateTaskStateFailCount++;
  531. }
  532. linstenCount++;
  533. String resultMsgStr = "侦听停止成功,停止目标侦听个数为【" + doFtpItems.size() + "】,成功停止个数为【" + linstenCount + "】";
  534. if (updateTaskStateFailCount > 0) {
  535. resultMsgStr += (",有【" + updateTaskStateFailCount + "】任务停止成功,更新数据失败");
  536. }
  537. resultMsg = this.createResultMsg("0000", resultMsgStr, taskIds);
  538. }
  539. } catch (Exception e) {
  540. // TODO Auto-generated catch block
  541. logger.error("调度器停止出错:" + ftpItemJson, e);
  542. resultMsg = createResultMsg("1999", "调度器停止出错:" + e, "");
  543. return "stopJob";
  544. }
  545. if (logger.isDebugEnabled()) {
  546. logger.debug("---【PrvncFtpToFileSystemConfigAction.startJob】出参为:" + resultMsg, resultMsg);
  547. }
  548. return "stopJob";
  549. }
  550. /**
  551. * 根据任务名称或任务ID模糊查询
  552. *
  553. * @return
  554. */
  555. public String searchTaskByNameOrId(HttpServletRequest request) {
  556. String ftpItemJson = request.getParameter("ftpItemJson");
  557. JSONObject ftpItemJsonObj = null;
  558. try {
  559. // 校验格式是否正确
  560. // {"taskName":"经办人照片同步处理"}
  561. ftpItemJsonObj = JSONObject.parseObject(ftpItemJson);
  562. } catch (Exception e) {
  563. logger.error("传入参数格式不正确:" + ftpItemJson, e);
  564. resultMsg = createResultMsg("1999", "传入参数格式不正确:" + ftpItemJson, "");
  565. return "addFtpItem";
  566. }
  567. // 将ftpItemJson装为Map保存操作
  568. Map paramIn = JSONObject.parseObject(ftpItemJsonObj.toJSONString(), Map.class);
  569. String taskNameOrTaskId = paramIn.get("taskName") == null ? "1" : paramIn.get("taskName").toString();
  570. taskNameOrTaskId = ValidatorUtils.getValueAsString(paramIn, "taskName");
  571. // 规则校验
  572. JSONObject data = new JSONObject();
  573. data.put("total", 1); // 搜索不进行分页处理
  574. data.put("currentPage", 1);
  575. List<Map> ftpItems = null;
  576. // 说明是taskId
  577. if (GenericValidator.isInt(taskNameOrTaskId) || GenericValidator.isLong(taskNameOrTaskId)) {
  578. // 根据taskId 查询记录
  579. paramIn.put("taskId", taskNameOrTaskId);
  580. Map ftpItem = iHcFtpFileDAO.queryFtpItemByTaskId(paramIn);
  581. if (ftpItem != null && ftpItem.containsKey("FTP_ITEM_ATTRS")) {
  582. ftpItem.remove("FTP_ITEM_ATTRS");// 前台暂时用不到,所以这里将属性移除
  583. ftpItems = new ArrayList<Map>();
  584. ftpItems.add(ftpItem);
  585. }
  586. } else {
  587. ftpItems = iHcFtpFileDAO.searchFtpItemByTaskName(paramIn);
  588. }
  589. JSONArray rows = new JSONArray();
  590. if (ftpItems != null && ftpItems.size() > 0) {
  591. DateFormat df = new SimpleDateFormat("yyyy-MM-dd");
  592. for (Map ftpItemMap : ftpItems) {
  593. // 处理时间显示和界面显示传输类型
  594. ftpItemMap.put("U_OR_D_NAME", CustomizedPropertyPlaceholderConfigurer.getContextProperty("task.tamplete.name." + ftpItemMap.get("U_OR_D")));// 暂且写死,最终还是读取配置
  595. ftpItemMap.put("CREATE_DATE", df.format(ftpItemMap.get("CREATE_DATE")));// 暂且写死,最终还是读取配置
  596. rows.add(JSONObject.parseObject(JSONObject.toJSONString(ftpItemMap)));
  597. }
  598. }
  599. data.put("rows", rows);
  600. resultMsg = data;
  601. return "searchTaskByNameOrId";
  602. }
  603. /**
  604. * 创建公用输出
  605. *
  606. * @return
  607. */
  608. private JSONObject createResultMsg(String resultCode, String resultMsg, String resultInfo) {
  609. JSONObject data = new JSONObject();
  610. data.put("RESULT_CODE", resultCode);
  611. data.put("RESULT_MSG", resultMsg);
  612. data.put("RESULT_INFO", resultInfo);
  613. return data;
  614. }
  615. }