day75_淘淘商城项目_08_同步索引库问题分析 + ActiveMQ介绍/安装/使用 + ActiveMQ整合spring + 使用ActiveMQ实现添加商品后同步索引库_匠心笔记
课程计划
- 1、同步索引库问题分析
- 2、什么是MQ
- 3、ActiveMQ的安装
- 4、ActiveMQ的使用方法
- 5、ActiveMQ整合spring
- 6、使用消息队列实现添加商品后同步索引库
1、同步索引库问题分析
方案一:在taotao-manager中,添加商品的业务逻辑中,添加一个同步索引库的业务逻辑。
缺点:业务逻辑耦合度
非常高,业务逻辑拆分不明确。(单一职能原则)
方案二:业务逻辑在taotao-search中实现,调用服务在taotao-manager实现。业务逻辑分开。既可以是服务的生产者,也可以是服务的消费者。
缺点:服务之间的耦合度变高
。服务的启动有先后顺序
。随着调用的服务会越来越多,服务之间的调用越来越复杂
,难以管理。
方案三:使用消息队列。MQ是一个消息中间件。
MQ就相当于一个秘书、一个运营商,如下图所示:
存在的问题:
1、如果MQ挂了,所有相关的服务都挂了。
2、MQ有性能的瓶颈,尽量减少消息的内容的大小。
注意:技术的选型和具体的业务有关,只选择合适的技术。
如果MQ挂了:
1、通过日志查找原因。
2、通知相关的人员修复。
3、关键的业务必须保证有备用方案。
2、ActiveM的介绍Q
MQ是一个消息中间件,比如:ActiveMQ(Java)、RabbitMQ(Erlang)、Kafka(大数据)都属于MQ,是MQ的产品。
2.1、什么是ActiveMQ
ActiveMQ 是 Apache 出品,是最流行的,能力强劲的
开源消息总线
。ActiveMQ 是一个完全支持JMS1.1
和J2EE 1.4
规范的JMS Provider 实现
,尽管JMS规范出台已经是很久的事情了,但是JMS在当今的J2EE应用中间仍然扮演着特殊的地位。
主要特点:
1. 多种语言和协议编写客户端。语言: Java, C, C++, C#, Ruby, Perl, Python, PHP。应用协议: OpenWire, Stomp REST, WS Notification, XMPP, AMQP。
2. 完全支持JMS1.1和J2EE 1.4规范 (持久化、XA消息、事务)。
3. 对spring的支持
,ActiveMQ可以很容易内嵌到使用Spring的系统里面去,而且也支持Spring2.0的特性。
4. 通过了常见J2EE服务器(如:Geronimo, JBoss 4, GlassFish, WebLogic)的测试,其中通过 JCA 1.5 resource adaptors 的配置,可以让ActiveMQ可以自动的部署到任何兼容J2EE 1.4 商业服务器上。
5. 支持多种传送协议:in-VM, TCP,SSL, NIO,UDP, JGroups, JXTA。
6. 支持通过 JDBC 和 journal 提供高速的消息持久化。
7. 从设计上保证了高性能的集群,客户端-服务器,点对点。
8. 支持Ajax。
9. 支持与Axis的整合。
10. 可以很容易调用内嵌 JMS provider,进行测试。
2.2、ActiveMQ的消息形式
对于消息的传递有两种类型:
一种是点对点的
,即一个生产者和一个消费者一一对应。(短信)
另一种是发布/订阅模式
,即一个生产者产生消息并进行发送后,可以由多个消费者进行接收。(广播)
JMS定义了五种不同的消息正文格式
,以及调用的消息类型,允许你发送并接收以一些不同形式的数据,提供现有消息格式的一些级别的兼容性。
StreamMessage –> Java原始值的数据流
MapMessage –> 一套名称-值对
TextMessage --> 一个字符串对象
ObjectMessage –> 一个序列化的 Java对象
BytesMessage –> 一个字节的数据流
3、ActiveMQ的安装
网址:http://activemq.apache.org/
下载ActiveMQ,使用的版本是5.13.0,最新的版本是5.15.8。
3.1、安装环境
1、activemq是java语言开发的,需要安装jdk。
2、安装Linux系统。生产环境都是Linux系统。
3.2、安装步骤
第一步:把 apache-activemq-5.13.0-bin.tar.gz 的压缩包上传到Linux系统。
第二步:解压缩后删除压缩包。
[root@itheima ~]# ll
总用量 65472
-rw-r--r--. 1 root root 50385720 11月 25 12:32 apache-activemq-5.13.0-bin.tar.gz
[root@itheima ~]# tar zxf apache-activemq-5.13.0-bin.tar.gz
[root@itheima ~]# ll
总用量 52360
drwxr-xr-x. 10 root root 193 11月 30 2015 apache-activemq-5.13.0
-rw-r--r--. 1 root root 53613208 1月 7 2016 apache-activemq-5.13.0-bin.tar.gz
[root@itheima ~]# rm -rf apache-activemq-5.13.0-bin.tar.gz
第三步:启动。
使用bin目录下的activemq命令
启动:
[root@itheima bin]# ./activemq start
关闭:
[root@itheima bin]# ./activemq stop
查看状态:
[root@itheima bin]# ./activemq status
查看运行端口号:
[root@itheima bin]# ps aux | grep activemq
第四步:本博主的是CentOS7,防火墙使用的是firewalld,我们使用修改配置文件的方式来添加用到的端口(修改后需要重启firewalld服务)
[root@itheima bin]# cd /etc/firewalld/zones/
[root@itheima zones]# pwd
/etc/firewalld/zones
[root@itheima zones]# vim public.xml
编辑public.xml文件,添加代码:<port protocol="tcp" port="8161"/>
保存退出后,然后我们需要重启firewalld服务:
[root@itheima zones]# service firewalld restart
Redirecting to /bin/systemctl restart firewalld.service
[root@itheima zones]#
第五步:进入activemq管理后台。
访问地址:http://192.168.25.168:8161/admin
用户名:admin
密码:admin
如下图所示:
解决ActiveMQ访问后台时出现503错误(点击Queues按钮时):
1、查看机器名
[root@itheima bin]# vim /etc/sysconfig/network
NETWORKING=yes
HOSTNAME=taotao.com
2、修改机器名和ip的映射关系,即修改hosts文件
[root@itheima bin]# vim /etc/hosts
127.0.0.1 localhost localhost.localdomain localhost4 localhost4.localdomain4 taotao
::1 localhost localhost.localdomain localhost6 localhost6.localdomain6
3、重启ActiveMQ服务
4、ActiveMQ的使用方法
点对点:
发布/订阅:
4.1、Queue(队列)
ActiveMQ的Queue方式默认在ActiveMQ的服务端是持久化缓存的。
4.1.1、Producer
本例中,我们使用taotao-manager-service工程作为生产者。
生产者:生产消息,发送端。
把jar包添加到工程中。使用5.11.2版本的jar包。版本号我们在taotao-parent中的pom.xml中配置。
第一步:创建ConnectionFactory对象,需要指定服务端ip及端口号61616。注意:8161是ActiveMQ的web服务的端口。
第二步:使用ConnectionFactory对象创建一个Connection对象。
第三步:开启连接,调用Connection对象的start()方法。
第四步:使用Connection对象创建一个Session对象。
第五步:使用Session对象创建一个Destination对象(topic、queue),此处创建一个Queue对象。
第六步:使用Session对象创建一个Producer对象。
第七步:创建一个Message对象,创建一个TextMessage对象。
第八步:使用Producer对象发送消息。
第九步:关闭资源。
我们在taotao-manager-service工程中新建测试类和测试方法。
测试代码如下:
/**
* ActiveMQ的Queue方式的Producer
* @throws Exception
*/
@Test
public void queueProducerTest() throws Exception {
// 第一步:创建ConnectionFactory对象,需要指定服务端ip及端口号61616。注意:8161是ActiveMQ的web服务的端口。
ConnectionFactory connectionFactory = new ActiveMQConnectionFactory("tcp://192.168.25.168:61616"); // 注意:虚拟机上的安装ActiveMQ服务要添加该端口
// 第二步:使用ConnectionFactory对象创建一个Connection对象。
Connection connection = connectionFactory.createConnection();
// 第三步:开启连接,调用Connection对象的start方法。
connection.start();
// 第四步:使用Connection对象创建一个Session对象。
// 第一个参数:是否开启分布式事务。true:开启分布式事务,当第一个参数为true时,第二个参数忽略。
// 一般我们不使用分布式事务,因为分布式事务复杂,性能不好。互联网项目要求快速响应给用户,其余的事情交给消息队列向各个服务发消息让他们去做。这样保证了数据的最终一致。
// 第二个参数:当第一个参数为false时,第二个参数才有意义,表示消息的应答模式。1、自动应答 2、手动应答。一般是自动应答。
Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
// 第五步:使用Session对象创建一个Destination对象(topic、queue),此处创建一个Queue对象。参数表示队列的名称。
Queue queue = session.createQueue("test-queue");
// 第六步:使用Session对象创建一个Producer对象。
MessageProducer producer = session.createProducer(queue);
// 第七步:创建一个Message对象,创建一个TextMessage对象。
// 方式一:
// TextMessage message = new ActiveMQTextMessage();
// message.setText("hello ActiveMQ, this is my first test.");
// 方式二:
TextMessage textMessage = session.createTextMessage("hello ActiveMQ, this is my first test.");
// 第八步:使用Producer对象发送消息。
producer.send(textMessage);
// 第九步:关闭资源。
producer.close();
session.close();
connection.close();
}
进入activemq管理后台查看新增的queues队列,如下图:
4.1.2、Consumer
本例中,我们使用taotao-search-service工程作为消费者。
消费者:接收消息。
第一步:创建一个ConnectionFactory对象。
第二步:从ConnectionFactory对象中获得一个Connection对象。
第三步:开启连接。调用Connection对象的start()方法。
第四步:使用Connection对象创建一个Session对象。
第五步:使用Session对象创建一个Destination对象。和发送端保持一致queue,并且队列的名称一致。
第六步:使用Session对象创建一个Consumer对象。
第七步:接收消息。
第八步:打印消息。
第九步:关闭资源。
测试代码如下:
/**
* ActiveMQ的Queue方式的Consumer
* @throws Exception
*/
@Test
public void queueConsumerTest() throws Exception {
// 第一步:创建一个ConnectionFactory对象。
ConnectionFactory connectionFactory = new ActiveMQConnectionFactory("tcp://192.168.25.168:61616");
// 第二步:从ConnectionFactory对象中获得一个Connection对象。
Connection connection = connectionFactory.createConnection();
// 第三步:开启连接。调用Connection对象的start方法。
connection.start();
// 第四步:使用Connection对象创建一个Session对象。
Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
// 第五步:使用Session对象创建一个Destination对象。和发送端保持一致queue,并且队列的名称一致。
Queue queue = session.createQueue("test-queue");
// 第六步:使用Session对象创建一个Consumer对象。
MessageConsumer consumer = session.createConsumer(queue);
// 第七步:接收消息。
consumer.setMessageListener(new MessageListener() {
@Override
public void onMessage(Message message) {
if (message instanceof TextMessage) {
TextMessage textMessage = (TextMessage) message;
try {
String text = null;
// 取消息的内容
text = textMessage.getText();
// 第八步:打印消息。
System.out.println(text);
} catch (JMSException e) {
e.printStackTrace();
}
}
}
});
// 系统等待接收消息:
// 方式一:
/*
while (true) {
Thread.sleep(100);
}*/
// 方式二:等待键盘输入
System.in.read();
// 第九步:关闭资源
consumer.close();
session.close();
connection.close();
}
4.2、Topic(话题)
ActiveMQ的Topic方式与ActiveMQ的Queue方式大部分都是一样的。ActiveMQ的Topic方式默认在ActiveMQ的服务端是不持久化的。
如果消费者客户端没有接受到消息内容,消息内容就会丢了。所以消费者客户端需要先在启动状态(Topic方式可以有多个消费者)。我们也可以修改配置将Topic方式持久化缓存。后面可以自己学习。
4.2.1、Producer
使用步骤:
第一步:创建ConnectionFactory对象,需要指定服务端ip及端口号。
第二步:使用ConnectionFactory对象创建一个Connection对象。
第三步:开启连接,调用Connection对象的start方法。
第四步:使用Connection对象创建一个Session对象。
第五步:使用Session对象创建一个Destination对象(topic、queue),此处创建一个Topic对象。
第六步:使用Session对象创建一个Producer对象。
第七步:创建一个Message对象,创建一个TextMessage对象。
第八步:使用Producer对象发送消息。
第九步:关闭资源。
测试代码如下:
/**
* ActiveMQ的Topic方式的Consumer
* @throws Exception
*/
@Test
public void topicProducerTest() throws Exception {
// 第一步:创建ConnectionFactory对象,需要指定服务端ip及端口号61616。注意:8161是ActiveMQ的web服务的端口。
ConnectionFactory connectionFactory = new ActiveMQConnectionFactory("tcp://192.168.25.168:61616"); // 注意:虚拟机上的安装ActiveMQ服务要添加该端口
// 第二步:使用ConnectionFactory对象创建一个Connection对象。
Connection connection = connectionFactory.createConnection();
// 第三步:开启连接,调用Connection对象的start方法。
connection.start();
// 第四步:使用Connection对象创建一个Session对象。
// 第一个参数:是否开启分布式事务。true:开启分布式事务,当第一个参数为true时,第二个参数忽略。
// 一般我们不使用分布式事务,因为分布式事务复杂,性能不好。互联网项目要求快速响应给用户,其余的事情交给消息队列向各个服务发消息让他们去做。这样保证了数据的最终一致。
// 第二个参数:当第一个参数为false时,第二个参数才有意义,表示消息的应答模式。1、自动应答 2、手动应答。一般是自动应答。
Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
// 第五步:使用Session对象创建一个Destination对象(topic、queue),此处创建一个Topic对象。参数表示话题的名称。
Topic topic = session.createTopic("test-topic");
// 第六步:使用Session对象创建一个Producer对象。
MessageProducer producer = session.createProducer(topic);
// 第七步:创建一个Message对象,创建一个TextMessage对象。
// 方式一:
// TextMessage message = new ActiveMQTextMessage();
// message.setText("hello ActiveMQ, this is my first test.");
// 方式二:
TextMessage textMessage = session.createTextMessage("hello ActiveMQ, this is my first test.");
// 第八步:使用Producer对象发送消息。
producer.send(textMessage);
// 第九步:关闭资源。
producer.close();
session.close();
connection.close();
}
4.2.2、Consumer
消费者:接收消息。
第一步:创建一个ConnectionFactory对象。
第二步:从ConnectionFactory对象中获得一个Connection对象。
第三步:开启连接。调用Connection对象的start方法。
第四步:使用Connection对象创建一个Session对象。
第五步:使用Session对象创建一个Destination对象。和发送端保持一致topic,并且话题的名称一致。
第六步:使用Session对象创建一个Consumer对象。
第七步:接收消息。
第八步:打印消息。
第九步:关闭资源。
测试代码如下:
/**
* ActiveMQ的Topic方式的Consumer
* @throws Exception
*/
@Test
public void topicConsumerTest() throws Exception {
// 第一步:创建一个ConnectionFactory对象。
ConnectionFactory connectionFactory = new ActiveMQConnectionFactory("tcp://192.168.25.168:61616");
// 第二步:从ConnectionFactory对象中获得一个Connection对象。
Connection connection = connectionFactory.createConnection();
// 第三步:开启连接。调用Connection对象的start方法。
connection.start();
// 第四步:使用Connection对象创建一个Session对象。
Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
// 第五步:使用Session对象创建一个Destination对象。和发送端保持一致Topic,并且队列的名称一致。
Topic topic = session.createTopic("test-topic");
// 第六步:使用Session对象创建一个Consumer对象。
MessageConsumer consumer = session.createConsumer(topic);
// 第七步:接收消息。
consumer.setMessageListener(new MessageListener() {
@Override
public void onMessage(Message message) {
if (message instanceof TextMessage) {
TextMessage textMessage = (TextMessage) message;
try {
String text = null;
// 取消息的内容
text = textMessage.getText();
// 第八步:打印消息。
System.out.println(text);
} catch (JMSException e) {
e.printStackTrace();
}
}
}
});
System.out.println("topic的消费者客户端01...");
// 系统等待接收消息:
// 方式一:
/*
while (true) {
Thread.sleep(100);
}*/
// 方式二:等待键盘输入
System.in.read();
// 第九步:关闭资源
consumer.close();
session.close();
connection.close();
}
注意:
Topic 默认是不存在于MQ服务器中的,一旦发送之后,如果没有订阅,就没了。
Queue 默认是存在于MQ的服务器中的,发送消息之后,随时取。但是一定是一个消费者取完就没了。
5、ActiveMQ整合spring
5.1、使用方法
第一步:把Activemq模板相关的jar包添加到工程中。在Maven工程中是添加依赖jar包。最开始配置spring时已经配置过了。
第二步:编写配置文件applicationContext-activemq.xml。
第三步:配置生产者。使用JMSTemplate对象。发送消息。
第四步:在spring容器中配置Destination。
<?xml version="1.0" encoding="UTF-8"?>
<beans xmlns="http://www.springframework.org/schema/beans"
xmlns:context="http://www.springframework.org/schema/context"
xmlns:p="http://www.springframework.org/schema/p"
xmlns:aop="http://www.springframework.org/schema/aop"
xmlns:tx="http://www.springframework.org/schema/tx"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://www.springframework.org/schema/beans
http://www.springframework.org/schema/beans/spring-beans-4.2.xsd
http://www.springframework.org/schema/context
http://www.springframework.org/schema/context/spring-context-4.2.xsd
http://www.springframework.org/schema/aop
http://www.springframework.org/schema/aop/spring-aop-4.2.xsd
http://www.springframework.org/schema/tx
http://www.springframework.org/schema/tx/spring-tx-4.2.xsd
http://www.springframework.org/schema/util
http://www.springframework.org/schema/util/spring-util-4.2.xsd">
<!-- 真正可以产生Connection的ConnectionFactory,由对应的 JMS服务厂商提供 -->
<bean id="targetConnectionFactory" class="org.apache.activemq.ActiveMQConnectionFactory">
<property name="brokerURL" value="tcp://192.168.25.168:61616"></property>
</bean>
<!-- Spring用于管理真正的ConnectionFactory的ConnectionFactory -->
<bean id="connectionFactory" class="org.springframework.jms.connection.SingleConnectionFactory">
<!-- 目标ConnectionFactory对应真实的可以产生JMS Connection的ConnectionFactory -->
<property name="targetConnectionFactory" ref="targetConnectionFactory"></property>
</bean>
<!-- 接收和发送消息时使用的类 -->
<!-- 配置消息的生产者 -->
<!-- Spring提供的JMS工具类,它可以进行消息发送、接收等 -->
<bean id="jmsTemplate" class="org.springframework.jms.core.JmsTemplate">
<!-- 这个connectionFactory对应的是我们定义的Spring提供的那个ConnectionFactory对象 -->
<property name="connectionFactory" ref="connectionFactory"></property>
</bean>
<!-- 配置消息的Destination对象 -->
<!--这个是队列目的地,点对点的 -->
<bean id="queueDestination" class="org.apache.activemq.command.ActiveMQQueue">
<constructor-arg name="name" value="spring-test-queue"></constructor-arg>
</bean>
<!-- 这个是话题目的地,一对多的 -->
<bean id="topicDestination" class="org.apache.activemq.command.ActiveMQTopic">
<constructor-arg name="name" value="spring-test-topic"></constructor-arg>
</bean>
</beans>
第五步:代码测试。
5.2、代码测试(Queue方式)
5.2.1、发送消息
taotao-manager-service中发送消息。
第一步:初始化一个spring容器。即加载配置文件applicationContext-activemq.xml。
第二步:从容器中获得JMSTemplate对象。
第三步:从容器中获得一个Destination对象。
第四步:使用JMSTemplate对象发送消息,需要知道Destination。
测试代码如下:
/**
* ActiveMQ整合spring,ActiveMQ的Queue方式的Producer
* @throws Exception
*/
@Test
public void queueProducerTest() throws Exception {
// 第一步:初始化一个spring容器,即加载配置文件applicationContext-activemq.xml。
ApplicationContext applicationContext = new ClassPathXmlApplicationContext("classpath:spring/applicationContext-activemq.xml");
// 第二步:从容器中获得JMSTemplate对象。
JmsTemplate jmsTemplate = applicationContext.getBean(JmsTemplate.class);
// 第三步:从容器中获得一个Destination对象
Queue queue = (Queue) applicationContext.getBean("queueDestination"); // 建议根据id取
// 第四步:使用JMSTemplate对象发送消息,需要知道Destination
jmsTemplate.send(queue, new MessageCreator() {
@Override
public Message createMessage(Session session) throws JMSException {
TextMessage textMessage = session.createTextMessage("spring activemq test");
return textMessage;
}
});
}
5.2.2、接收消息
taotao-search-service中接收消息。
第一步:把jar包添加到工程中。使用5.11.2版本的jar包。版本号我们在taotao-parent中的pom.xml中配置。
第二步:创建一个MessageListener的实现类。用于接收ActiveMQ发送的消息。
/**
* 自定义的MessageListener的实现类,用于接收ActiveMQ发送的消息
* @author chenmingjun
* @date 2018年11月25日下午11:30:56
* @version 1.0
*/
public class MyMessageListener implements MessageListener {
@Override
public void onMessage(Message message) {
// 接收到消息,进行消费
try {
TextMessage textMessage = (TextMessage) message;
// 取消息内容
String text = textMessage.getText();
System.out.println(text);
} catch (JMSException e) {
e.printStackTrace();
}
}
}
第三步:编写配置文件applicationContext-activemq.xml,配置spring和Activemq整合,配置消息的消费者。
<?xml version="1.0" encoding="UTF-8"?>
<beans xmlns="http://www.springframework.org/schema/beans"
xmlns:context="http://www.springframework.org/schema/context"
xmlns:p="http://www.springframework.org/schema/p"
xmlns:aop="http://www.springframework.org/schema/aop"
xmlns:tx="http://www.springframework.org/schema/tx"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://www.springframework.org/schema/beans
http://www.springframework.org/schema/beans/spring-beans-4.2.xsd
http://www.springframework.org/schema/context
http://www.springframework.org/schema/context/spring-context-4.2.xsd
http://www.springframework.org/schema/aop
http://www.springframework.org/schema/aop/spring-aop-4.2.xsd
http://www.springframework.org/schema/tx
http://www.springframework.org/schema/tx/spring-tx-4.2.xsd
http://www.springframework.org/schema/util
http://www.springframework.org/schema/util/spring-util-4.2.xsd">
<!-- 真正可以产生Connection的ConnectionFactory,由对应的 JMS服务厂商提供 -->
<bean id="targetConnectionFactory" class="org.apache.activemq.ActiveMQConnectionFactory">
<property name="brokerURL" value="tcp://192.168.25.168:61616"></property>
</bean>
<!-- Spring用于管理真正的ConnectionFactory的ConnectionFactory -->
<bean id="connectionFactory" class="org.springframework.jms.connection.SingleConnectionFactory">
<!-- 目标ConnectionFactory对应真实的可以产生JMS Connection的ConnectionFactory -->
<property name="targetConnectionFactory" ref="targetConnectionFactory"></property>
</bean>
<!-- 接收和发送消息时使用的类 -->
<!-- 配置消息的Destination对象 -->
<!--这个是队列目的地,点对点的 -->
<bean id="queueDestination" class="org.apache.activemq.command.ActiveMQQueue">
<constructor-arg name="name" value="spring-test-queue"></constructor-arg>
</bean>
<!-- 这个是话题目的地,一对多的 -->
<bean id="topicDestination" class="org.apache.activemq.command.ActiveMQTopic">
<constructor-arg name="name" value="spring-test-topic"></constructor-arg>
</bean>
<!-- 配置消息的消费者 -->
<!-- 先配置监听器 -->
<bean id="myMessageListener" class="com.taotao.search.listener.MyMessageListener" />
<!-- 再配置消息监听容器 -->
<bean class="org.springframework.jms.listener.DefaultMessageListenerContainer">
<property name="connectionFactory" ref="connectionFactory" />
<property name="destination" ref="queueDestination" />
<property name="messageListener" ref="myMessageListener" />
</bean>
</beans>
第四步:测试代码。
测试代码如下:
/**
* ActiveMQ整合spring,ActiveMQ的Queue方式的Consumer
* @throws Exception
*/
@Test
public void queueConsumerTest() throws Exception {
// 第一步:初始化一个spring容器,即加载配置文件applicationContext-activemq.xml。
ApplicationContext applicationContext = new ClassPathXmlApplicationContext("classpath:spring/applicationContext-activemq.xml");
// 第二步:系统等待接收消息:
// 方式一:
/*
while (true) {
Thread.sleep(100);
}*/
// 方式二:等待键盘输入
System.in.read();
}
Topic方式同理!
6、ActiveMQ整合到项目中
需要在商品的添加、修改、删除的时候,同步索引库。将数据从数据库中查询出来导入到索引库更新。
消息的发送方为:taotao-manager-service
消息的接收方为:taotao-search-service
两个工程都需要依赖activmq:
6.1、Producer
在taotao-manager-service工程中发送消息。
需要加入对activmq的依赖。
功能分析:
发送端当商品添加完成后发送一个TextMessage,包含一个商品id即可。
接收端接收到商品id通过数据库查询到商品的信息(注意:搜索的结果商品的信息,此商品非彼商品)再去同步索引库。
6.1.1、applicationContext-activemq.xml的配置
taotao-manager-service工程中的applicationContext-activemq.xml
<?xml version="1.0" encoding="UTF-8"?>
<beans xmlns="http://www.springframework.org/schema/beans"
xmlns:context="http://www.springframework.org/schema/context"
xmlns:p="http://www.springframework.org/schema/p"
xmlns:aop="http://www.springframework.org/schema/aop"
xmlns:tx="http://www.springframework.org/schema/tx"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://www.springframework.org/schema/beans
http://www.springframework.org/schema/beans/spring-beans-4.2.xsd
http://www.springframework.org/schema/context
http://www.springframework.org/schema/context/spring-context-4.2.xsd
http://www.springframework.org/schema/aop
http://www.springframework.org/schema/aop/spring-aop-4.2.xsd
http://www.springframework.org/schema/tx
http://www.springframework.org/schema/tx/spring-tx-4.2.xsd
http://www.springframework.org/schema/util
http://www.springframework.org/schema/util/spring-util-4.2.xsd">
<!-- 真正可以产生Connection的ConnectionFactory,由对应的 JMS服务厂商提供 -->
<bean id="targetConnectionFactory" class="org.apache.activemq.ActiveMQConnectionFactory">
<property name="brokerURL" value="tcp://192.168.25.168:61616"></property>
</bean>
<!-- Spring用于管理真正的ConnectionFactory的ConnectionFactory -->
<bean id="connectionFactory" class="org.springframework.jms.connection.SingleConnectionFactory">
<!-- 目标ConnectionFactory对应真实的可以产生JMS Connection的ConnectionFactory -->
<property name="targetConnectionFactory" ref="targetConnectionFactory"></property>
</bean>
<!-- 接收和发送消息时使用的类 -->
<!-- 配置消息的生产者 -->
<!-- Spring提供的JMS工具类,它可以进行消息发送、接收等 -->
<bean id="jmsTemplate" class="org.springframework.jms.core.JmsTemplate">
<!-- 这个connectionFactory对应的是我们定义的Spring提供的那个ConnectionFactory对象 -->
<property name="connectionFactory" ref="connectionFactory"></property>
</bean>
<!-- 配置消息的Destination对象 -->
<!--这个是队列目的地,点对点的 -->
<bean id="queueDestination" class="org.apache.activemq.command.ActiveMQQueue">
<constructor-arg name="name" value="spring-test-queue"></constructor-arg>
</bean>
<!-- 由于新增商品,对应的商品搜索索引库要同步、要生成订单页面、要同步缓存等,即很多地方要监听商品添加这个事件,所以我们使用Topic -->
<!-- 这个是话题目的地,一对多的 -->
<bean id="itemAddTopic" class="org.apache.activemq.command.ActiveMQTopic">
<constructor-arg name="name" value="item-add-topic"></constructor-arg>
</bean>
</beans>
6.1.2、代码添加发送消息逻辑的实现
在ItemServiceImpl中添加业务逻辑。在添加商品后,返回结果之前发送消息。
代码如下:
@Override
public TaotaoResult saveItem(TbItem item, String desc) {
// 1、生成商品id,本例中使用工具类IDUtils生成商品id
final Long itemId = IDUtils.genItemId();
item.setId(itemId);
// 2、补全商品表TbItem的其他属性
// 商品状态,1-正常,2-下架,3-删除
item.setStatus((byte) 1);
Date date = new Date();
item.setCreated(date);
item.setUpdated(date);
// 3、向商品表中插入数据
itemMapper.insert(item);
// 4、创建一个商品描述表TbItemDesc对象
TbItemDesc itemDesc = new TbItemDesc();
// 5、补全商品描述表TbItemDesc的其他属性
itemDesc.setItemId(itemId);
itemDesc.setItemDesc(desc);
itemDesc.setCreated(date);
itemDesc.setUpdated(date);
// 6、向商品描述表中插入数据
itemDescMapper.insert(itemDesc);
// 向ActiveMQ发送一个商品添加的消息,使用JmsTemplate对象,需要注入进来
jmsTemplate.send(topicDestination, new MessageCreator() {
@Override
public Message createMessage(Session session) throws JMSException {
// 发送商品id
TextMessage textMessage = session.createTextMessage(itemId + "");
return textMessage;
}
});
// 7、返回TaotaoResult.ok()
return TaotaoResult.ok();
}
6.2、Consumer
在taotao-search-service工程中消费消息。
需要加入对activmq的依赖。
功能分析:
1、接收消息。需要创建MessageListener接口的实现类。
2、取消息,取商品id。
3、根据商品id查询数据库。
4、创建SolrInputDocument对象。
5、使用SolrServer对象写入索引库。
6、返回成功,返回TaotaoResult。
6.2.1、Dao层
根据商品id查询商品搜索信息(注意:是从3张表中查,此商品非彼商品)。
返回一个SearchItem。
SearchItemMapper接口中添加如下的方法:
SearchItemMapper.java
/**
* 根据商品id查询搜索商品信息。(注意:是从3张表中查,此商品非彼商品)
* @param iteamId
* @return
*/
SearchItem getSearchItemById(Long iteamId);
映射文件中添加如下内容:
SearchItemMapper.xml
<select id="getSearchItemById" parameterType="Long" resultType="com.taotao.common.pojo.SearchItem">
SELECT
a.id,
a.title,
a.sell_point,
a.price,
a.image,
b. NAME AS category_name,
c.item_desc
FROM
tb_item a
LEFT JOIN tb_item_cat b ON a.cid = b.id
LEFT JOIN tb_item_desc c ON a.id = c.item_id
WHERE
a.`status` = 1
AND
a.id = #{itemId}
</select>
6.2.2、Service层
业务逻辑可以在MessageListener中实现,也可以在MessageListener中实现。本例中我们在MessageListener中实现。
6.2.3、MessageListener
ItemAddListenerjava
/**
* 监听商品添加事件,同步索引库
* @author chenmingjun
* @date 2018年11月26日上午10:41:25
* @version 1.0
*/
public class ItemAddListener implements MessageListener {
@Autowired
private SearchItemMapper searchItemMapper;
// 注入SolrServer
@Autowired
private SolrServer solrServer;
@Override
public void onMessage(Message message) {
try {
// 从消息中取出商品id
TextMessage textMessage = null;
Long iteamId = null;
// 判断消息类型是否是TextMessage类型
if (message instanceof TextMessage) {
// 是TextMessage类型,就获取该消息
textMessage = (TextMessage) message;
iteamId = Long.parseLong(textMessage.getText());
}
// 在查询数据库之前,需要注意:需要等待一下“服务层的消息生产者”的事务提交,否则会报空指针异常;
// 或者将发消息的时机放到“表现层”,表现层调服务层服务,那么事务一定提交了。
Thread.sleep(1000);
// 根据商品id查询数据库,取出商品消息(是搜索的结果商品的信息,此商品非彼商品)
SearchItem searchItem = searchItemMapper.getSearchItemById(iteamId);
// 创建文档对象SolrInputDocument
// 向文档对象中添加域
// 把文档对象写入索引库
// 提交
SolrInputDocument document = new SolrInputDocument();
document.addField("id", searchItem.getId());
document.addField("item_title", searchItem.getTitle());
document.addField("item_sell_point", searchItem.getSell_point());
document.addField("item_price", searchItem.getPrice());
document.addField("item_image", searchItem.getImage());
document.addField("item_category_name", searchItem.getCategory_name());
document.addField("item_desc", searchItem.getItem_desc());
solrServer.add(document);
solrServer.commit();
} catch (Exception e) {
e.printStackTrace();
}
}
}
6.2.4、spring中配置消息监听容器
taotao-search-service工程中applicationContext-activemq.xml
<?xml version="1.0" encoding="UTF-8"?>
<beans xmlns="http://www.springframework.org/schema/beans"
xmlns:context="http://www.springframework.org/schema/context"
xmlns:p="http://www.springframework.org/schema/p"
xmlns:aop="http://www.springframework.org/schema/aop"
xmlns:tx="http://www.springframework.org/schema/tx"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://www.springframework.org/schema/beans
http://www.springframework.org/schema/beans/spring-beans-4.2.xsd
http://www.springframework.org/schema/context
http://www.springframework.org/schema/context/spring-context-4.2.xsd
http://www.springframework.org/schema/aop
http://www.springframework.org/schema/aop/spring-aop-4.2.xsd
http://www.springframework.org/schema/tx
http://www.springframework.org/schema/tx/spring-tx-4.2.xsd
http://www.springframework.org/schema/util
http://www.springframework.org/schema/util/spring-util-4.2.xsd">
<!-- 真正可以产生Connection的ConnectionFactory,由对应的 JMS服务厂商提供 -->
<bean id="targetConnectionFactory" class="org.apache.activemq.ActiveMQConnectionFactory">
<property name="brokerURL" value="tcp://192.168.25.168:61616"></property>
</bean>
<!-- Spring用于管理真正的ConnectionFactory的ConnectionFactory -->
<bean id="connectionFactory" class="org.springframework.jms.connection.SingleConnectionFactory">
<!-- 目标ConnectionFactory对应真实的可以产生JMS Connection的ConnectionFactory -->
<property name="targetConnectionFactory" ref="targetConnectionFactory"></property>
</bean>
<!-- 接收和发送消息时使用的类 -->
<!-- 配置消息的Destination对象 -->
<!--这个是队列目的地,点对点的 -->
<bean id="queueDestination" class="org.apache.activemq.command.ActiveMQQueue">
<constructor-arg name="name" value="spring-test-queue"></constructor-arg>
</bean>
<!-- 这个是话题目的地,一对多的 -->
<bean id="itemAddTopic" class="org.apache.activemq.command.ActiveMQTopic">
<constructor-arg name="name" value="item-add-topic"></constructor-arg>
</bean>
<!-- 配置消息的消费者 -->
<!-- 先配置监听器 -->
<bean id="myMessageListener" class="com.taotao.search.listener.MyMessageListener" />
<!-- 再配置消息监听容器 -->
<bean class="org.springframework.jms.listener.DefaultMessageListenerContainer">
<property name="connectionFactory" ref="connectionFactory" />
<property name="destination" ref="queueDestination" />
<property name="messageListener" ref="myMessageListener" />
</bean>
<!-- 先配置监听器 -->
<bean id="itemAddListener" class="com.taotao.search.listener.ItemAddListener" />
<!-- 再配置消息监听容器 -->
<bean class="org.springframework.jms.listener.DefaultMessageListenerContainer">
<property name="connectionFactory" ref="connectionFactory" />
<property name="destination" ref="itemAddTopic" />
<property name="messageListener" ref="itemAddListener" />
</bean>
</beans>
截图说明如下: