JAVA分布式--ActiveMQ 消息中间件(下)

简介: JAVA分布式--ActiveMQ 消息中间件(下)

2. ActiveMQ 示例

1). P2P 示例

I. 导包–activemq-all-5.15.3.jar

II. Producer


/**
 * 定义消息的生产者
 * @author mazaiting
 */
public class Producer {
    // 用户名
    private static final String USERNAME = ActiveMQConnection.DEFAULT_USER;
    // 密码
    private static final String PASSWORD = ActiveMQConnection.DEFAULT_PASSWORD;
    // 链接
    private static final String BROKENURL = ActiveMQConnection.DEFAULT_BROKER_URL;
    /**
     * 定义消息并发送,等待消息的接收者(消费者)消费此消息
     * @param args
     * @throws JMSException 
     */
    public static void main(String[] args) throws JMSException {
        // 消息中间件的链接工厂
        ConnectionFactory connectionFactory = new ActiveMQConnectionFactory(
                USERNAME, PASSWORD, BROKENURL);
        // 连接
        Connection connection = null;
        // 会话
        Session session = null;
        // 消息的目的地
        Destination destination = null;
        // 消息生产者
        MessageProducer messageProducer = null;
        try {
            // 通过连接工厂获取链接
            connection = connectionFactory.createConnection();
            // 创建会话,进行消息的发送
            // 参数一:是否启用事务
            // 参数二:设置自动签收
            session = connection.createSession(true, Session.AUTO_ACKNOWLEDGE);
            // 创建消息队列
            destination = session.createQueue("talkWithMo");
            // 创建一个消息生产者
            messageProducer = session.createProducer(destination);
            // 设置持久化/非持久化, 如果非持久化,MQ重启后可能后导致消息丢失
            messageProducer.setDeliveryMode(DeliveryMode.NON_PERSISTENT);
            // 模拟发送消息
            for (int i = 0; i < 5; i++) {
                TextMessage textMessage = session.createTextMessage("给妈妈发送的消息:"+i);
                System.out.println("textMessage: " + textMessage);
                messageProducer.send(textMessage);
            }
            // 如果设置了事务,会话就必须提交
            session.commit();
        } catch (JMSException e) {
            e.printStackTrace();
        } finally {
            if (null != connection) {
                connection.close();
            }
        }
    }
}


III. Consumer


/**
 * 定义消息的消费者
 * @author mazaiting
 */
public class Consumer {
    // 用户名
    private static final String USERNAME = ActiveMQConnection.DEFAULT_USER;
    // 密码
    private static final String PASSWORD = ActiveMQConnection.DEFAULT_PASSWORD;
    // 链接
    private static final String BROKENURL = ActiveMQConnection.DEFAULT_BROKER_URL;
    /**
     * 接收消息
     * @param args
     * @throws JMSException 
     */
    public static void main(String[] args) throws JMSException {
        // 消息中间件的链接工厂
        ConnectionFactory connectionFactory = null;
        // 链接
        Connection connection = null;
        // 会话
        Session session = null;
        // 消息的目的地
        Destination destination = null;
        // 消息的消费者
        MessageConsumer messageConsumer = null;
        // 实例化链接工厂,创建一个链接
        connectionFactory = new ActiveMQConnectionFactory(USERNAME, PASSWORD, BROKENURL);
        try {
            // 通过工厂获取链接
            connection = connectionFactory.createConnection();
            // 启动链接
            connection.start();
            // 创建会话,进行消息的接收
            session = connection.createSession(true, Session.AUTO_ACKNOWLEDGE);
            // 创建消息队列
            destination = session.createQueue("talkWithMo");
            // 创建一个消息的消费者
            messageConsumer = session.createConsumer(destination);
            // 模拟接收消息
            while (true) {
                TextMessage textMessage = (TextMessage) messageConsumer.receive(10000);
                if (null != textMessage) {
                    System.out.println("收到消息: " + textMessage);
                } else {
                    break;
                }
            }
            // 提交
            session.commit();
        } catch (JMSException e) {
            e.printStackTrace();
        } finally {
            if (null != connection) {
                connection.close();
            }
        }
    }
}


IV. 测试

先运行生产者Producer


image.png


ActiveMQ控制台


image.png


再运行消费者Consumer


image.png


ActiveMQ控制台


