线程池与mq的简单结合使用

线程池与mq的简单结合使用线程池与mq的简单结合使用

大家好,又见面了,我是你们的朋友全栈君。本文考虑发送方和接收方有多个线程发布消息和多个线程接收消息的情况:
 

1.生产者

package com.activemq3;
import javax.jms.Connection;
import javax.jms.ConnectionFactory;
import javax.jms.Destination;
import javax.jms.JMSException;
import javax.jms.MessageProducer;
import javax.jms.Session;
import javax.jms.TextMessage;
import javax.naming.InitialContext;

import org.apache.activemq.ActiveMQConnection;
import org.apache.activemq.ActiveMQConnectionFactory;

/**
 * 消息生产者-消息发布者(多线程发送)
 * 
 */
public class JMSProducerThread {

    private static final String USERNAME=ActiveMQConnection.DEFAULT_USER; // 默认的连接用户名
    private static final String PASSWORD=ActiveMQConnection.DEFAULT_PASSWORD; // 默认的连接密码
    private static final String BROKEURL=ActiveMQConnection.DEFAULT_BROKER_URL; // 默认的连接地址
    private static final int SENDNUM=10; // 发送的消息数量
    ConnectionFactory connectionFactory=null; // 连接工厂
    private Connection connection = null;
    private Session session = null;
    private Destination destination=null; // 消息的目的地
    public void init(){
        // 实例化连接工厂
        connectionFactory=new ActiveMQConnectionFactory(JMSProducerThread.USERNAME, JMSProducerThread.PASSWORD, JMSProducerThread.BROKEURL);
        try {
            connection=connectionFactory.createConnection(); // 通过连接工厂获取连接
            connection.start();
        } catch (JMSException e) {
            e.printStackTrace();
        }
    }
    public void produce(){
        try {
            MessageProducer messageProducer; // 消息生产者
            session=connection.createSession(Boolean.TRUE, Session.AUTO_ACKNOWLEDGE); // 创建Session
            destination=session.createQueue("queue1");
            messageProducer=session.createProducer(destination); // 创建消息生产者
            for(int i=0;i<JMSProducerThread.SENDNUM;i++){
                TextMessage message=session.createTextMessage("ActiveMQ中"+Thread.currentThread().getName()+"线程发送的数据"+":"+i);
                System.out.println(Thread.currentThread().getName()+"线程"+"发送消息:"+"ActiveMQ 发布的消息"+":"+i);
                messageProducer.send(message);
                session.commit();
            }
        } catch (JMSException e) {
            e.printStackTrace();
        }
    }


    /**
     * 发送消息
     * @param session
     * @param messageProducer
     * @throws Exception
     */
    public static void sendMessage(Session session,MessageProducer messageProducer)throws Exception{
        for(int i=0;i<JMSProducerThread.SENDNUM;i++){
            TextMessage message=session.createTextMessage("ActiveMQ 发送的消息"+i);
            System.out.println("发送消息:"+"ActiveMQ 发布的消息"+i);
            messageProducer.send(message);
        }
    }

}

2.消费者

package com.activemq3;

import javax.jms.Connection;
import javax.jms.ConnectionFactory;
import javax.jms.Destination;
import javax.jms.JMSException;
import javax.jms.MessageConsumer;
import javax.jms.MessageProducer;
import javax.jms.Session;
import javax.jms.TextMessage;
import javax.naming.InitialContext;

import org.apache.activemq.ActiveMQConnection;
import org.apache.activemq.ActiveMQConnectionFactory;

/**
 * 消息生产者-消息发布者(多线程发送消息)
 *
 */
public class JMSConsumerThread {

    private static final String USERNAME=ActiveMQConnection.DEFAULT_USER; // 默认的连接用户名
    private static final String PASSWORD=ActiveMQConnection.DEFAULT_PASSWORD; // 默认的连接密码
    private static final String BROKEURL=ActiveMQConnection.DEFAULT_BROKER_URL; // 默认的连接地址

    ConnectionFactory connectionFactory=null; // 连接工厂
    private Connection connection = null;
    private Session session = null;
    private Destination destination=null; // 消息的目的地
    public void init(){
        // 实例化连接工厂
        connectionFactory=new ActiveMQConnectionFactory(JMSConsumerThread.USERNAME, JMSConsumerThread.PASSWORD, JMSConsumerThread.BROKEURL);
        try {
            connection=connectionFactory.createConnection(); // 通过连接工厂获取连接
            connection.start();
        } catch (JMSException e) {
            e.printStackTrace();
        }
    }
    public void consumer(){
        MessageConsumer messageConsumer; // 消息的消费者

        try {
            session=connection.createSession(Boolean.FALSE, Session.AUTO_ACKNOWLEDGE); // 创建Session
            destination=session.createQueue("queue1");
            messageConsumer=session.createConsumer(destination); // 创建消息消费者
            messageConsumer.setMessageListener(new Listener3()); // 注册消息监听
        } catch (JMSException e) {
            e.printStackTrace();
        }
    }


}

package com.activemq3;

import java.util.Random;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

import javax.jms.JMSException;
import javax.jms.Message;
import javax.jms.MessageListener;
import javax.jms.TextMessage;

/**
 * 消息监听(多线程接收)
 *
 */
