添加计划外航班消息处理
This commit is contained in:
@@ -0,0 +1,38 @@
|
||||
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.msg.MSG;
|
||||
import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD.FLTR;
|
||||
import com.gzzn.omms.msgexchangeapi.msghandler.HandlerResult;
|
||||
import com.gzzn.omms.msgexchangeapi.msghandler.schd.base.SchdBaseHandler;
|
||||
|
||||
/**
|
||||
* 添加计划外航班计划
|
||||
* @author zhouxiunai
|
||||
*
|
||||
*/
|
||||
public class ADFTHandler extends SchdBaseHandler{
|
||||
private static Logger logger = LoggerFactory.getLogger(ADFTHandler.class);
|
||||
|
||||
@Override
|
||||
public HandlerResult run(Cminmsg cminmsg) {
|
||||
MSG msg = exchangeService.xmlstrToObject(cminmsg.getCminmsgsClobMsg(), MSG.class);
|
||||
if(msg.getSCHD().getFLTR().size() != 1)
|
||||
{
|
||||
logger.error("预期每条消息只会包含一个动态航班,当前{}",msg.getSCHD().getFLTR().size());
|
||||
return HandlerResult.failure();
|
||||
}
|
||||
|
||||
//每条消息只会包含一个动态航班
|
||||
FLTR fltr = msg.getSCHD().getFLTR().get(0);
|
||||
flightInfoService.saveFltr(fltr);
|
||||
|
||||
this.sendFltrAndMsg(fltr, cminmsg);
|
||||
this.updateCminmsgToProcessed(cminmsg);
|
||||
|
||||
return HandlerResult.success();
|
||||
} //end function
|
||||
}
|
||||
@@ -3,12 +3,14 @@ package com.gzzn.omms.msgexchangeapi.msghandler.schd.base;
|
||||
import java.util.Arrays;
|
||||
|
||||
import com.gzzn.omms.msgexchangeapi.entity.Cminmsg;
|
||||
import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD;
|
||||
import com.gzzn.omms.msgexchangeapi.msghandler.HandlerResult;
|
||||
import com.gzzn.omms.msgexchangeapi.msghandler.IBaseHandler;
|
||||
import com.gzzn.omms.msgexchangeapi.service.ICminmsgService;
|
||||
import com.gzzn.omms.msgexchangeapi.service.IKafkaService;
|
||||
import com.gzzn.omms.msgexchangeapi.service.exchange.IExchangeService;
|
||||
import com.gzzn.omms.msgexchangeapi.service.flightInfo.IFlightInfoService;
|
||||
import com.gzzn.omms.msgexchangeapi.utils.JsonUtil;
|
||||
import com.gzzn.omms.msgexchangeapi.utils.SpringUtil;
|
||||
|
||||
/**
|
||||
@@ -47,6 +49,25 @@ public class SchdBaseHandler implements IBaseHandler{
|
||||
kafkaservice.msgSend("msg", msg);
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
* 发送航班信息到schd topic ,同时发送动态消息到msg topic
|
||||
* @param fltr
|
||||
* @param cminmsg
|
||||
*/
|
||||
protected void sendFltrAndMsg(SCHD.FLTR fltr,Cminmsg cminmsg)
|
||||
{
|
||||
//send schd
|
||||
String schdMsg = JsonUtil.getString(fltr);
|
||||
kafkaservice.msgSend("schd", schdMsg);
|
||||
|
||||
//send msg
|
||||
String msg = exchangeService.xmlstrToJson(
|
||||
cminmsg.getCminmsgsClobMsg()
|
||||
);
|
||||
kafkaservice.msgSend("msg", msg);
|
||||
}
|
||||
|
||||
/**
|
||||
* 更新Cminmsg状态为已处理
|
||||
* @param cminmsg
|
||||
|
||||
Reference in New Issue
Block a user