image.png


V. 消息类型


StreamMessage Java原始值的数据流

MapMessage 一套名称-键值对

TextMessage 一个字符串对象

ObjectMessage 一个序列号的Java对象

BytesMessage 一个未解释字节的数据流

VI. 控制台 Queue

Messages Enqueued:表示生产了多少条消息,记做P

Messages Dequeued:表示消费了多少条消息,记做C

Number Of Consumers:表示在该队列上还有多少消费者在等待接受消息

Number Of Pending Messages:表示还有多少条消息没有被消费,实际上是表示消息的积压程度,就是P-C

VII. 签收

签收就是消费者接受到消息后,需要告诉消息服务器,我收到消息了。当消息服务器收到回执后,本条消息将失效。因此签收将对PTP模式产生很大影响。如果消费者收到消息后,并不签收,那么本条消息继续有效,很可能会被其他消费者消费掉!

AUTO_ACKNOWLEDGE:表示在消费者receive消息的时候自动的签收

CLIENT_ACKNOWLEDGE:表示消费者receive消息后必须手动的调用acknowledge()方法进行签收

DUPS_OK_ACKNOWLEDGE:签不签收无所谓了,只要消费者能够容忍重复的消息接受,当然这样会降低Session的开销


2). request/reply模型

I. 实现思路


image.png


Client的Producer发出一个JMS message形式的request,request上附加了一些额外的属性:


correlation ID(用来和返回的correlation ID对比进行验证),

JMSReplyTo属性(放置jms message的destination,这样worker的Consumer获得jms message就能得到destination)


Worker的consumer收到requset,处理request并用producer发出reply,destination就从requset的JMSReplyTo属性中得到。


II. Server代码


public class Server implements MessageListener {
    // 经纪人链接
    private static final String BROKER_URL = "tcp://localhost:61616";
    // 请求队列
    private static final String REQUEST_QUEUE = "requestQueue";
    // 经纪人服务
    private BrokerService brokerService;
    // 会话
    private Session session;
    // 生产者
    private MessageProducer producer;
    // 消费者
    private MessageConsumer consumer;
    private void start() throws Exception {
        createBroker();
        setUpConsumer();
    }
    /**
     * 创建经纪人
     * @throws Exception 
     */
    private void createBroker() throws Exception {
        // 创建经纪人服务
        brokerService = new BrokerService();
        // 设置是否持久化
        brokerService.setPersistent(false);
        // 设置是否使用JMX
        brokerService.setUseJmx(false);
        // 添加链接
        brokerService.addConnector(BROKER_URL);
        // 启动
        brokerService.start();
    }
    /**
     * 设置消费者
     * @throws JMSException 
     */
    private void setUpConsumer() throws JMSException {
        // 创建连接工厂
        ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory(BROKER_URL);
        // 创建连接
        Connection connection = connectionFactory.createConnection();
        // 启动连接
        connection.start();
        // 创建Session
        session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
        // 创建队列
        Destination adminQueue = session.createQueue(REQUEST_QUEUE);
        // 创建生产者
        producer = session.createProducer(null);
        // 设置持久化模式
        producer.setDeliveryMode(DeliveryMode.NON_PERSISTENT);
        // 创建消费者
        consumer = session.createConsumer(adminQueue);
        // 消费者设置消息监听
        consumer.setMessageListener(this);
    }
    public void stop() throws Exception {
        producer.close();
        consumer.close();
        session.close();
        brokerService.stop();
    }
    @Override
    public void onMessage(Message message) {
        try {
            // 创建新消息
            TextMessage response = this.session.createTextMessage();
            // 判断消息是否是文本消息
            if (message instanceof TextMessage) {
                // 强转为文本消息 
                TextMessage textMessage = (TextMessage) message;
                // 获取消息内容
                String text = textMessage.getText();
                // 设置消息
                response.setText(handleRequest(text));
            }
            response.setJMSCorrelationID(message.getJMSCorrelationID());
            producer.send(message.getJMSReplyTo(), response);
        } catch (JMSException e) {
            e.printStackTrace();
        }
    }
    /**
     * 构建消息内容
     * @param text 文本
     * @return
     */
    private String handleRequest(String text) {
        return "Response to '" + text + "'";
    }
    public static void main(String[] args) throws Exception {
        Server server = new Server();
        // 启动
        server.start();
        System.out.println();
        System.out.println("Press any key to stop the server");
        System.out.println();
        System.in.read();
        server.stop();
    }
}


