一.消息生成者
1.1消息生成者配置
1.2 消息发送端代码
1.3 创建交换机,队列,并建立关系

二.消费者
2.1消费者
:::success
@RabbitListener与@RabbitHandler
:::
1.@RabbitListener 注解是指定某方法作为消息消费的方法,例如监听某 Queue 里面的消息。
2.@RabbitListener标注在方法上,直接监听指定的队列,此时接收的参数需要与发送市类型一致
@Componentpublic class PointConsumer {//监听的队列名@RabbitListener(queues = "point.to.point")public void processOne(String name) {System.out.println("point.to.point:" + name);}}
3.@RabbitListener 可以标注在类上面,需配合 @RabbitHandler 注解一起使用
@RabbitListener 标注在类上面表示当有收到消息的时候,就交给 @RabbitHandler 的方法处理,根据接受的参数类型进入具体的方法中。
@Component@RabbitListener(queues = "consumer_queue")public class Receiver {@RabbitHandlerpublic void processMessage1(String message) {System.out.println(message);}@RabbitHandlerpublic void processMessage2(byte[] message) {System.out.println(new String(message));}}
三.限流配置
3.1配置文件
#在单个请求中处理的消息个数,他应该大于等于事务数量(unack的最大数量) spring.rabbitmq.listener.simple.prefetch=2 #在@RabbitListener(queues = { HighDeviceMessage.QUEUE_NAME },concurrency = “${spring.rabbitmq.highdevice.concurrency}”)配置的占位符配置 spring.rabbitmq.highdevice.concurrency=2-5
3.2消费者配置
@Componentpublic class HighDeviceMessageHandler {// @RabbitListener(queues = { HighDeviceMessage.QUEUE_NAME },ackMode ="MANUAL",concurrency = "1-10")@RabbitListener(queues = { HighDeviceMessage.QUEUE_NAME },ackMode ="MANUAL",concurrency = "${spring.rabbitmq.highdevice.concurrency}")public void handle(String msgStr, Channel channel,@Headers Map<String,Object> headers) {try {log.info("休息3秒");log.info("handle HighDeviceMessage:{}---",msgStr);TimeUnit.SECONDS.sleep(3);long deliveryTag = (Long)headers.get(AmqpHeaders.DELIVERY_TAG);//手工ackchannel.basicAck(deliveryTag,true);} catch (Exception e) {log.error("handle HighDeviceMessage err", e);}}}
