修改xml消息在分发前一次转成MSG对象,避免后面业务反复做不必要的转换操作

This commit is contained in:
zhouxiunai
2019-01-08 15:44:59 +08:00
parent 3c7e55c2cc
commit f7fd218f93
11 changed files with 71 additions and 40 deletions
@@ -4,6 +4,7 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.gzzn.omms.msgexchangeapi.entity.Cminmsg;
import com.gzzn.omms.msgexchangeapi.entity.CminmsgWapper;
import com.gzzn.omms.msgexchangeapi.entity.msg.MSG;
import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD.FLTR;
import com.gzzn.omms.msgexchangeapi.msghandler.HandlerResult;
@@ -18,8 +19,8 @@ public class ADFTHandler extends SchdBaseHandler{
private static Logger logger = LoggerFactory.getLogger(ADFTHandler.class);
@Override
public HandlerResult run(Cminmsg cminmsg) {
MSG msg = exchangeService.xmlToMsg(cminmsg.getCminmsgsClobMsg());
public HandlerResult run(CminmsgWapper cminmsg) {
MSG msg = cminmsg.getMsg();
if(msg.getSCHD().getFLTR().size() != 1)
{
logger.error("预期每条消息只会包含一个动态航班,当前{}",msg.getSCHD().getFLTR().size());
@@ -3,7 +3,7 @@ package com.gzzn.omms.msgexchangeapi.msghandler.schd;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.gzzn.omms.msgexchangeapi.entity.Cminmsg;
import com.gzzn.omms.msgexchangeapi.entity.CminmsgWapper;
import com.gzzn.omms.msgexchangeapi.msghandler.HandlerResult;
import com.gzzn.omms.msgexchangeapi.msghandler.schd.base.SchdBaseHandler;
@@ -16,7 +16,7 @@ public class DNLDHandler extends SchdBaseHandler{
private static Logger logger = LoggerFactory.getLogger(DNLDHandler.class);
@Override
public HandlerResult run(Cminmsg cminmsg) {
public HandlerResult run(CminmsgWapper cminmsg) {
Boolean updateResult = flightInfoService.updateByDaySchd(cminmsg);
if(false == updateResult)
{
@@ -3,7 +3,7 @@ package com.gzzn.omms.msgexchangeapi.msghandler.schd;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.gzzn.omms.msgexchangeapi.entity.Cminmsg;
import com.gzzn.omms.msgexchangeapi.entity.CminmsgWapper;
import com.gzzn.omms.msgexchangeapi.msghandler.HandlerResult;
import com.gzzn.omms.msgexchangeapi.msghandler.schd.base.SchdBaseHandler;
@@ -16,7 +16,7 @@ public class RESPHandler extends SchdBaseHandler {
private static Logger logger = LoggerFactory.getLogger(RESPHandler.class);
@Override
public HandlerResult run(Cminmsg cminmsg) {
public HandlerResult run(CminmsgWapper cminmsg) {
Boolean updateResult = flightInfoService.updateByDaySchd(cminmsg);
if(false == updateResult)
{
@@ -1,16 +1,15 @@
package com.gzzn.omms.msgexchangeapi.msghandler.schd.base;
import java.util.Arrays;
import java.util.List;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.gzzn.omms.msgexchangeapi.dto.ResponseDto;
import com.gzzn.omms.msgexchangeapi.entity.Cminmsg;
import com.gzzn.omms.msgexchangeapi.entity.CminmsgWapper;
import com.gzzn.omms.msgexchangeapi.entity.msg.MSG;
import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD;
import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD.FLTR;
import com.gzzn.omms.msgexchangeapi.msghandler.HandlerResult;
import com.gzzn.omms.msgexchangeapi.msghandler.IBaseHandler;
import com.gzzn.omms.msgexchangeapi.service.ICminmsgService;
@@ -42,7 +41,7 @@ public class SchdBaseHandler implements IBaseHandler{
}
@Override
public HandlerResult run(Cminmsg cminmsg) {
public HandlerResult run(CminmsgWapper cminmsg) {
return null;
}
@@ -50,11 +49,10 @@ public class SchdBaseHandler implements IBaseHandler{
* 发送日计划到kafka队列
* @param cminmsg
*/
protected void sendDschd(Cminmsg cminmsg)
protected void sendDschd(CminmsgWapper cminmsg)
{
//发送日航班消息的通知给前端,让前端重新加载全部动态
String clobMsg = cminmsg.getCminmsgsClobMsg();
MSG dnldMsg = exchangeService.xmlToMsg(clobMsg);
MSG dnldMsg = cminmsg.getMsg();
int recs = dnldMsg.getSCHD().getRECS();
logger.info("发送日计划到kafka,条数:{}",recs);
@@ -70,12 +68,10 @@ public class SchdBaseHandler implements IBaseHandler{
* 仅发送动态消息
* @param cminmsg
*/
protected void sendMsg(Cminmsg cminmsg)
protected void sendMsg(CminmsgWapper cminmsg)
{
//send msg
String msg = exchangeService.xmlMsgToJsonMsg(
cminmsg.getCminmsgsClobMsg()
);
String msg = exchangeService.msgToJson(cminmsg.getMsg());
kafkaservice.msgSend("msg", msg);
}
@@ -85,16 +81,14 @@ public class SchdBaseHandler implements IBaseHandler{
* @param fltr
* @param cminmsg
*/
protected void sendFltrAndMsg(SCHD.FLTR fltr,Cminmsg cminmsg)
protected void sendFltrAndMsg(SCHD.FLTR fltr,CminmsgWapper cminmsg)
{
//send schd
String schdMsg = JsonUtil.getString(fltr);
kafkaservice.msgSend("schd", schdMsg);
//send msg
String msg = exchangeService.xmlMsgToJsonMsg(
cminmsg.getCminmsgsClobMsg()
);
String msg = exchangeService.msgToJson(cminmsg.getMsg());
kafkaservice.msgSend("msg", msg);
}