III. Client代码


public class Client implements MessageListener {
    // 经纪人链接
    private static final String BROKER_URL = "tcp://localhost:61616";
    // 请求队列
    private static final String REQUEST_QUEUE = "requestQueue";
    // 连接
    private Connection connection;
    // 会话
    private Session session;
    // 生产者
    private MessageProducer producer;
    // 消费者
    private MessageConsumer consumer;
    // 请求队列
    private Queue tempDest;
    public void start() throws JMSException {
        // 连接工厂
        ActiveMQConnectionFactory activeMQConnectionFactory = new ActiveMQConnectionFactory(BROKER_URL);
        // 创建连接
        connection = activeMQConnectionFactory.createConnection();
        // 开启连接
        connection.start();
        // 创建会话
        session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
        // 创建队列
        Destination adminQueue = session.createQueue(REQUEST_QUEUE);
        // 创建生产者
        producer = session.createProducer(adminQueue);
        // 设置持久化模式
        producer.setDeliveryMode(DeliveryMode.NON_PERSISTENT);
        // 创建模板队列
        tempDest = session.createTemporaryQueue();
        // 创建消费者
        consumer = session.createConsumer(tempDest);
        // 设置消息监听
        consumer.setMessageListener(this);      
    }
    /**
     * 停止
     * @throws JMSException 
     */
    public void stop() throws JMSException {
        producer.close();
        consumer.close();
        session.close();
    }
    /**
     * 请求
     * @param request
     * @throws JMSException 
     */
    public void request(String request) throws JMSException {
        System.out.println("Request: " + request);
        // 创建文本消息
        TextMessage textMessage = session.createTextMessage();
        // 设置文本内容
        textMessage.setText(request);
        // 设置回复
        textMessage.setJMSReplyTo(tempDest);
        // 获取UUID
        String correlationId = UUID.randomUUID().toString();
        // 设置JMS id
        textMessage.setJMSCorrelationID(correlationId);
        // 发送消息
        this.producer.send(textMessage);
    }
    @Override
    public void onMessage(Message message) {
        try {
            System.out.println("Received response for: " + ((TextMessage)message).getText());
        } catch (JMSException e) {
            e.printStackTrace();
        }
    }
    public static void main(String[] args) throws JMSException, InterruptedException {
        Client client = new Client();
        // 启动
        client.start();
        int i = 0;
        while(i++ < 10) {
            client.request("REQUEST- " + i);
        }
        Thread.sleep(3000);
        client.stop();
    }
}


IV. 测试

启动Server


image.png


启动Client


image.png



