9 changed files with 158 additions and 115 deletions
@ -1,30 +1,51 @@ |
|||||||
package com.logpm.factory.receiver; |
package com.logpm.factory.receiver; |
||||||
|
|
||||||
|
import com.logpm.factory.snm.dto.OrderInfoDTO; |
||||||
|
import com.logpm.factory.snm.service.IPanFactoryDataService; |
||||||
|
import com.rabbitmq.client.Channel; |
||||||
|
import lombok.extern.slf4j.Slf4j; |
||||||
|
import org.springblade.common.constant.RabbitConstant; |
||||||
|
import org.springframework.amqp.core.Message; |
||||||
|
import org.springframework.amqp.rabbit.annotation.RabbitHandler; |
||||||
|
import org.springframework.amqp.rabbit.annotation.RabbitListener; |
||||||
|
import org.springframework.beans.factory.annotation.Autowired; |
||||||
|
import org.springframework.stereotype.Component; |
||||||
|
|
||||||
|
import java.io.IOException; |
||||||
|
import java.util.Map; |
||||||
|
|
||||||
|
|
||||||
/** |
/** |
||||||
* 直接队列1 处理器 |
* 直接队列1 处理器 |
||||||
* |
* |
||||||
* @author yangkai.shen |
* @author yangkai.shen |
||||||
*/ |
*/ |
||||||
//@Slf4j
|
@Slf4j |
||||||
//@RabbitListener(queues = RabbitConstant.DIRECT_MODE_QUEUE_ONE)
|
@RabbitListener(queues = RabbitConstant.ORDER_STATUS_QUEUE) |
||||||
//@Component
|
@Component |
||||||
public class DirectQueueOneHandler { |
public class DirectQueueOneHandler { |
||||||
|
|
||||||
// @RabbitHandler
|
@Autowired |
||||||
// public void directHandlerManualAck(PanFactoryOrderDTO factoryOrderDTO, Message message, Channel channel) {
|
private IPanFactoryDataService panFactoryDataService; |
||||||
// // 如果手动ACK,消息会被监听消费,但是消息在队列中依旧存在,如果 未配置 acknowledge-mode 默认是会在消费完毕后自动ACK掉
|
|
||||||
// final long deliveryTag = message.getMessageProperties().getDeliveryTag();
|
@RabbitHandler |
||||||
// try {
|
public void directHandlerManualAck(Map map, Message message, Channel channel) { |
||||||
// log.info("直接队列1,手动ACK,接收消息:{}", factoryOrderDTO);
|
// 如果手动ACK,消息会被监听消费,但是消息在队列中依旧存在,如果 未配置 acknowledge-mode 默认是会在消费完毕后自动ACK掉
|
||||||
// //通知 MQ 消息已被成功消费,可以ACK了
|
final long deliveryTag = message.getMessageProperties().getDeliveryTag(); |
||||||
// channel.basicAck(deliveryTag, false);
|
try { |
||||||
// } catch (IOException e) {
|
log.info("直接队列1,手动ACK,接收消息:{}", map); |
||||||
// try {
|
OrderInfoDTO orderInfoDTO = (OrderInfoDTO) map.get("messageData"); |
||||||
// // 处理失败,重新压入MQ
|
|
||||||
// channel.basicRecover();
|
panFactoryDataService.handleData(orderInfoDTO); |
||||||
// } catch (IOException e1) {
|
|
||||||
// e1.printStackTrace();
|
channel.basicAck(deliveryTag, false); |
||||||
// }
|
} catch (IOException e) { |
||||||
// }
|
try { |
||||||
// }
|
// 处理失败,重新压入MQ
|
||||||
|
channel.basicRecover(); |
||||||
|
} catch (IOException e1) { |
||||||
|
e1.printStackTrace(); |
||||||
|
} |
||||||
|
} |
||||||
|
} |
||||||
} |
} |
||||||
|
Loading…
Reference in new issue