`
Donald_Draper
  • 浏览: 999846 次
社区版块
存档分类
最新评论

ActiveMQ会话初始化详解

阅读更多
JMS(ActiveMQ) PTP和PUB/SUB模式实例:http://donald-draper.iteye.com/blog/2347445
ActiveMQ连接工厂、连接详解:http://donald-draper.iteye.com/blog/2348070
ActiveMQ会话初始化:http://donald-draper.iteye.com/blog/2348341
ActiveMQ生产者:http://donald-draper.iteye.com/blog/2348381
ActiveMQ消费者:http://donald-draper.iteye.com/blog/2348389
ActiveMQ启动过程详解:http://donald-draper.iteye.com/blog/2348399
ActiveMQ Broker发送消息给消费者过程详解:http://donald-draper.iteye.com/blog/2348440
Spring与ActiveMQ的集成:http://donald-draper.iteye.com/blog/2347638
Spring与ActiveMQ的集成详解一:http://donald-draper.iteye.com/blog/2348449
Spring与ActiveMQ的集成详解二:http://donald-draper.iteye.com/blog/2348461
在上面一篇我们说过,ActiveMQ的ActiveMQConnectionFactory,ActiveMQConnection,TcpTransport,
ActiveMQConnectionFactory的创建过程,主要为初始化为broker url,用户密码,是否压缩、异步发送消息、支持消息优先级、非阻塞传输,最大线程数,生产窗口大小等属性;从ActiveMQConnectionFactory创建连接,首先通过TcpTransportFacotory创建TcpTransport,然后保证成待锁机制的MutexTransport,最后包装成ResponseCorrelator;根据TcpTransport和ActiveMQConnection状态管理器JMSStatsImpl创建ActiveMQConnection,创建ActiveMQConnection过程中,主要是是否异步分发消息,线程执行器,连接状态管理器,调度器等;然后
设置连接用户密码通过ConnectionInfo,配置是否支持消息优先级、非阻塞传输,最大线程数,生产窗口大小,Transport监听器transportListener;最后启动TcpTransport和Connection,启动TcpTransport主要是初始化socket,ip,端口,输入输出缓存区,输入输出流DataI/OnputStream,启动连接主要启动会话ActiveMQSession。

实例主要生产者代码片段:
ConnectionFactory :连接工厂,JMS 用它创建连接  
ConnectionFactory connectionFactory = new ActiveMQConnectionFactory(user,password,url);  
Connection :JMS 客户端到JMS Provider 的连接  
Connection connection = connectionFactory.createConnection();  
Connection 启动  
connection.start();  
System.out.println("Connection is start...");  
//创建一个session
//第一个参数:是否支持事务,如果为true,则会忽略第二个参数,被jms服务器设置为SESSION_TRANSACTED
//第二个参数为false时,paramB的值可为Session.AUTO_ACKNOWLEDGE,Session.CLIENT_ACKNOWLEDGE,DUPS_OK_ACKNOWLEDGE其中一个。
//Session.AUTO_ACKNOWLEDGE为自动确认,客户端发送和接收消息不需要做额外的工作。哪怕是接收端发生异常,也会被当作正常发送成功。
//Session.CLIENT_ACKNOWLEDGE为客户端确认。客户端接收到消息后,必须调用javax.jms.Message的acknowledge方法。jms服务器才会当作发送成功,并删除消息。
//DUPS_OK_ACKNOWLEDGE允许副本的确认模式。一旦接收方应用程序的方法调用从处理消息处返回,会话对象就会确认消息的接收;而且允许重复确认。
Session session = connection.createSession(Boolean.TRUE,Session.AUTO_ACKNOWLEDGE);  
Queue :消息的目的地;消息发送给谁.  
Queue  destination = session.createQueue(qname);  
MessageProducer:消息发送者  
MessageProducer producer = session.createProducer(destination);  
//设置生产者的模式,有两种可选
//DeliveryMode.PERSISTENT 当activemq关闭的时候,队列数据将会被保存
//DeliveryMode.NON_PERSISTENT 当activemq关闭的时候,队列里面的数据将会被清空 
producer.setDeliveryMode(DeliveryMode.PERSISTENT);  
构造消息,此处写死,项目就是参数,或者方法获取  
sendMessage(session, producer);  
session.commit();  
connection.close();  


今天来看一下,ActiveMQSession会话,消息队列和订阅主题,生产者及发送消息
从下面这句开始:
Session session = connection.createSession(Boolean.TRUE,Session.AUTO_ACKNOWLEDGE);  

//ActiveMQConnection
 public Session createSession(boolean transacted, int acknowledgeMode)
        throws JMSException
    {
        //检查连接有没关闭
        checkClosedOrFailed();
        ensureConnectionInfoSent();
        if(!transacted)
        {   
	    //如果transacted为非事务,而acknowledgeMode为事务SESSION_TRANSACTED,抛出异常
            if(acknowledgeMode == 0)
                throw new JMSException("acknowledgeMode SESSION_TRANSACTED cannot be used for an non-transacted Session");
            //acknowledgeMode不在0-3范围之内,这
	    if(acknowledgeMode < 0 || acknowledgeMode > 4)
                throw new JMSException((new StringBuilder()).append("invalid acknowledgeMode: ").append(acknowledgeMode).append(". Valid values are Session.AUTO_ACKNOWLEDGE (1), ").append("Session.CLIENT_ACKNOWLEDGE (2), Session.DUPS_OK_ACKNOWLEDGE (3), ActiveMQSession.INDIVIDUAL_ACKNOWLEDGE (4) or for transacted sessions Session.SESSION_TRANSACTED (0)").toString());
        }
        return new ActiveMQSession(this, getNextSessionId(), transacted ? 0 : acknowledgeMode != 0 ? acknowledgeMode : 1, isDispatchAsync(), isAlwaysSessionAsync());
    }
//来看ActiveMQSession的构造
public class ActiveMQSession
    implements Session, QueueSession, TopicSession, StatsCapable, ActiveMQDispatcher
{
    public static final int INDIVIDUAL_ACKNOWLEDGE = 4;
    public static final int MAX_ACK_CONSTANT = 4;
    private static final Logger LOG = LoggerFactory.getLogger(org/apache/activemq/ActiveMQSession);
    private final ThreadPoolExecutor connectionExecutor;//连接线程执行器
    protected int acknowledgementMode; //通知模式
    protected final ActiveMQConnection connection;//MQ连接
    protected final SessionInfo info;//会话信息
    protected final LongSequenceGenerator consumerIdGenerator;//消费者id产生器
    protected final LongSequenceGenerator producerIdGenerator;//生产者id产生器
    protected final LongSequenceGenerator deliveryIdGenerator;
    protected final ActiveMQSessionExecutor executor;
    protected final AtomicBoolean started;  //是否启动
    protected final CopyOnWriteArrayList consumers;//消费者
    protected final CopyOnWriteArrayList producers;//生产者
    protected boolean closed;
    private volatile boolean synchronizationRegistered;
    protected boolean asyncDispatch;
    protected boolean sessionAsyncDispatch;
    protected final boolean debug;
    protected final Object sendMutex;//发送互斥锁
    protected final Object redeliveryGuard;
    private final AtomicBoolean clearInProgress;
    private MessageListener messageListener;//消息监听器
    private final JMSSessionStatsImpl stats;
    private TransactionContext transactionContext;
    private DeliveryListener deliveryListener;
    private MessageTransformer transformer;
    private BlobTransferPolicy blobTransferPolicy;
    private long lastDeliveredSequenceId;
    final AtomicInteger clearRequestsCounter;
    protected ActiveMQSession(ActiveMQConnection connection, SessionId sessionId, int acknowledgeMode, boolean asyncDispatch, boolean sessionAsyncDispatch)
        throws JMSException
    {
        consumerIdGenerator = new LongSequenceGenerator();
        producerIdGenerator = new LongSequenceGenerator();
        deliveryIdGenerator = new LongSequenceGenerator();
        started = new AtomicBoolean(false);//启动状态
        consumers = new CopyOnWriteArrayList();//消费者
        producers = new CopyOnWriteArrayList();//生产者
        sendMutex = new Object();//发送互斥锁
        redeliveryGuard = new Object();
        clearInProgress = new AtomicBoolean();
        lastDeliveredSequenceId = -2L;
        clearRequestsCounter = new AtomicInteger(0);
        debug = LOG.isDebugEnabled();
        this.connection = connection;
        acknowledgementMode = acknowledgeMode;//消息确认模式
        this.asyncDispatch = asyncDispatch;
        this.sessionAsyncDispatch = sessionAsyncDispatch;
        info = new SessionInfo(connection.getConnectionInfo(), sessionId.getValue());
        setTransactionContext(new TransactionContext(connection));//设置连接事务上下文
        stats = new JMSSessionStatsImpl(producers, consumers);
        this.connection.asyncSendPacket(info);//异步发送会话信息
        setTransformer(connection.getTransformer());
        setBlobTransferPolicy(connection.getBlobTransferPolicy());
        connectionExecutor = connection.getExecutor();//获取连接执行器
        executor = new ActiveMQSessionExecutor(this);//新建消息会话执行器
        connection.addSession(this);//将会话添加到ActiveMQ连接的的会话队列CopyOnWriteArrayList
        if(connection.isStarted())
	     //启动
            start();
    }
}

从上面可以看出会话的创建主要最的工作是,初始化消费者,生产者id产生器,会话消费者,生产者队列,消息确认模式,是否异步分发,设置连接事务上下文,异步发送会话信息,新建消息会话执行器,会话添加到ActiveMQConnection的会话队列CopyOnWriteArrayList。
结合上一篇和这一篇说以一下ActiveMQConnection与ActiveMQSession,ActiveMQMessageConsumer,
ActiveMQMessageProducer的关系,连接管理会话(1-n),会话管理消息者与生产者(1-n)。

下面再看一下会话的启动
启动
protected void start()
        throws JMSException
    {
        started.set(true);
        ActiveMQMessageConsumer c;
	//启动消费者
        for(Iterator iter = consumers.iterator(); iter.hasNext(); c.start())
            c = (ActiveMQMessageConsumer)iter.next();
        //启动消息会话执行器
        executor.start();
    }

会话启动分两部分,启动消费者,让消费者消费消息,启动会话执行器
先来看启动消费者
来看ActiveMQMessageConsumer的启动
public class ActiveMQMessageConsumer
    implements MessageAvailableConsumer, StatsCapable, ActiveMQDispatcher
{
   protected final ActiveMQSession session;
    protected final ConsumerInfo info;//消费者信息
    protected final MessageDispatchChannel unconsumedMessages;//消息分发通道(未消费消费的消息)
    protected final LinkedList deliveredMessages = new LinkedList();//需要传输消息
    private PreviouslyDeliveredMap previouslyDeliveredMessages;
    private int deliveredCounter;//传输消息计数器
    private int additionalWindowSize;
    private long redeliveryDelay;//消息传输延时
    private int ackCounter;//消息回复计时器
    private int dispatchedCount;//消息分发计时器
    private final AtomicReference messageListener = new AtomicReference();//消息监听器引用
    private final JMSConsumerStatsImpl stats;//消费者状态信息管理器
    private final String selector;选择器
    private boolean synchronizationRegistered;
    private final AtomicBoolean started = new AtomicBoolean(false);//启动状态
    private MessageAvailableListener availableListener;
    private RedeliveryPolicy redeliveryPolicy;//传输策略
    private boolean optimizeAcknowledge;
    private final AtomicBoolean deliveryingAcknowledgements = new AtomicBoolean();
    private ExecutorService executorService;//执行器
    private MessageTransformer transformer;
    private boolean clearDeliveredList;
    AtomicInteger inProgressClearRequiredFlag;
    private MessageAck pendingAck;//消息回复
    private long lastDeliveredSequenceId;
    private IOException failureError;
    private long optimizeAckTimestamp;
    private long optimizeAcknowledgeTimeOut;
    private long optimizedAckScheduledAckInterval;
    private Runnable optimizedAckTask;//消息回复优化任务
    private long failoverRedeliveryWaitPeriod;//当Master宕机时,消息传输等待时间
    private boolean transactedIndividualAck;
    private boolean nonBlockingRedelivery;//是否非阻塞传输
    private boolean consumerExpiryCheckEnabled;
   public void start()
        throws JMSException
    {
        if(unconsumedMessages.isClosed())
        {
            return;
        } else
        {
            started.set(true);
	    //启动消息分发通道
            unconsumedMessages.start();
	    //唤醒会话执行器
            session.executor.wakeup();
            return;
        }
    }


在ActiveMQMessageConsumer的构造中有对消息通道unconsumedMessages的初始化
如果支持消息优先级,则为SimplePriorityMessageDispatchChannel,否则则为
FifoMessageDispatchChannel
 if(session.connection.isMessagePrioritySupported())
            unconsumedMessages = new SimplePriorityMessageDispatchChannel();
 else
            unconsumedMessages = new FifoMessageDispatchChannel();


我们来看一下FifoMessageDispatchChannel
public class FifoMessageDispatchChannel
    implements MessageDispatchChannel
{
    private final Object mutex = new Object();//消息分发通道互斥量
    private final LinkedList list = new LinkedList();
    private boolean closed;
    private boolean running;//运行状态
    public void start()
    {
        synchronized(mutex)
        {
            running = true;
	    //唤醒所有等待消息分发锁的线程
            mutex.notifyAll();
        }
    }
}

再看SimplePriorityMessageDispatchChannel

public class SimplePriorityMessageDispatchChannel
    implements MessageDispatchChannel
{
    private static final Integer MAX_PRIORITY = Integer.valueOf(10);
    private final Object mutex = new Object();
    private final LinkedList lists[];
    private boolean closed;
    private boolean running;
    private int size;

   public void start()
    {
        synchronized(mutex)
        {
            running = true;
            mutex.notifyAll();
        }
    }
}

SimplePriorityMessageDispatchChannel的启动与FifoMessageDispatchChannel相同;

回到会话启动中的执行器启动
//ActiveMQSession
protected void start()
        throws JMSException
    {
        started.set(true);
        ActiveMQMessageConsumer c;
	//启动消费者
        for(Iterator iter = consumers.iterator(); iter.hasNext(); c.start())
            c = (ActiveMQMessageConsumer)iter.next();
        //启动消息会话执行器
        executor.start();
    }

//ActiveMQSessionExecutor
public class ActiveMQSessionExecutor
    implements Task
{
   private final ActiveMQSession session;//MQ会话
    private final MessageDispatchChannel messageQueue;//消息分发通道
    private boolean dispatchedBySessionPool;//是否依靠会话池分发消息
    private volatile TaskRunner taskRunner;//任务执行器
    private boolean startedOrWarnedThatNotStarted;
   
    ActiveMQSessionExecutor(ActiveMQSession session)
    {
        this.session = session;
        if(this.session.connection != null && this.session.connection.isMessagePrioritySupported())
            messageQueue = new SimplePriorityMessageDispatchChannel();
        else
            messageQueue = new FifoMessageDispatchChannel();
    }
   
}
//ActiveMQSessionExecutor
 synchronized void start()
    {
        //如果消息分发通道处于未启动状态,则启动消息分化通道
        if(!messageQueue.isRunning())
        {
	   //这个在上面已经看到
            messageQueue.start();
	    //如果有为消费的消息,则唤醒
            if(hasUncomsumedMessages())
                wakeup();
        }
    }
//判断是否有未消费的消息
public boolean hasUncomsumedMessages()
    {
        return !messageQueue.isClosed() && messageQueue.isRunning() && !messageQueue.isEmpty();
    }
 //唤醒
 public void wakeup()
    {
label0:
        {
            if(dispatchedBySessionPool)
                break MISSING_BLOCK_LABEL_134;
            if(!session.isSessionAsyncDispatch())
                break label0;
            TaskRunner taskRunner;
            try
            {
label1:
                {
                    taskRunner = this.taskRunner;
                    if(taskRunner != null)
                        break MISSING_BLOCK_LABEL_105;
                    synchronized(this)
                    {
                        if(this.taskRunner != null)
                            break MISSING_BLOCK_LABEL_90;
                        if(isRunning())
                            break label1;
                    }
                    return;
                }
            }
            catch(InterruptedException e)
            {
                Thread.currentThread().interrupt();
                break MISSING_BLOCK_LABEL_134;
            }
        }
	//创建会话任务执运行器
        this.taskRunner = session.connection.getSessionTaskRunner().createTaskRunner(this, (new StringBuilder()).append("ActiveMQ Session: ").append(session.getSessionId()).toString());
        taskRunner = this.taskRunner;
        activemqsessionexecutor;
        JVM INSTR monitorexit ;
        break MISSING_BLOCK_LABEL_105;
        exception;
        throw exception;
        taskRunner.wakeup();
        break MISSING_BLOCK_LABEL_134;
    }

从这一句可以看出通过ActiveMQConnection来创建的,为TaskRunnerFactory
session.connection.getSessionTaskRunner().createTaskRunner

public TaskRunnerFactory getSessionTaskRunner()
    {
        synchronized(this)
        {
            if(sessionTaskRunner == null)
            {
                sessionTaskRunner = new TaskRunnerFactory("ActiveMQ Session Task", 7, false, 1000, isUseDedicatedTaskRunner(), maxThreadPoolSize);
                sessionTaskRunner.setRejectedTaskHandler(rejectedTaskHandler);
            }
        }
        return sessionTaskRunner;
    }

再来看为TaskRunnerFactory如何创建任务运行器
public class TaskRunnerFactory
    implements Executor
{ 
    private ExecutorService executor;//执行器
    private int maxIterationsPerRun;
    private String name;
    private int priority;//优先级
    private boolean daemon;
    private final AtomicLong id;
    private boolean dedicatedTaskRunner;
    private long shutdownAwaitTermination;
    private final AtomicBoolean initDone;//是否初始化
    private int maxThreadPoolSize;//最大线程池大小
    private RejectedExecutionHandler rejectedTaskHandler;//当线程达到,线程池大小的拒绝策略
    private ClassLoader threadClassLoader;
 public TaskRunnerFactory(String name, int priority, boolean daemon, int maxIterationsPerRun, boolean dedicatedTaskRunner, int maxThreadPoolSize)
    {
        id = new AtomicLong(0L);
        shutdownAwaitTermination = 30000L;
        initDone = new AtomicBoolean(false);
        this.maxThreadPoolSize = 2147483647;
        rejectedTaskHandler = null;
        this.name = name;
        this.priority = priority;
        this.daemon = daemon;
        this.maxIterationsPerRun = maxIterationsPerRun;
        this.dedicatedTaskRunner = dedicatedTaskRunner;
        this.maxThreadPoolSize = maxThreadPoolSize;
    }
}
TaskRunnerFactory

public TaskRunner createTaskRunner(Task task, String name)
    {
       //初始化
        init();
	//如果线程池为空,则创建线程执行器
        if(executor != null)
            return new PooledTaskRunner(executor, task, maxIterationsPerRun);
        else
            return new DedicatedTaskRunner(task, name, priority, daemon);
    }

  public void init()
    {
       //如果未初始化,则设置状态为已初始化
        if(initDone.compareAndSet(false, true))
        {
	   //判断是否用专业执行器,是则创建DedicatedTaskRunner
            if(dedicatedTaskRunner || "true".equalsIgnoreCase(System.getProperty("org.apache.activemq.UseDedicatedTaskRunner")))
                executor = null;
            else
            if(executor == null)
	       //创建默认执行器
                executor = createDefaultExecutor();
            LOG.debug("Initialized TaskRunnerFactory[{}] using ExecutorService: {}", name, executor);
        }
    }
    //创建默认执行器
protected ExecutorService createDefaultExecutor()
    {
        ThreadPoolExecutor rc = new ThreadPoolExecutor(0, getMaxThreadPoolSize(), getDefaultKeepAliveTime(), TimeUnit.SECONDS, new SynchronousQueue(), new ThreadFactory() {

            public Thread newThread(Runnable runnable)
            {
                String threadName = (new StringBuilder()).append(name).append("-").append(id.incrementAndGet()).toString();
                Thread thread = new Thread(runnable, threadName);
                thread.setDaemon(daemon);
                thread.setPriority(priority);
                if(threadClassLoader != null)
                    thread.setContextClassLoader(threadClassLoader);
                thread.setUncaughtExceptionHandler(new Thread.UncaughtExceptionHandler() {

                    public void uncaughtException(Thread t, Throwable e)
                    {
                        TaskRunnerFactory.LOG.error("Error in thread '{}'", t.getName(), e);
                    }

                    final _cls1 this$1;

                    
                    {
                        this$1 = _cls1.this;
                        super();
                    }
                });
                TaskRunnerFactory.LOG.trace("Created thread[{}]: {}", threadName, thread);
                return thread;
            }

            final TaskRunnerFactory this$0;

            
            {
                this$0 = TaskRunnerFactory.this;
                super();
            }
        });
        if(rejectedTaskHandler != null)
            rc.setRejectedExecutionHandler(rejectedTaskHandler);
        return rc;
    }

回到createTaskRunner
public TaskRunner createTaskRunner(Task task, String name)
    {
       //初始化
        init();
	//如果线程池为空,则创建线程执行器
        if(executor != null)
            return new PooledTaskRunner(executor, task, maxIterationsPerRun);
        else
            return new DedicatedTaskRunner(task, name, priority, daemon);
    }

我们来看PooledTaskRunner
class PooledTaskRunner
    implements TaskRunner
{
    private final int maxIterationsPerRun;//每次执行允许最大线程数
    private final Executor executor;//执行器
    private final Task task;
    private final Runnable runable;
    private boolean queued;//是否为队列
    private boolean shutdown;
    private boolean iterating;
    private volatile Thread runningThread;

    public PooledTaskRunner(Executor executor, final Task task, int maxIterationsPerRun)
    {
        this.executor = executor;
        this.maxIterationsPerRun = maxIterationsPerRun;
        this.task = task;
        runable = new Runnable() {

            public void run()
            {
                runningThread = Thread.currentThread();
                runTask();
                PooledTaskRunner.LOG.trace("Run task done: {}", task);
                runningThread = null;
                break MISSING_BLOCK_LABEL_70;
                Exception exception;
                exception;
                PooledTaskRunner.LOG.trace("Run task done: {}", task);
                runningThread = null;
                throw exception;
            }

            final Task val$task;
            final PooledTaskRunner this$0;

            
            {
                this$0 = PooledTaskRunner.this;
                task = task1;
                super();
            }
        };
    }
}

//PooledTaskRunner
  
 public void wakeup()
        throws InterruptedException
    {
label0:
        {
            synchronized(runable)
            {
                if(!queued && !shutdown)
                    break label0;
            }
            return;
        }
        queued = true;
        if(!iterating)
	    //执行runable
            executor.execute(runable);
        runnable;
        JVM INSTR monitorexit ;
          goto _L1
        exception;
        throw exception;
_L1:
    }

在PooledTaskRunner的构造中可以看到
runable = new Runnable() {

            public void run()
            {
                runningThread = Thread.currentThread();
                //运行任务
                runTask();
                PooledTaskRunner.LOG.trace("Run task done: {}", task);
                runningThread = null;
            }

            final Task val$task;
            final PooledTaskRunner this$0;

            
            {
                this$0 = PooledTaskRunner.this;
                task = task1;
                super();
            }
        };

//运行任务
final void runTask()
    {
        boolean done = false;
        int i = 0;
        do
        {
            if(i >= maxIterationsPerRun)
                break;
            LOG.trace("Running task iteration {} - {}", Integer.valueOf(i), task);
	    //运行任务的iterate,而任务实际为ActiveMQSessionExecutor
            if(!task.iterate())
            {
                done = true;
                break;
            }
            i++;
        } while(true);


//ActiveMQSessionExecutor implements Task
public boolean iterate()
    {
        for(Iterator i$ = session.consumers.iterator(); i$.hasNext();)
        {
            ActiveMQMessageConsumer consumer = (ActiveMQMessageConsumer)i$.next();
            if(consumer.iterate())
                return true;
        }
        //从消息队列中获取消息,消息出消息队列
        MessageDispatch message = messageQueue.dequeueNoWait();
        if(message == null)
        {
            return false;
        } else
        {
	    //分发消息
            dispatch(message);
            return !messageQueue.isEmpty();
        }
    }
//ActiveMQSessionExecutor
 void dispatch(MessageDispatch message)
    {
        Iterator i$ = session.consumers.iterator();
        do
        {
            if(!i$.hasNext())
                break;
            ActiveMQMessageConsumer consumer = (ActiveMQMessageConsumer)i$.next();
            ConsumerId consumerId = message.getConsumerId();
            if(!consumerId.equals(consumer.getConsumerId()))
                continue;
	    //遍历会话的消费者,分发消息,实际调用的是ActiveMQMessageConsumer的dispatch(message)
            consumer.dispatch(message);
            break;
        } while(true);
    }

//ActiveMQMessageConsumer
 public boolean iterate()
    { 
       //private final AtomicReference messageListener = new AtomicReference();
        MessageListener listener = (MessageListener)messageListener.get();
        if(listener != null)
        {
	   //如果消息消费者注册的消息监听器,则未消费消息通道则从消息队列取出消息,分发消息
            MessageDispatch md = unconsumedMessages.dequeueNoWait();
            if(md != null)
            {
	        //分发消息
                dispatch(md);
                return true;
            }
        }
        return false;
    }


再看ActiveMQMessageConsumer的dispatch(message)

public void dispatch(MessageDispatch md)
    {
        //包装分发消息
        ActiveMQMessage message = createActiveMQMessage(md);
        beforeMessageIsConsumed(md);
        try
        {
            boolean expired = isConsumerExpiryCheckEnabled() && message.isExpired();
            if(!expired)
	        //如果消息没有过期,则消息消费者通过消息监听器MessageListener消费消息
                listener.onMessage(message);
            afterMessageIsConsumed(md, expired);
        }
	if(!unconsumedMessages.isRunning())
            session.connection.rollbackDuplicate(this, md.getMessage());
	//将消息添加到未分发消息通道
        unconsumedMessages.enqueue(md);
	if(redeliveryExpectedInCurrentTransaction(md, true))
        {
            LOG.debug("{} tracking transacted redelivery {}", getConsumerId(), md.getMessage());
            if(transactedIndividualAck)
                immediateIndividualTransactedAck(md);
            else
	        //发送回复消息
                session.sendAck(new MessageAck(md, (byte)0, 1));
        } 
  }

总结:
会话的创建主要最的工作是,初始化消费者,生产者id产生器,会话消费者,生产者队列,消息确认模式,是否异步分发,设置连接事务上下文,异步发送会话信息,新建消息会话执行器,会话添加到ActiveMQConnection的会话队列CopyOnWriteArrayList;然后启动消费则,主要是启动消息分发通道,唤醒会话执行器ActiveMQSessionExecutor;最后启动会话执行器ActiveMQSessionExecutor,在启动会话执行器时,如果消息分发通道处于未启动状态,则启动消息分发通道,如果有未消费的消息,唤醒消息执行器,唤醒主要做的做工作是
ActiveMQConnection创建任务执行TaskRunnerFactory,有任务执行工厂TaskRunnerFactory,有任务执行工厂创建执行任务PooledTaskRunner,PooledTaskRunner是ActiveMQSessionExecutor的包装,PooledTaskRunner执行就是执行ActiveMQSessionExecutor
iterate的函数,这个过程主要是ActiveMQSessionExecutor从ActiveMQSession获取会话消费者consumer,然后遍历消费者,消费者通过MessageListener消费消息。
ActiveMQConnection与ActiveMQSession,ActiveMQMessageConsumer,
ActiveMQMessageProducer的关系,连接管理会话(1-n),会话管理消息者与生产者(1-n)。ActiveMQSession关联一个ActiveMQSessionExecutor,由会话执行器,消费者消费消息。


//MessageListener
public interface MessageListener
{
    public abstract void onMessage(Message message);
}


//MessageDispatch
public class MessageDispatch extends BaseCommand
{
    protected ConsumerId consumerId;//消费id
    protected ActiveMQDestination destination;//目的地
    protected Message message;//消息
    protected int redeliveryCounter;//消息传输计数器
    protected transient long deliverySequenceId;
    protected transient Object consumer;//消费者
    protected transient TransmitCallback transmitCallback;
    protected transient Throwable rollbackCause;
}

//Session
public interface Session
    extends Runnable
{   
    public static final int AUTO_ACKNOWLEDGE = 1;
    public static final int CLIENT_ACKNOWLEDGE = 2;
    public static final int DUPS_OK_ACKNOWLEDGE = 3;
    public static final int SESSION_TRANSACTED = 0;
     public abstract ObjectMessage createObjectMessage()
        throws JMSException;

    public abstract ObjectMessage createObjectMessage(Serializable serializable)
        throws JMSException;

    public abstract StreamMessage createStreamMessage()
        throws JMSException;

    public abstract TextMessage createTextMessage()
        throws JMSException;

    public abstract TextMessage createTextMessage(String s)
        throws JMSException;

    public abstract boolean getTransacted()
        throws JMSException;

    public abstract int getAcknowledgeMode()
        throws JMSException;

    public abstract void commit()
        throws JMSException;

    public abstract void rollback()
        throws JMSException;

    public abstract void close()
        throws JMSException;

    public abstract void recover()
        throws JMSException;

    public abstract MessageListener getMessageListener()
        throws JMSException;

    public abstract void setMessageListener(MessageListener messagelistener)
        throws JMSException;

    public abstract void run();

    public abstract MessageProducer createProducer(Destination destination)
        throws JMSException;

    public abstract MessageConsumer createConsumer(Destination destination)
        throws JMSException;

    public abstract MessageConsumer createConsumer(Destination destination, String s)
        throws JMSException;

    public abstract MessageConsumer createConsumer(Destination destination, String s, boolean flag)
        throws JMSException;

    public abstract Queue createQueue(String s)
        throws JMSException;

    public abstract Topic createTopic(String s)
        throws JMSException;

    public abstract TopicSubscriber createDurableSubscriber(Topic topic, String s)
        throws JMSException;

    public abstract TopicSubscriber createDurableSubscriber(Topic topic, String s, String s1, boolean flag)
        throws JMSException;

    public abstract QueueBrowser createBrowser(Queue queue)
        throws JMSException;

    public abstract QueueBrowser createBrowser(Queue queue, String s)
        throws JMSException;

    public abstract TemporaryQueue createTemporaryQueue()
        throws JMSException;

    public abstract TemporaryTopic createTemporaryTopic()
        throws JMSException;

    public abstract void unsubscribe(String s)
        throws JMSException;
}
0
0
分享到:
评论
1 楼 di1984HIT 2017-07-31  
xuexile ~~~~   

相关推荐

    Apache ActiveMQ教程 JMS 整合Tomcat

    1. **创建Connection**:初始化JMS连接,设定URL、用户名和密码。 2. **创建Session**:在连接基础上创建会话,定义是否支持事务以及确认模式。 3. **创建Destination**:指定消息目的地,即Queue或Topic名称,...

    Java面试框架高频问题2019

    - **初始化**:调用bean的初始化方法,如`@PostConstruct`或`init-method`属性指定的方法。 - **运行时**:bean处于可用状态。 - **销毁**:调用bean的销毁方法,如`@PreDestroy`或`destroy-method`属性指定的方法。...

    java面试题,180多页,绝对良心制作,欢迎点评,涵盖各种知识点,排版优美,阅读舒心

    类什么时候才被初始化: 58 类的初始化步骤: 59 【*JVM】什么是JVM线程死锁?JVM线程死锁,你该如何判断是因为什么?如果用VisualVM,dump线程信息出来,会有哪些信息? 59 【*JVM】查看jvm虚拟机里面堆、线程的...

    基于模糊故障树的工业控制系统可靠性分析与Python实现

    内容概要:本文探讨了模糊故障树(FFTA)在工业控制系统可靠性分析中的应用,解决了传统故障树方法无法处理不确定数据的问题。文中介绍了模糊数的基本概念和实现方式,如三角模糊数和梯形模糊数,并展示了如何用Python实现模糊与门、或门运算以及系统故障率的计算。此外,还详细讲解了最小割集的查找方法、单元重要度的计算,并通过实例说明了这些方法的实际应用场景。最后,讨论了模糊运算在处理语言变量方面的优势,强调了在可靠性分析中处理模糊性和优化计算效率的重要性。 适合人群:从事工业控制系统设计、维护的技术人员,以及对模糊数学和可靠性分析感兴趣的科研人员。 使用场景及目标:适用于需要评估复杂系统可靠性的场合,特别是在面对不确定数据时,能够提供更准确的风险评估。目标是帮助工程师更好地理解和预测系统故障,从而制定有效的预防措施。 其他说明:文中提供的代码片段和方法可用于初步方案验证和技术探索,但在实际工程项目中还需进一步优化和完善。

    风力发电领域双馈风力发电机(DFIG)Simulink模型的构建与电流电压波形分析

    内容概要:本文详细探讨了双馈风力发电机(DFIG)在Simulink环境下的建模方法及其在不同风速条件下的电流与电压波形特征。首先介绍了DFIG的基本原理,即定子直接接入电网,转子通过双向变流器连接电网的特点。接着阐述了Simulink模型的具体搭建步骤,包括风力机模型、传动系统模型、DFIG本体模型和变流器模型的建立。文中强调了变流器控制算法的重要性,特别是在应对风速变化时,通过实时调整转子侧的电压和电流,确保电流和电压波形的良好特性。此外,文章还讨论了模型中的关键技术和挑战,如转子电流环控制策略、低电压穿越性能、直流母线电压脉动等问题,并提供了具体的解决方案和技术细节。最终,通过对故障工况的仿真测试,验证了所建模型的有效性和优越性。 适用人群:从事风力发电研究的技术人员、高校相关专业师生、对电力电子控制系统感兴趣的工程技术人员。 使用场景及目标:适用于希望深入了解DFIG工作原理、掌握Simulink建模技能的研究人员;旨在帮助读者理解DFIG在不同风速条件下的动态响应机制,为优化风力发电系统的控制策略提供理论依据和技术支持。 其他说明:文章不仅提供了详细的理论解释,还附有大量Matlab/Simulink代码片段,便于读者进行实践操作。同时,针对一些常见问题给出了实用的调试技巧,有助于提高仿真的准确性和可靠性。

    基于西门子S7-200 PLC和组态王的八层电梯控制系统设计与实现

    内容概要:本文详细介绍了基于西门子S7-200 PLC和组态王软件构建的八层电梯控制系统。首先阐述了系统的硬件配置,包括PLC的IO分配策略,如输入输出信号的具体分配及其重要性。接着深入探讨了梯形图编程逻辑,涵盖外呼信号处理、轿厢运动控制以及楼层判断等关键环节。随后讲解了组态王的画面设计,包括动画效果的实现方法,如楼层按钮绑定、轿厢移动动画和门开合效果等。最后分享了一些调试经验和注意事项,如模拟困人场景、防抖逻辑、接线艺术等。 适合人群:从事自动化控制领域的工程师和技术人员,尤其是对PLC编程和组态软件有一定基础的人群。 使用场景及目标:适用于需要设计和实施小型电梯控制系统的工程项目。主要目标是帮助读者掌握PLC编程技巧、组态画面设计方法以及系统联调经验,从而提高项目的成功率。 其他说明:文中提供了详细的代码片段和调试技巧,有助于读者更好地理解和应用相关知识点。此外,还强调了安全性和可靠性方面的考量,如急停按钮的正确接入和硬件互锁设计等。

    CarSim与Simulink联合仿真:基于MPC模型预测控制实现智能超车换道

    内容概要:本文介绍了如何将CarSim的动力学模型与Simulink的智能算法相结合,利用模型预测控制(MPC)实现车辆的智能超车换道。主要内容包括MPC控制器的设计、路径规划算法、联合仿真的配置要点以及实际应用效果。文中提供了详细的代码片段和技术细节,如权重矩阵设置、路径跟踪目标函数、安全超车条件判断等。此外,还强调了仿真过程中需要注意的关键参数配置,如仿真步长、插值设置等,以确保系统的稳定性和准确性。 适合人群:从事自动驾驶研究的技术人员、汽车工程领域的研究人员、对联合仿真感兴趣的开发者。 使用场景及目标:适用于需要进行自动驾驶车辆行为模拟的研究机构和企业,旨在提高超车换道的安全性和效率,为自动驾驶技术研发提供理论支持和技术验证。 其他说明:随包提供的案例文件已调好所有参数,可以直接导入并运行,帮助用户快速上手。文中提到的具体参数和配置方法对于初学者非常友好,能够显著降低入门门槛。

    基于单片机的鱼缸监测设计(51+1602+AD0809+18B20+UART+JKx2)#0107

    包括:源程序工程文件、Proteus仿真工程文件、论文材料、配套技术手册等 1、采用51单片机作为主控; 2、采用AD0809(仿真0808)检测"PH、氨、亚硝酸盐、硝酸盐"模拟传感; 3、采用DS18B20检测温度; 4、采用1602液晶显示检测值; 5、检测值同时串口上传,调试助手监看; 6、亦可通过串口指令对加热器、制氧机进行控制;

    风电领域双馈永磁风电机组并网仿真及短路故障分析与MPPT控制

    内容概要:本文详细介绍了双馈永磁风电机组并网仿真模型及其短路故障分析方法。首先构建了一个9MW风电场模型,由6台1.5MW双馈风机构成,通过升压变压器连接到120kV电网。文中探讨了风速模块的设计,包括渐变风、阵风和随疾风的组合形式,并提供了相应的Python和MATLAB代码示例。接着讨论了双闭环控制策略,即功率外环和电流内环的具体实现细节,以及MPPT控制用于最大化风能捕获的方法。此外,还涉及了短路故障模块的建模,包括三相电压电流特性和离散模型与phasor模型的应用。最后,强调了永磁同步机并网模型的特点和注意事项。 适合人群:从事风电领域研究的技术人员、高校相关专业师生、对风电并网仿真感兴趣的工程技术人员。 使用场景及目标:适用于风电场并网仿真研究,帮助研究人员理解和优化风电机组在不同风速条件下的性能表现,特别是在短路故障情况下的应对措施。目标是提高风电系统的稳定性和可靠性。 其他说明:文中提供的代码片段和具体参数设置有助于读者快速上手并进行实验验证。同时提醒了一些常见的错误和需要注意的地方,如离散化步长的选择、初始位置对齐等。

    空手道训练测试系统BLE106版本

    适用于空手道训练和测试场景

    【音乐创作领域AI提示词】AI音乐提示词(deepseek,豆包,kimi,chatGPT,扣子空间,manus,AI训练师)

    内容概要:本文介绍了金牌音乐作词大师的角色设定、背景经历、偏好特点、创作目标、技能优势以及工作流程。金牌音乐作词大师凭借深厚的音乐文化底蕴和丰富的创作经验,能够为不同风格的音乐创作歌词,擅长将传统文化元素与现代流行文化相结合,创作出既富有情感又触动人心的歌词。在创作过程中,会严格遵守社会主义核心价值观,尊重用户需求,提供专业修改建议,确保歌词内容健康向上。; 适合人群:有歌词创作需求的音乐爱好者、歌手或音乐制作人。; 使用场景及目标:①为特定主题或情感创作歌词,如爱情、励志等;②融合传统与现代文化元素创作独特风格的歌词;③对已有歌词进行润色和优化。; 阅读建议:阅读时可以重点关注作词大师的创作偏好、技能优势以及工作流程,有助于更好地理解如何创作出高质量的歌词。同时,在提出创作需求时,尽量详细描述自己的情感背景和期望,以便获得更贴合心意的作品。

    linux之用户管理教程.md

    linux之用户管理教程.md

    基于单片机的搬运机器人设计(51+1602+L298+BZ+KEY6)#0096

    包括:源程序工程文件、Proteus仿真工程文件、配套技术手册等 1、采用51/52单片机作为主控芯片; 2、采用1602液晶显示设置及状态; 3、采用L298驱动两个电机,模拟机械臂动力、移动底盘动力; 3、首先按键配置-待搬运物块的高度和宽度(为0不能开始搬运); 4、按下启动键开始搬运,搬运流程如下: 机械臂先把物块抓取到机器车上, 机械臂减速 机器车带着物块前往目的地 机器车减速 机械臂把物块放下来 机械臂减速 机器车回到物块堆积处(此时机器车是空车) 机器车减速 蜂鸣器提醒 按下复位键,结束本次搬运

    基于下垂控制的三相逆变器电压电流双闭环仿真及MATLAB/Simulink/PLECS实现

    内容概要:本文详细介绍了基于下垂控制的三相逆变器电压电流双闭环控制的仿真方法及其在MATLAB/Simulink和PLECS中的具体实现。首先解释了下垂控制的基本原理,即有功调频和无功调压,并给出了相应的数学表达式。随后讨论了电压环和电流环的设计与参数整定,强调了两者带宽的差异以及PI控制器的参数选择。文中还提到了一些常见的调试技巧,如锁相环的响应速度、LC滤波器的谐振点处理、死区时间设置等。此外,作者分享了一些实用的经验,如避免过度滤波、合理设置采样周期和下垂系数等。最后,通过突加负载测试展示了系统的动态响应性能。 适合人群:从事电力电子、微电网研究的技术人员,尤其是有一定MATLAB/Simulink和PLECS使用经验的研发人员。 使用场景及目标:适用于希望深入了解三相逆变器下垂控制机制的研究人员和技术人员,旨在帮助他们掌握电压电流双闭环控制的具体实现方法,提高仿真的准确性和效率。 其他说明:本文不仅提供了详细的理论讲解,还结合了大量的实战经验和调试技巧,有助于读者更好地理解和应用相关技术。

    光伏并网逆变器全栈开发资料:硬件设计、控制算法及实战经验

    内容概要:本文详细介绍了光伏并网逆变器的全栈开发资料,涵盖了从硬件设计到控制算法的各个方面。首先,文章深入探讨了功率接口板的设计,包括IGBT缓冲电路、PCB布局以及EMI滤波器的具体参数和设计思路。接着,重点讲解了主控DSP板的核心控制算法,如MPPT算法的实现及其注意事项。此外,还详细描述了驱动扩展板的门极驱动电路设计,特别是光耦隔离和驱动电阻的选择。同时,文章提供了并联仿真的具体实现方法,展示了环流抑制策略的效果。最后,分享了许多宝贵的实战经验和调试技巧,如主变压器绕制、PWM输出滤波、电流探头使用等。 适合人群:从事电力电子、光伏系统设计的研发工程师和技术爱好者。 使用场景及目标:①帮助工程师理解和掌握光伏并网逆变器的硬件设计和控制算法;②提供详细的实战经验和调试技巧,提升产品的可靠性和性能;③适用于希望深入了解光伏并网逆变器全栈开发的技术人员。 其他说明:文中不仅提供了具体的电路设计和代码实现,还分享了许多宝贵的实际操作经验和常见问题的解决方案,有助于提高开发效率和产品质量。

    机器人轨迹规划中粒子群优化与3-5-3多项式结合的时间最优路径规划

    内容概要:本文详细介绍了粒子群优化(PSO)算法与3-5-3多项式相结合的方法,在机器人轨迹规划中的应用。首先解释了粒子群算法的基本原理及其在优化轨迹参数方面的作用,随后阐述了3-5-3多项式的数学模型,特别是如何利用不同阶次的多项式确保轨迹的平滑过渡并满足边界条件。文中还提供了具体的Python代码实现,展示了如何通过粒子群算法优化时间分配,使3-5-3多项式生成的轨迹达到时间最优。此外,作者分享了一些实践经验,如加入惩罚项以避免超速,以及使用随机扰动帮助粒子跳出局部最优。 适合人群:对机器人运动规划感兴趣的科研人员、工程师和技术爱好者,尤其是有一定编程基础并对优化算法有初步了解的人士。 使用场景及目标:适用于需要精确控制机器人运动的应用场合,如工业自动化生产线、无人机导航等。主要目标是在保证轨迹平滑的前提下,尽可能缩短运动时间,提高工作效率。 其他说明:文中不仅给出了理论讲解,还有详细的代码示例和调试技巧,便于读者理解和实践。同时强调了实际应用中需要注意的问题,如系统的建模精度和安全性考量。

    【KUKA 机器人资料】:kuka机器人压铸欧洲标准.pdf

    KUKA机器人相关资料

    光子晶体中BIC与OAM激发的模拟及三维Q值计算

    内容概要:本文详细探讨了光子晶体中的束缚态在连续谱中(BIC)及其与轨道角动量(OAM)激发的关系。首先介绍了光子晶体的基本概念和BIC的独特性质,随后展示了如何通过Python代码模拟二维光子晶体中的BIC,并解释了BIC在光学器件中的潜在应用。接着讨论了OAM激发与BIC之间的联系,特别是BIC如何增强OAM激发效率。文中还提供了使用有限差分时域(FDTD)方法计算OAM的具体步骤,并介绍了计算本征态和三维Q值的方法。此外,作者分享了一些实验中的有趣发现,如特定条件下BIC表现出OAM特征,以及不同参数设置对Q值的影响。 适合人群:对光子晶体、BIC和OAM感兴趣的科研人员和技术爱好者,尤其是从事微纳光子学研究的专业人士。 使用场景及目标:适用于希望通过代码模拟深入了解光子晶体中BIC和OAM激发机制的研究人员。目标是掌握BIC和OAM的基础理论,学会使用Python和其他工具进行模拟,并理解这些现象在实际应用中的潜力。 其他说明:文章不仅提供了详细的代码示例,还分享了许多实验心得和技巧,帮助读者避免常见错误,提高模拟精度。同时,强调了物理离散化方式对数值计算结果的重要影响。

    C#联合Halcon 17.12构建工业视觉项目的配置与应用

    内容概要:本文详细介绍了如何使用C#和Halcon 17.12构建一个功能全面的工业视觉项目。主要内容涵盖项目配置、Halcon脚本的选择与修改、相机调试、模板匹配、生产履历管理、历史图像保存以及与三菱FX5U PLC的以太网通讯。文中不仅提供了具体的代码示例,还讨论了实际项目中常见的挑战及其解决方案,如环境配置、相机控制、模板匹配参数调整、PLC通讯细节、生产数据管理和图像存储策略等。 适合人群:从事工业视觉领域的开发者和技术人员,尤其是那些希望深入了解C#与Halcon结合使用的专业人士。 使用场景及目标:适用于需要开发复杂视觉检测系统的工业应用场景,旨在提高检测精度、自动化程度和数据管理效率。具体目标包括但不限于:实现高效的视觉处理流程、确保相机与PLC的无缝协作、优化模板匹配算法、有效管理生产和检测数据。 其他说明:文中强调了框架整合的重要性,并提供了一些实用的技术提示,如避免不同版本之间的兼容性问题、处理实时图像流的最佳实践、确保线程安全的操作等。此外,还提到了一些常见错误及其规避方法,帮助开发者少走弯路。

    基于Matlab的9节点配电网中分布式电源接入对节点电压影响的研究

    内容概要:本文探讨了分布式电源(DG)接入对9节点配电网节点电压的影响。首先介绍了9节点配电网模型的搭建方法,包括定义节点和线路参数。然后,通过在特定节点接入分布式电源,利用Matlab进行潮流计算,模拟DG对接入点及其周围节点电压的影响。最后,通过绘制电压波形图,直观展示了不同DG容量和接入位置对配电网电压分布的具体影响。此外,还讨论了电压越限问题以及不同线路参数对电压波动的影响。 适合人群:电力系统研究人员、电气工程学生、从事智能电网和分布式能源研究的专业人士。 使用场景及目标:适用于研究分布式电源接入对配电网电压稳定性的影响,帮助优化分布式电源的规划和配置,确保电网安全稳定运行。 其他说明:文中提供的Matlab代码和图表有助于理解和验证理论分析,同时也为后续深入研究提供了有价值的参考资料。

Global site tag (gtag.js) - Google Analytics