目录
相关文章
|
11月前
|
消息中间件 存储 Kafka
分布式消息中间件设计与实现
本文深入探讨了消息中间件的核心功能实现与高并发、高可用设计。在生产者设计中,涵盖消息构造、序列化、路由策略及可靠性保障(如ACK机制)。消费者部分分析了拉取/推送模式、分区分配与消息确认机制。同时,Broker作为核心组件,负责消息路由、存储和投递,并通过索引技术实现快速检索。 高并发设计方面,重点讨论了文件存储(顺序写入、分段存储)、日志结构存储及负载均衡策略(如哈希分区、轮询分区)。为确保高可用性,文章详细解析了主从复制、故障转移机制以及同城/异地多活容灾方案。
392 13
|
11月前
|
消息中间件 存储 中间件
分布式消息中间件基础
消息中间件是一种基于异步消息传递的分布式系统通信工具,核心功能包括消息传输、存储、路由与投递,能够实现系统解耦、异步处理和流量削峰。其主要组件包括生产者、消费者、Broker、主题/队列等,支持点对点和发布-订阅两种消息模型。主流中间件如Kafka(高吞吐)、RabbitMQ(灵活路由)、RocketMQ(事务支持)各有特色,适用于不同场景。此外,中间件还涉及多种协议(AMQP、MQTT等)、可靠性传输机制(持久化、确认机制)、顺序性与重复性问题解决以及事务支持(两阶段提交、本地消息表等)。选择中间件需根据业务需求权衡性能、功能和运维成本。
503 6
|
12月前
|
人工智能 安全 Java
智慧工地源码,Java语言开发,微服务架构,支持分布式和集群部署,多端覆盖
智慧工地是“互联网+建筑工地”的创新模式,基于物联网、移动互联网、BIM、大数据、人工智能等技术,实现对施工现场人员、设备、材料、安全等环节的智能化管理。其解决方案涵盖数据大屏、移动APP和PC管理端,采用高性能Java微服务架构,支持分布式与集群部署,结合Redis、消息队列等技术确保系统稳定高效。通过大数据驱动决策、物联网实时监测预警及AI智能视频监控,消除数据孤岛,提升项目可控性与安全性。智慧工地提供专家级远程管理服务,助力施工质量和安全管理升级,同时依托可扩展平台、多端应用和丰富设备接口,满足多样化需求,推动建筑行业数字化转型。
391 5
|
存储 人工智能 算法
解锁分布式文件分享的 Java 一致性哈希算法密码
在数字化时代,文件分享成为信息传播与协同办公的关键环节。本文深入探讨基于Java的一致性哈希算法,该算法通过引入虚拟节点和环形哈希空间,解决了传统哈希算法在分布式存储中的“哈希雪崩”问题,确保文件分配稳定高效。文章还展示了Java实现代码,并展望了其在未来文件分享技术中的应用前景,如结合AI优化节点布局和区块链增强数据安全。
|
存储 缓存 Java
Java中的分布式缓存与Memcached集成实战
通过在Java项目中集成Memcached,可以显著提升系统的性能和响应速度。合理的缓存策略、分布式架构设计和异常处理机制是实现高效缓存的关键。希望本文提供的实战示例和优化建议能够帮助开发者更好地应用Memcached,实现高性能的分布式缓存解决方案。
275 9
|
存储 分布式计算 Hadoop
基于Java的Hadoop文件处理系统:高效分布式数据解析与存储
本文介绍了如何借鉴Hadoop的设计思想,使用Java实现其核心功能MapReduce,解决海量数据处理问题。通过类比图书馆管理系统,详细解释了Hadoop的两大组件:HDFS(分布式文件系统)和MapReduce(分布式计算模型)。具体实现了单词统计任务,并扩展支持CSV和JSON格式的数据解析。为了提升性能,引入了Combiner减少中间数据传输,以及自定义Partitioner解决数据倾斜问题。最后总结了Hadoop在大数据处理中的重要性,鼓励Java开发者学习Hadoop以拓展技术边界。
503 7
|
Java 数据库
在Java中使用Seata框架实现分布式事务的详细步骤
通过以上步骤,利用 Seata 框架可以实现较为简单的分布式事务处理。在实际应用中,还需要根据具体业务需求进行更详细的配置和处理。同时,要注意处理各种异常情况,以确保分布式事务的正确执行。
|
消息中间件 Java Kafka
在Java中实现分布式事务的常用框架和方法
总之,选择合适的分布式事务框架和方法需要综合考虑业务需求、性能、复杂度等因素。不同的框架和方法都有其特点和适用场景,需要根据具体情况进行评估和选择。同时,随着技术的不断发展,分布式事务的解决方案也在不断更新和完善,以更好地满足业务的需求。你还可以进一步深入研究和了解这些框架和方法,以便在实际应用中更好地实现分布式事务管理。
1233 161
|
存储 NoSQL Java
Java调度任务如何使用分布式锁保证相同任务在一个周期里只执行一次?
【10月更文挑战第29天】Java调度任务如何使用分布式锁保证相同任务在一个周期里只执行一次?
541 1
|
消息中间件 中间件 数据库
NServiceBus:打造企业级服务总线的利器——深度解析这一面向消息中间件如何革新分布式应用开发与提升系统可靠性
【10月更文挑战第9天】NServiceBus 是一个面向消息的中间件,专为构建分布式应用程序设计,特别适用于企业级服务总线(ESB)。它通过消息队列实现服务间的解耦,提高系统的可扩展性和容错性。在 .NET 生态中,NServiceBus 提供了强大的功能,支持多种传输方式如 RabbitMQ 和 Azure Service Bus。通过异步消息传递模式,各组件可以独立运作,即使某部分出现故障也不会影响整体系统。 示例代码展示了如何使用 NServiceBus 发送和接收消息,简化了系统的设计和维护。
343 3