见文章:
Linux安装rabbitMq及elrang环境_qq837993702的博客-CSDN博客
RabbitMQ 集群环境搭建_qq837993702的博客-CSDN博客
二、pom文件配置三、配置文件org.springframework.boot spring-boot-starter-parent2.2.6.RELEASE org.springframework.boot spring-boot-starter-amqp
spring.rabbitmq.host= 10.10.240.206 spring.rabbitmq.port= 5672 #集群模式ip:端口用,隔开 spring.rabbitmq.addresses=ip:端口,ip:端口,ip:端口 spring.rabbitmq.username= admin spring.rabbitmq.password= 123456 #集群模式 spring.rabbitmq.mode=cluster #单机模式 #spring.rabbitmq.mode=single ##是否启用 发布确认 spring.rabbitmq.publisher-confirms=true ## 是否启用发布返回 spring.rabbitmq.publisher-returns=true # 表示消息确认方式,其有三种配置方式,分别是none、manual和auto;默认auto spring.rabbitmq.listener.direct.acknowledge-mode=manual spring.rabbitmq.listener.simple.acknowledge-mode=manual四、RabbitConfig
1、RabbitTemplate 配置,实现集群模式和单点模式自由切换
package cn.edu.nfu.jw.conf;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.amqp.core.*;
import org.springframework.amqp.rabbit.config.SimpleRabbitListenerContainerFactory;
import org.springframework.amqp.rabbit.connection.CachingConnectionFactory;
import org.springframework.amqp.rabbit.connection.ConnectionFactory;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
import org.springframework.boot.autoconfigure.amqp.RabbitProperties;
import org.springframework.boot.autoconfigure.amqp.SimpleRabbitListenerContainerFactoryConfigurer;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Scope;
import javax.annotation.Resource;
@Configuration
public class RabbitConfig {
private final Logger logger = LoggerFactory.getLogger(this.getClass());
// @Value("${spring.rabbitmq.host}")
// private String host;
//
// @Value("${spring.rabbitmq.port}")
// private int port;
// @Value("${spring.rabbitmq.username}")
// private String username;
//
// @Value("${spring.rabbitmq.password}")
// private String password;
@Resource
private RabbitProperties rabbitProperties;
@Bean("customContainerFactory")
public SimpleRabbitListenerContainerFactory
containerFactory(SimpleRabbitListenerContainerFactoryConfigurer configurer, ConnectionFactory
connectionFactory) {
SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
factory.setConcurrentConsumers(DEFAULT_CONCURRENT);
factory.setMaxConcurrentConsumers(DEFAULT_CONCURRENT);
configurer.configure(factory, connectionFactory);
return factory;
}
//单机节点
@Bean
@ConditionalOnProperty(name = "spring.rabbitmq.mode", havingValue = "single")
public ConnectionFactory connectionFactorySingle(@Value("${spring.rabbitmq.host}") String host,@Value("${spring.rabbitmq.port}") int port) {
CachingConnectionFactory connectionFactory = new CachingConnectionFactory(rabbitProperties.getHost(),rabbitProperties.getPort());
connectionFactory.setUsername(rabbitProperties.getUsername());
connectionFactory.setPassword(rabbitProperties.getPassword());
connectionFactory.setVirtualHost("/");
connectionFactory.setPublisherConfirms(rabbitProperties.isPublisherConfirms());
connectionFactory.setPublisherReturns(rabbitProperties.isPublisherReturns());
return connectionFactory;
}
//集群节点
@Bean
@ConditionalOnProperty(name = "spring.rabbitmq.mode", havingValue = "cluster")
public ConnectionFactory connectionFactoryCluster(@Value("${spring.rabbitmq.addresses}") String addresses) {
CachingConnectionFactory cachingConnectionFactory = new CachingConnectionFactory();
cachingConnectionFactory.setAddresses(rabbitProperties.getAddresses());
cachingConnectionFactory.setUsername(rabbitProperties.getUsername());
cachingConnectionFactory.setPassword(rabbitProperties.getPassword());
cachingConnectionFactory.setVirtualHost("/");
cachingConnectionFactory.setPublisherConfirms(rabbitProperties.isPublisherConfirms());
cachingConnectionFactory.setPublisherReturns(rabbitProperties.isPublisherReturns());
logger.info("集群连接工厂设置完成,连接地址{}"+rabbitProperties.getAddresses());
logger.info("集群连接工厂设置完成,连接用户{}"+rabbitProperties.getUsername());
return cachingConnectionFactory;
}
@Bean
//必须是prototype类型
@Scope("prototype")
public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) {
RabbitTemplate template = new RabbitTemplate(connectionFactory);
return template;
}
}
2、rabbitmq交换机,队列等信息配置
package cn.edu.nfu.jw.eval.constant;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.BindingBuilder;
import org.springframework.amqp.core.DirectExchange;
import org.springframework.amqp.core.Queue;
import org.springframework.amqp.rabbit.config.SimpleRabbitListenerContainerFactory;
import org.springframework.amqp.rabbit.connection.ConnectionFactory;
import org.springframework.boot.autoconfigure.amqp.SimpleRabbitListenerContainerFactoryConfigurer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class RabbitMqConstant {
private final Logger logger = LoggerFactory.getLogger(this.getClass());
//学生评教写入用的交换机(直连交换机是一种带路由功能的交换机,一个队列会和一个交换机绑定)
public static final String STU_eval_EXCHANGE_DIRECT = "stu-eval-mq-exchange_A";
public static final String STU_eval_QUEUE_DIRECT = "StuevalDirectQueue";
public static final String STU_eval_ROUTINGKEY_DIRECT = "stu-eval-mq-routingKey_A";
public static final int DEFAULT_ConCURRENT = 10;//多线程消费者数量,默认10
@Bean("stuevalContainerFactory")
public SimpleRabbitListenerContainerFactory
containerFactory(SimpleRabbitListenerContainerFactoryConfigurer configurer, ConnectionFactory
connectionFactory) {
SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
factory.setConcurrentConsumers(DEFAULT_CONCURRENT);
factory.setMaxConcurrentConsumers(DEFAULT_CONCURRENT);
configurer.configure(factory, connectionFactory);
return factory;
}
@Bean
public DirectExchange stuevalExchange() {
return new DirectExchange(STU_eval_EXCHANGE_DIRECT);
}
@Bean
public Queue queueStueval() {
return new Queue(STU_eval_QUEUE_DIRECT, true); //队列持久
}
//一个交换机可以绑定多个消息队列,也就是消息通过一个交换机,可以分发到不同的队列当中去。
@Bean
public Binding bindingStueval() {
return BindingBuilder.bind(queueStueval()).to(stuevalExchange()).with(STU_eval_ROUTINGKEY_DIRECT);
}
}
五、生产者实现(带发布确认)
Map六、消费者实现(实现消息手动确认)mqMessage = new HashMap<>(); mqMessage.put("saveOrUpdate", "update"); //写数据这里改为发送mq mqMessage.put("evalJxbUser", evalJxbUser); mqMessage.put("evalJxbUserRecord", obj); // 消息返回, yml需要配置 publisher-returns: true rabbitTemplate.setReturnCallback((message, replyCode, replyText, exchange, routingKey) -> { String correlationId = message.getMessageProperties().getCorrelationIdString(); log.info("学生评教发送mq失败:"); throw new ServiceException(-1, "评教失败,消息发送失败"); }); // 消息确认, 需要配置 publisher-/confirm/is: true rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> { if (!ack) { log.info("stuevalDirectSender生产者消息发送失败" + cause + correlationData.toString()); throw new ServiceException(-1, "评教失败,消息发送失败"); } else { //发送成功后标记 valueOps.set("stueval:"+obj.getUserNo()+":"+obj.getevalJxbInfoId(),"true"); log.info("stuevalDirectSender生产者消息发送成功 "); } }); rabbitTemplate.convertAndSend(RabbitMqConstant.STU_eval_EXCHANGE_DIRECT, RabbitMqConstant.STU_eval_ROUTINGKEY_DIRECT, JSON.toJSonString(mqMessage));
package cn.edu.nfu.jw.eval.listener;
import cn.edu.nfu.jw.eval.constant.RabbitMqConstant;
import cn.edu.nfu.jw.eval.dao.evalJxbUserMapper;
import cn.edu.nfu.jw.eval.dao.evalJxbUserRecordMapper;
import cn.edu.nfu.jw.eval.dao.evalJxbUserRecordTargetMapper;
import cn.edu.nfu.jw.eval.domain.evalJxbUser;
import cn.edu.nfu.jw.eval.domain.evalJxbUserRecord;
import cn.edu.nfu.jw.eval.domain.evalJxbUserRecordTarget;
import cn.edu.nfu.jw.exception.ServiceException;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import com.rabbitmq.client.Channel;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitHandler;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;
import javax.annotation.Resource;
import java.io.IOException;
import java.util.Map;
@Component
@RabbitListener(queues = RabbitMqConstant.STU_eval_QUEUE_DIRECT
, containerFactory = "stuevalContainerFactory")//监听的队列名称 StuevalDirectQueue
@Transactional(transactionManager = "jwTransactionManager")
public class StuevalDirectReceiver {
private static final Log log = LogFactory.getLog(StuevalDirectReceiver.class);
@Resource
private evalJxbUserRecordMapper dao;
@Resource
private evalJxbUserRecordTargetMapper evalJxbUserRecordTargetMapper;
@Resource
private evalJxbUserMapper evalJxbUserMapper;
@Resource
private RedisTemplate redisTemplate;
@RabbitHandler
public void process(String jsonStr, Channel channel, Message message) throws IOException {
log.info("StuevalDirectReceiver消费者收到消息 : " + jsonStr);
try {
Map mqMessage = JSON.parseObject(jsonStr, Map.class);
Boolean o = false,o2 = false;
evalJxbUser evalJxbUser = JSONObject.parseObject(((JSONObject) mqMessage.get("evalJxbUser")).toJSonString(),evalJxbUser.class);
evalJxbUserRecord obj = JSONObject.parseObject(((JSONObject) mqMessage.get("evalJxbUserRecord")).toJSonString(),evalJxbUserRecord.class);
if (mqMessage.get("saveOrUpdate") != null&&"save".equals(mqMessage.get("saveOrUpdate"))) {
o = evalJxbUserMapper.save(evalJxbUser);
} else {
o = evalJxbUserMapper.update(evalJxbUser);
}
o2 = dao.save(obj);
//保存评教选项记录
for (evalJxbUserRecordTarget evalJxbXsRecordTarget : obj.getevalRecordTargetList()) {
evalJxbXsRecordTarget.setevalRecordId(obj.getId());
evalJxbXsRecordTarget.setUuser(obj.getUuser());
}
evalJxbUserRecordTargetMapper.bitchSave(obj.getevalRecordTargetList());
if (o&&o2) {//消费完成,删除标记
redisTemplate.delete("stueval:"+obj.getUserNo()+":"+obj.getevalJxbInfoId());
log.info(obj.getUserNo() + ":学生评教保存成功!");
channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
log.info("StuevalDirectReceiver消费者消费消息 : " + jsonStr);
} else {
log.error("保存学生评教信息失败!"+obj.getUserNo());
}
} catch (Exception e) {
//消费者处理出了问题,需要告诉队列信息消费失败,失败后重新入队
try {
//nack返回false,并重新回到队列
channel.basicNack(message.getMessageProperties().getDeliveryTag(),
false, true);
} catch (IOException ioException) {//防止内存溢出,这里不再重新入队
log.error("重新放入队列失败,失败原因:"+e.getMessage(),e);
}
log.error("保存学生评教信息,从MQ队列里面获取信息保存异常,异常信息为:" + e.getMessage(), e);
throw new ServiceException(-1, "保存学生评教信息,从MQ队列里面获取信息保存异常,异常信息为:" + e.getMessage());
}
}
}



