版权声明:LeifChen原创,转载请注明出处,谢谢~ https://blog.csdn.net/leifchen90/article/details/84248762
Java 消息中间件
文章目录
消息中间件:关注于数据的发送与接收,利用高效可靠的异步消息传递机制集成分布式系统。
常见消息中间件
1. ActiveMQ
ActiveMQ 是 Apache 出品,最流行,能力强劲的开源消息总线。
- 完全支持 JMS 1.1 和 J2EE 1.4规范(持久化,XA消息,事务)
- 支持多种语言和协议编写客户端
- 虚拟主题、组合目的、镜像队列
下载解压后,执行 /bin/win64/ 路径下 的 activemq.bat 批处理文件,并打开 http://localhost:8161
查看是否安装成功。
2. RabbitMQ
RabbitMQ 是一个开源的 AMQP 实现,服务器端用 Erlang 语言编写。用于在分布式系统中存储转发消息,在易用性、扩展性、高可用性等方面表现不俗。
- 支持多种客户端
- AMQP 的完整实现(vhost、Exchange、Binding、Routing Key 等)
- 事务支持/发布确认
- 消息持久化
3. Kafka
Kafka 是一种高吞吐量的分布式发布订阅消息系统,是一个分布式的、分区的、可靠的分布式日志存储服务。
- 通过O(1)的磁盘数据结构提供消息的持久化,对以TB的消息存储也能够保持长时间的稳定性能
- 高吞吐量:即使是非常普通的硬件,Kafka 也可以支持每秒数百万的消息
- Partition、Consumer Group
4. 综合比较
Propertity | ActiveMQ | RabbitMQ | Kafka |
---|---|---|---|
跨语言 | 支持(Java 优先) | 语言无关 | 支持(Java 优先) |
支持协议 | OpenWire,Stomp,XMPP,AMQP | AMQP | |
优点 | 遵循 JMS 规范;安装部署方便 | 继承 Erlang 天生的并发性,稳定性,安全性有保障 | |
缺点 | 消息丢失;社区不活跃 | Erlang 语言难度较大;不支持动态扩展 | 严格的顺序机制,不支持消息优先级;不支持标准的消息协议,不利于平台迁移 |
综合评价 | 适合中小企业级消息应用场景,不适合上千个队列的应用场景 | 适合对稳定性要求高的企业级应用 | 一般应用在大数据日志处理或对实时性、可靠性要求稍低的场景 |
规范与协议
1. JMS 规范
Java 消息服务(Java Message Service)即 JMS,是一个 Java 平台关于面向消息中间件的 API,用于在两个应用程序之间,或分布式系统中发送消息,进行异步通信。
- 提供者:实现 JMS规范的消息中间件服务器
- 客户端:发送或接收消息的应用程序
- 生产者/发布者:创建并发送消息的客户端
- 消费者/订阅者:接收并处理消息的客户端
- 消息:应用层序之间传递的数据内容
- 消息模式:传递消息的方式,JMS 中定义了队列和主题两种模式
1.1 消息模式
1.1.1 队列模型
- 客户端包括生产者和消费者
- 队列中消息只能被一个消费者消费
- 消费者可以随时消费队列中的消息
代码
- 生产者
package com.chen.jms.queue;
import org.apache.activemq.ActiveMQConnectionFactory;
import javax.jms.*;
/**
* 队列模式:生产者
*
* @Author LeifChen
* @Date 2018-11-16
*/
public class QueueProducer {
private static final String URL = "tcp://localhost:61616";
private static final String QUEUE_NAME = "queue-test";
private static final int COUNT = 100;
public static void main(String[] args) throws JMSException {
// 1.创建 ConnectionFactory
ConnectionFactory connectionFactory = new ActiveMQConnectionFactory(URL);
// 2.创建 Connection
Connection connection = connectionFactory.createConnection();
// 3.启动连接
connection.start();
// 4.创建会话
Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
// 5.创建一个目标
Destination destination = session.createQueue(QUEUE_NAME);
// 6.创建一个发布者
MessageProducer producer = session.createProducer(destination);
for (int i = 0; i < COUNT; i++) {
// 7.创建消息
TextMessage textMessage = session.createTextMessage("test" + i);
// 8.发布消息
producer.send(textMessage);
System.out.println("发送消息:" + textMessage.getText());
}
// 9.关闭连接
connection.close();
}
}
- 消费者
package com.chen.jms.queue;
import org.apache.activemq.ActiveMQConnectionFactory;
import javax.jms.*;
/**
* 队列模式:消费者
*
* @Author LeifChen
* @Date 2018-11-16
*/
public class QueueConsumer {
private static final String URL = "tcp://localhost:61616";
private static final String QUEUE_NAME = "queue-test";
public static void main(String[] args) throws JMSException {
// 1.创建 ConnectionFactory
ConnectionFactory connectionFactory = new ActiveMQConnectionFactory(URL);
// 2.创建 Connection
Connection connection = connectionFactory.createConnection();
// 3.启动连接
connection.start();
// 4.创建会话
Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
// 5.创建一个目标
Destination destination = session.createQueue(QUEUE_NAME);
// 6.创建一个消费者
MessageConsumer consumer = session.createConsumer(destination);
// 7.创建一个监听器
consumer.setMessageListener(message -> {
TextMessage textMessage = (TextMessage) message;
try {
System.out.println("接收消息:" + textMessage.getText());
} catch (JMSException e) {
e.printStackTrace();
}
});
}
}
1.1.2 主题模型
- 客户端包括发布者和订阅者
- 主题中的消息被所有订阅者消费
- 消费者不能消费订阅之前就发送到主题中的消息
代码
- 发布者
package com.chen.jms.topic;
import org.apache.activemq.ActiveMQConnectionFactory;
import javax.jms.*;
/**
* 主题模式:发布者
*
* @Author LeifChen
* @Date 2018-11-16
*/
public class TopicProducer {
private static final String URL = "tcp://localhost:61616";
private static final String TOPIC_NAME = "topic-test";
private static final int COUNT = 100;
public static void main(String[] args) throws JMSException {
// 1.创建 ConnectionFactory
ConnectionFactory connectionFactory = new ActiveMQConnectionFactory(URL);
// 2.创建 Connection
Connection connection = connectionFactory.createConnection();
// 3.启动连接
connection.start();
// 4.创建会话
Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
// 5.创建一个目标
Destination destination = session.createTopic(TOPIC_NAME);
// 6.创建一个发布者
MessageProducer producer = session.createProducer(destination);
for (int i = 0; i < COUNT; i++) {
// 7.创建消息
TextMessage textMessage = session.createTextMessage("test" + i);
// 8.发布消息
producer.send(textMessage);
System.out.println("发送消息:" + textMessage.getText());
}
// 9.关闭连接
connection.close();
}
}
- 订阅者
package com.chen.jms.topic;
import org.apache.activemq.ActiveMQConnectionFactory;
import javax.jms.*;
/**
* 主题模式:订阅者
*
* @Author LeifChen
* @Date 2018-11-16
*/
public class TopicConsumer {
private static final String URL = "tcp://localhost:61616";
private static final String TOPIC_NAME = "topic-test";
public static void main(String[] args) throws JMSException {
// 1.创建 ConnectionFactory
ConnectionFactory connectionFactory = new ActiveMQConnectionFactory(URL);
// 2.创建 Connection
Connection connection = connectionFactory.createConnection();
// 3.启动连接
connection.start();
// 4.创建会话
Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
// 5.创建一个目标
Destination destination = session.createTopic(TOPIC_NAME);
// 6.创建一个消费者
MessageConsumer consumer = session.createConsumer(destination);
// 7.创建一个监听器
consumer.setMessageListener(message -> {
TextMessage textMessage = (TextMessage) message;
try {
System.out.println("接收消息:" + textMessage.getText());
} catch (JMSException e) {
e.printStackTrace();
}
});
}
}
1.2 JMS 编码接口
- ConnectionFactory:用于创建连接到消息中间件的连接工厂
- Connection:代表了应用程序和消息服务器之间的通信链路
- MessageConsumer:由会话创建,用于接收发送到目标的消息
- MessageProducer:由会话创建,用于发送消息到目标
- Meeage:是在消费者和生产者之间传送的对象,消息头,一组消息属性,一个消息体
- Destination:指消息发布和接收的地点,包括队列或主题
2. AMQP 协议
AMQP(advanced message queuing protocol)是一个提供消息服务的应用层标准协议,基于此协议的客户端与消息中间件可传递消息,并不受客户端/中间件不同产品、不同开发语言等条件的限制。
3. JMS 与 AMQP 比较
Propertity | JMS 规范 | AMQP 协议 |
---|---|---|
定义 | Java API | Wire-protocol |
跨语言 | 否 | 是 |
消息模型 | 提供两种消息模型: p2p pub/sub |
提供五种消息模型: direct fanout topic headers system |
消息类型 | TextMessage MapMessage BytesMessage StreamMessage ObjectMessage Message |
byte[] |
综合评价 | JMS 定义了 Java API 层面的标准 | AMQP 是面向详细、队列、路由(包括点对点的发布/订阅)、可靠性、安全 |