spring整合消息队列rabbitmq
博客专区 > Fe-Fe 的博客 > 博客详情
spring整合消息队列rabbitmq
Fe-Fe 发表于4年前
spring整合消息队列rabbitmq
  • 发表于 4年前
  • 阅读 58636
  • 收藏 74
  • 点赞 8
  • 评论 22

腾讯云 十分钟定制你的第一个小程序>>>   

spring大家太熟,就不多说了

rabbitmq一个amqp的队列服务实现,具体介绍请参考本文http://lynnkong.iteye.com/blog/1699684

本文侧重介绍如何将rabbitmq整合到项目中

ps:本文只是简单一个整合介绍,属于抛砖引玉,具体实现还需大家深入研究哈..


 

1.首先是生产者配置

<?xml version="1.0" encoding="UTF-8"?>
<beans xmlns="http://www.springframework.org/schema/beans"
       xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
       xmlns:context="http://www.springframework.org/schema/context"
       xmlns:rabbit="http://www.springframework.org/schema/rabbit"
       xsi:schemaLocation="
            http://www.springframework.org/schema/beans
                http://www.springframework.org/schema/beans/spring-beans.xsd
            http://www.springframework.org/schema/context
                http://www.springframework.org/schema/context/spring-context.xsd
            http://www.springframework.org/schema/rabbit
                http://www.springframework.org/schema/rabbit/spring-rabbit-1.0.xsd">

    
   <!-- 连接服务配置  -->
   <rabbit:connection-factory id="connectionFactory" host="localhost" username="guest"
        password="guest" port="5672"  />
        
   <rabbit:admin connection-factory="connectionFactory"/>
   
   <!-- queue 队列声明-->
   <rabbit:queue id="queue_one" durable="true" auto-delete="false" exclusive="false" name="queue_one"/>
   
   
   <!-- exchange queue binging key 绑定 -->
    <rabbit:direct-exchange name="my-mq-exchange" durable="true" auto-delete="false" id="my-mq-exchange">
        <rabbit:bindings>
            <rabbit:binding queue="queue_one" key="queue_one_key"/>
        </rabbit:bindings>
    </rabbit:direct-exchange>
    
    <-- spring amqp默认的是jackson 的一个插件,目的将生产者生产的数据转换为json存入消息队列,由于fastjson的速度快于jackson,这里替换为fastjson的一个实现 -->
    <bean id="jsonMessageConverter"  class="mq.convert.FastJsonMessageConverter"></bean>
    
    <-- spring template声明-->
    <rabbit:template exchange="my-mq-exchange" id="amqpTemplate"  connection-factory="connectionFactory"  message-converter="jsonMessageConverter"/>
</beans>

2.fastjson messageconver插件实现


import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageProperties;
import org.springframework.amqp.support.converter.AbstractMessageConverter;
import org.springframework.amqp.support.converter.MessageConversionException;

import fe.json.FastJson;

public class FastJsonMessageConverter  extends AbstractMessageConverter {
	private static Log log = LogFactory.getLog(FastJsonMessageConverter.class);

	public static final String DEFAULT_CHARSET = "UTF-8";

	private volatile String defaultCharset = DEFAULT_CHARSET;
	
	public FastJsonMessageConverter() {
		super();
		//init();
	}
	
	public void setDefaultCharset(String defaultCharset) {
		this.defaultCharset = (defaultCharset != null) ? defaultCharset
				: DEFAULT_CHARSET;
	}
	
	public Object fromMessage(Message message)
			throws MessageConversionException {
		return null;
	}
	
	public <T> T fromMessage(Message message,T t) {
		String json = "";
		try {
			json = new String(message.getBody(),defaultCharset);
		} catch (UnsupportedEncodingException e) {
			e.printStackTrace();
		}
		return (T) FastJson.fromJson(json, t.getClass());
	}	
	

	protected Message createMessage(Object objectToConvert,
			MessageProperties messageProperties)
			throws MessageConversionException {
		byte[] bytes = null;
		try {
			String jsonString = FastJson.toJson(objectToConvert);
			bytes = jsonString.getBytes(this.defaultCharset);
		} catch (UnsupportedEncodingException e) {
			throw new MessageConversionException(
					"Failed to convert Message content", e);
		} 
		messageProperties.setContentType(MessageProperties.CONTENT_TYPE_JSON);
		messageProperties.setContentEncoding(this.defaultCharset);
		if (bytes != null) {
			messageProperties.setContentLength(bytes.length);
		}
		return new Message(bytes, messageProperties);

	}
}

3.生产者端调用


import java.util.List;

import org.springframework.amqp.core.AmqpTemplate;


public class MyMqGatway {
	
	@Autowired
	private AmqpTemplate amqpTemplate;
	
	public void sendDataToCrQueue(Object obj) {
		amqpTemplate.convertAndSend("queue_one_key", obj);
	}	
}

4.消费者端配置(与生产者端大同小异)