public class Listener3 implements MessageListener{
    private ExecutorService threadPool =Executors.newFixedThreadPool(8);
    @Override
    public void onMessage(final Message message) {
    threadPool.execute(new Runnable() {

            @Override
            public void run() {
                try {
                    try {
                        Thread.sleep(new Random().nextInt(2)*500);
                    } catch (InterruptedException e) {
                        e.printStackTrace();
                    }
                    System.out.println("接收线程..."+Thread.currentThread().getName()+"收到的消息是:"+((TextMessage)message).getText());
                } catch (JMSException e) {
                    e.printStackTrace();
                }               
            }
        }); 
    }
}

生产者执行发送的主函数:
 

在main方法中创建了一个固定3个 线程的线程池,去处理5个线程任务。

package com.activemq3;

import java.util.Random;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
/**
 * 多线程发送消息
 */

public class MainProducer {
          public static void main(String[] args) {

              ExecutorService threadPool=Executors.newFixedThreadPool(3);
              for (int i = 0; i < 5; i++) {
                  threadPool.submit(new Runnable() {

                    @Override
                    public void run() {
                        JMSProducerThread jph=new JMSProducerThread();
                        try {
                            Thread.sleep(new Random().nextInt(5)*500);
                        } catch (Exception e) {
                            e.printStackTrace();
                        }
                        jph.init();
                        jph.produce();                      
                    }
                });

            }
             /* 
            for (int i = 0; i <5; i++) {

               new Thread(new Runnable() {

                @Override
                public void run() {
                    JMSProducerThread js=new JMSProducerThread();
                    try {
                        Thread.sleep(new Random().nextInt(5)*500);
                    } catch (InterruptedException e) {
                        e.printStackTrace();
                    }
                    js.init();
                    js.produce();

                }
            }).start();

            }*/         
            }
}

消费者执行发送的主函数:
 

消费者中监听器中创建了固定8个线程的线程池去接收消息。

package com.activemq3;
public class MainConsumer {
          public static void main(String[] args) {
            JMSConsumerThread jch=new JMSConsumerThread();
            jch.init();
            jch.consumer();             
  }
}

版权声明:本文内容由互联网用户自发贡献,该文观点仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌侵权/违法违规的内容, 请发送邮件至 举报,一经查实,本站将立刻删除。

发布者:全栈程序员-用户IM,转载请注明出处:https://javaforall.cn/106177.html原文链接:https://javaforall.cn

【正版授权,激活自己账号】: Jetbrains全家桶Ide使用,1年售后保障,每天仅需1毛

【官方授权 正版激活】: 官方授权 正版激活 支持Jetbrains家族下所有IDE 使用个人JB账号...

(0)


相关推荐

  • vbs整人小代码大全

    事前准备:创建一个文件夹,后缀名为vbs,在里面编东西即可参考很多资料,如有不足请谅解文章目录vbs整人代码大全1.循环类1.1.一直说Youarefoolish!1.2.打开无数个计算器(慎用)1.3.重复说话100次1.4.不是循环但是很烦人1.5.给我按100次回车1.6.非常恶心的Alt+F42.问答类2.1.简单问答2.2.简单问答2.3.恐怖问答3.威胁类3.1.挟持回答3.2.挟持回答3.3虚惊一场结束vbs的代码完结vbs整人代码大全1.循环类1.1.一直说Youaref

  • 安捷伦示波器使用说明书_安捷伦信号发生器使用方法

    安捷伦示波器使用说明书_安捷伦信号发生器使用方法本帖最后由god_blessme于2017-9-1913:45编辑小弟最近在搞一个程序,是要读取安捷伦示波器每一屏数据并储存,网上貌似对于tek示波器连接的比较多,安捷伦的超级少,所以大部分是自己看着改的命令,现在碰到的问题很奇葩,运行程序后一个figure显示的数据是正确的,一个figure显示的是错误的。然而我在循环里把waveform_YIncrement变量的注释去掉的话,fig…

    2022年10月12日
  • Android ListView滚动条配置完全解析

    Android ListView滚动条配置完全解析详细介绍ListView中和滚动条相关的所有配置,包括普通滚动条和快速滚动条。

  • java测试面试问题_struts2面试题

    java测试面试问题_struts2面试题Javashiro面试题1、简单介绍一下Shiro框架?ApacheShiro是Java的一个安全框架。使用Shiro可以非常容易的开发出足够好的应用。其不仅可以用在JavaSE环境,也可以用在JavaEE环境。Shiro可以帮助我们完成功能:认证、授权、加密、会话管理、与Web集成、缓存等。三个核心组件:Subject,SecurityManager和Realms。●Subject:即“当…

    2022年10月14日
  • Idea 替换 区分大小写「建议收藏」

    Idea 替换 区分大小写「建议收藏」全文替换的时候,没有忽略大小写,导致替换类名和参数不一样猜测idea有没有区分大小写功能biu了一下Shift+Ctrl+RCc有点像,试试到达要求

  • 移动端禁用长按复制js兼容css样式_手机为什么长按不能复制

    移动端禁用长按复制js兼容css样式_手机为什么长按不能复制添加全局禁止选择文本的CSS属性*{-webkit-touch-callout:none;/*系统默认菜单*/-webkit-user-select:none;/*webkit浏览器*//*noinspectionCssUnknownProperty*/-khtml-user-select:none;/*早期浏览器*/-moz-user-select:none;/*火狐浏览器*/-ms-user-select:

发表回复

您的电子邮箱地址不会被公开。

关注全栈程序员社区公众号