消息队列

This commit is contained in:
yuezhihang
2022-07-05 11:19:49 +08:00
parent 07303047fa
commit 667b3efff4
@@ -18,6 +18,7 @@ import org.springframework.amqp.rabbit.annotation.Exchange;
import org.springframework.amqp.rabbit.annotation.Queue; import org.springframework.amqp.rabbit.annotation.Queue;
import org.springframework.amqp.rabbit.annotation.QueueBinding; import org.springframework.amqp.rabbit.annotation.QueueBinding;
import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.amqp.support.AmqpHeaders;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.Message; import org.springframework.messaging.Message;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
@@ -54,6 +55,7 @@ public class SendQdUpdateMQService {
key = "qd-update-key_SQ_GSAR_RELEASE")) key = "qd-update-key_SQ_GSAR_RELEASE"))
public void createMQ(Map<String, Object> bussMap, Message message, Channel channel) throws Exception { public void createMQ(Map<String, Object> bussMap, Message message, Channel channel) throws Exception {
try {
//数据字典数据组装 //数据字典数据组装
Map<String, Object> dicTypeEO = dicTypeEOService.getDicTypeListCode(); Map<String, Object> dicTypeEO = dicTypeEOService.getDicTypeListCode();
Map<String, String> dict = new HashMap<>(); Map<String, String> dict = new HashMap<>();
@@ -200,6 +202,13 @@ public class SendQdUpdateMQService {
} }
} }
}); });
}catch(Exception e){
logger.error(e.getMessage());
}finally {
Long tag = (Long) message.getHeaders().get(AmqpHeaders.DELIVERY_TAG);
channel.basicAck(tag,false);
logger.debug("消息确认成功!!!!!!!!!");
}
} }
@@ -210,6 +219,7 @@ public class SendQdUpdateMQService {
key = "standard-key_updateText_RELEASE")) key = "standard-key_updateText_RELEASE"))
public void createStandardMQ(Map<String, Object> bussMap, Message message, Channel channel) throws Exception { public void createStandardMQ(Map<String, Object> bussMap, Message message, Channel channel) throws Exception {
try {
//数据字典数据组装 //数据字典数据组装
//数据字典数据组装 //数据字典数据组装
Map<String, Object> dicTypeEO = dicTypeEOService.getDicTypeListCode(); Map<String, Object> dicTypeEO = dicTypeEOService.getDicTypeListCode();
@@ -224,7 +234,6 @@ public class SendQdUpdateMQService {
List<DSarStandardDTO> changeXXYX = (List<DSarStandardDTO>) bussMap.get("changeXXYX"); List<DSarStandardDTO> changeXXYX = (List<DSarStandardDTO>) bussMap.get("changeXXYX");
//写修改逻辑 //写修改逻辑
Thread.sleep(5000); Thread.sleep(5000);
for (DSarStandardDTO jjss : changeJJSS) { for (DSarStandardDTO jjss : changeJJSS) {
@@ -269,9 +278,13 @@ public class SendQdUpdateMQService {
} }
} }
} }
}catch(Exception e){
logger.error(e.getMessage());
}finally {
Long tag = (Long) message.getHeaders().get(AmqpHeaders.DELIVERY_TAG);
channel.basicAck(tag,false);
logger.debug("消息确认成功!!!!!!!!!");
}
} }