<?xml version="1.0" encoding="UTF-8"?>
<beans xmlns="http://www.springframework.org/schema/beans"
       xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
       xmlns:context="http://www.springframework.org/schema/context"
       xmlns:rabbit="http://www.springframework.org/schema/rabbit"
       xsi:schemaLocation="
            http://www.springframework.org/schema/beans
                http://www.springframework.org/schema/beans/spring-beans.xsd
            http://www.springframework.org/schema/context
                http://www.springframework.org/schema/context/spring-context.xsd
            http://www.springframework.org/schema/rabbit
                http://www.springframework.org/schema/rabbit/spring-rabbit-1.0.xsd">

    
   <!-- 连接服务配置  -->
   <rabbit:connection-factory id="connectionFactory" host="localhost" username="guest"
        password="guest" port="5672"  />
        
   <rabbit:admin connection-factory="connectionFactory"/>
   
   <!-- queue 队列声明-->
   <rabbit:queue id="queue_one" durable="true" auto-delete="false" exclusive="false" name="queue_one"/>
   
   
   <!-- exchange queue binging key 绑定 -->
    <rabbit:direct-exchange name="my-mq-exchange" durable="true" auto-delete="false" id="my-mq-exchange">
        <rabbit:bindings>
            <rabbit:binding queue="queue_one" key="queue_one_key"/>
        </rabbit:bindings>
    </rabbit:direct-exchange>

    
     
    <!-- queue litener  观察 监听模式 当有消息到达时会通知监听在对应的队列上的监听对象-->
    <rabbit:listener-container connection-factory="connectionFactory" acknowledge="auto" task-executor="taskExecutor">
        <rabbit:listener queues="queue_one" ref="queueOneLitener"/>
    </rabbit:listener-container>
</beans>


5.消费者端调用

import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageListener;

public class QueueOneLitener implements  MessageListener{
	@Override
	public void onMessage(Message message) {
		System.out.println(" data :" + message.getBody());
	}
}

6.由于消费端当队列有数据到达时,对应监听的对象就会被通知到,无法做到批量获取,批量入库,因此可以在消费端缓存一个临时队列,将mq取出来的数据存入本地队列,后台线程定时批量处理即可


打个广告:创业项目www.1zheke.com,有培训诉求的人可以支持一下,或者周围有培训需求的人,帮忙扩散一下,限量1折课,奖学拿到手软,非常感谢

共有 人打赏支持
粉丝 6
博文 5
码字总数 1079
评论 (22)
杨的
按您的配置会报错,报: No bean named 'taskExecutor' is defined
Fe-Fe

引用来自“杨的”的评论

按您的配置会报错,报: No bean named 'taskExecutor' is defined

你定义一个线程池的bean就可以了
优雅先生
fromMessage的这一行:
json = new String(message.getBody(),"UTF-8");
建议改成:
json = new String(message.getBody(),defaultCharset);
Fe-Fe

引用来自“美好的2014”的评论

fromMessage的这一行:
json = new String(message.getBody(),"UTF-8");
建议改成:
json = new String(message.getBody(),defaultCharset);

谢谢29
王建平winstar
exchange和queue的绑定操作两边都要执行吗?
小小小满呢
引入spring-rabbit的jar包后 报 java.lang.ClassNotFoundException: org.springframework.amqp.core.AmqpTemplate 可是反编译jar文件,这个类是存在的
开源中国首席脑科主任
为什么我是报connect error 连接不上???
EricTang

引用来自“Kings__”的评论

为什么我是报connect error 连接不上???
用的如果是guest用户的话,默认只能localhost的,检查一下,具体看报错信息
菘鬆
fe.json.FastJson 这个给我一份 可以吗
bowen0801
e.json.FastJson可以发一下吗 bowen0801@126.com
小矮子人
online033_login@126.com 求源码一份,谢谢 本地运行报错
现在_NOW
lgoodbook@163.com 求源码一份,谢谢 参考参考
你我他有个梦
liuzha13@163.com 求源码一份,谢谢 本地跑不通
你我他有个梦

引用来自“菘鬆”的评论

fe.json.FastJson 这个给我一份 可以吗
给你发过吗?给我发一份吧,谢谢
你我他有个梦

引用来自“菘鬆”的评论

fe.json.FastJson 这个给我一份 可以吗
liuzha13@163.com
尐葮阿譽
看过rabbitmq官网的人都知道,queue,exchange,routingKey三者的绑定只需要在消费端进行就可以了,无需在生产端进行.
zhuganlai
方便给俺发一份么,113617814@qq.com
logen2014

引用来自“菘鬆”的评论

fe.json.FastJson 这个给我一份 可以吗
312451021@qq.com
老范的自留地
/**
* 把json数据转换成类
* @param jsonString
* @param clazz
* @param <T>
* @return
*/
public static <T> T fromJson(String json, Class<T> clazz) {
return JSON.parseObject(json,clazz);
}


/**
* po类转换成json String
* @param po
* @return
*/
public static String toJson(Object obj) {
//String result = JSON.toJSONString(obj);
String result=JSON.toJSONString(obj, SerializerFeature.WriteMapNullValue);
return result;
}
}
学习使人上进
一个业务系统存在多个消费connector,这种情况,有遇到过吗,怎么处理呢
×
Fe-Fe
如果觉得我的文章对您有用,请随意打赏。您的支持将鼓励我继续创作!
* 金额(元)
¥1 ¥5 ¥10 ¥20 其他金额
打赏人
留言
* 支付类型
微信扫码支付
打赏金额:
已支付成功
打赏金额: