parent
b37acbfcab
commit
3f7643f0e8
@ -0,0 +1,94 @@
|
||||
package cn.estsh.i3plus.core.apiservice.mq;
|
||||
|
||||
import cn.estsh.i3plus.core.api.iservice.busi.ISysMessageService;
|
||||
import cn.estsh.i3plus.platform.common.tool.TimeTool;
|
||||
import cn.estsh.i3plus.platform.common.util.PlatformConstWords;
|
||||
import cn.estsh.i3plus.pojo.base.enumutil.CommonEnumUtil;
|
||||
import cn.estsh.i3plus.pojo.base.enumutil.ImppEnumUtil;
|
||||
import cn.estsh.i3plus.pojo.platform.bean.SysMessage;
|
||||
import cn.estsh.i3plus.pojo.platform.bean.SysRefUserMessage;
|
||||
import com.alibaba.fastjson.JSON;
|
||||
import com.alibaba.fastjson.JSONObject;
|
||||
import com.rabbitmq.client.Channel;
|
||||
import org.apache.commons.lang3.StringUtils;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.amqp.core.Message;
|
||||
import org.springframework.amqp.rabbit.annotation.RabbitListener;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* @Description : SWEB通知处理队列
|
||||
* @Reference :
|
||||
* @Author : yunhao
|
||||
* @CreateDate : 2018-11-15 22:15
|
||||
* @Modify:
|
||||
**/
|
||||
@ConditionalOnExpression("'${impp.mq.queue.sweb.notice}' == 'true'")
|
||||
@Component
|
||||
public class MessageSWebNoticeQueueReceiver {
|
||||
private static final Logger LOGGER = LoggerFactory.getLogger(MessageSWebNoticeQueueReceiver.class);
|
||||
|
||||
@Autowired
|
||||
private ISysMessageService sysMessageService;
|
||||
|
||||
/**
|
||||
* SWEB通知处理队列
|
||||
*
|
||||
* @param msg
|
||||
* @param channel
|
||||
* @param message
|
||||
*/
|
||||
@RabbitListener(queues = PlatformConstWords.SWEB_NOTICE_QUEUE)
|
||||
public void processImppMail(SysMessage msg, Channel channel, Message message) {
|
||||
LOGGER.info("【MQ-{}】 数据接收成功:{}", PlatformConstWords.SWEB_NOTICE_QUEUE, msg);
|
||||
try {
|
||||
msg = sysMessageService.insertSysMessage(msg);
|
||||
|
||||
if (StringUtils.isNotBlank(msg.getMessageReceiversId())) {
|
||||
// 获取供应商信息 由string转换为json对象
|
||||
JSONObject userJsonObject = JSON.parseObject(msg.getMessageReceiversId());
|
||||
|
||||
List<SysRefUserMessage> insertList = new ArrayList<>(userJsonObject.size());
|
||||
SysRefUserMessage refUserMessage;
|
||||
|
||||
for (String user : userJsonObject.keySet()) {
|
||||
refUserMessage = new SysRefUserMessage();
|
||||
refUserMessage.setMessageId(msg.getId());
|
||||
refUserMessage.setMessageSoftType(CommonEnumUtil.SOFT_TYPE.SWEB.getValue());
|
||||
refUserMessage.setMessageTitleRdd(msg.getMessageTitle());
|
||||
refUserMessage.setMessageTypeRdd(msg.getMessageType());
|
||||
refUserMessage.setMessageSenderNameRdd(msg.getMessageSenderNameRdd());
|
||||
refUserMessage.setReceiverId(Long.parseLong(user));
|
||||
refUserMessage.setReceiverNameRdd(userJsonObject.get(user).toString());
|
||||
refUserMessage.setMessageStatus(ImppEnumUtil.MESSAGE_STATUS.UNREAD.getValue());
|
||||
refUserMessage.setReceiverTime(TimeTool.getNowTime(true));
|
||||
refUserMessage.setIsUrgent(msg.getIsUrgent());
|
||||
|
||||
insertList.add(refUserMessage);
|
||||
}
|
||||
|
||||
sysMessageService.insertSysRefUserMessage(insertList);
|
||||
}
|
||||
|
||||
// 消息处理完成
|
||||
LOGGER.info("【MQ-{}】站内信{} DeliveryTag:{} 处理成功", PlatformConstWords.IMPP_MESSAGE_LETTER_QUEUE,
|
||||
msg, message.getMessageProperties().getDeliveryTag());
|
||||
channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
|
||||
} catch (IOException e) {
|
||||
try {
|
||||
// 未成功处理,重新发送
|
||||
channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true);
|
||||
} catch (IOException e1) {
|
||||
e1.printStackTrace();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
}
|
Loading…
Reference in New Issue