【问题标题】:Concurrent JMS messages cause NULL并发 JMS 消息导致 NULL
【发布时间】:2014-08-18 18:24:09
【问题描述】:

我已经测试了串行和并发发送 JMS 消息(5 个线程同时从生产者发送 jms 消息)。

当我同时发送 100 条消息时,接收端的消息负载很少是 NULL。串行发送时没有问题。

我是否需要在消费者端设置会话池或使用 MDB 来同时处理消息? JMS 的设置很好,因为我们正在接收消息。我在这里有什么遗漏吗?

pproject 设置的简短描述:

  1. Publisher 是无状态会话 bean
  2. Weblogic 8.1 jms 服务器连接工厂和目标通过检索 JNDI
  3. Consumer 是订阅此服务器 JMS 队列的 java 类 并执行任务。 (这不是 MDB 或 Threaded 类,监听队列 异步)

已编辑

JmsConsumer

package net;

import java.io.BufferedWriter;
import java.io.FileWriter;
import java.io.IOException;
import java.io.PrintWriter;
import java.util.Date;
import java.util.Hashtable;

import javax.jms.ExceptionListener;
import javax.jms.JMSException;
import javax.jms.Message;
import javax.jms.MessageListener;
import javax.jms.Queue;
import javax.jms.QueueConnection;
import javax.jms.QueueConnectionFactory;
import javax.jms.QueueReceiver;
import javax.jms.QueueSession;
import javax.jms.Session;
import javax.jms.TextMessage;
import javax.naming.Context;
import javax.naming.InitialContext;

public class ReadJMS implements MessageListener, ExceptionListener {

    public final static String JNDI_FACTORY = "weblogic.jndi.WLInitialContextFactory";
    public final static String PROVIDER_URL = "t3://address:7003";

    public final static String JMS_FACTORY = "MSS.QueueConnectionFactory";
    public final static String QUEUE = "jms.queue";

    @SuppressWarnings("null")
    public void receiveMessage() throws Exception {
        // System.out.println("receiveMessage()..");
        Hashtable env = new Hashtable();
        env.put(Context.INITIAL_CONTEXT_FACTORY, JNDI_FACTORY);
        env.put(Context.PROVIDER_URL, PROVIDER_URL);

        // Define queue
        QueueReceiver qreceiver = null;
        QueueSession qsession = null;
        QueueConnection qcon = null;
        ReadJMS async = new ReadJMS();
        try {
            InitialContext ctx = new InitialContext(env);

            QueueConnectionFactory qconFactory = (QueueConnectionFactory) ctx
                    .lookup(JMS_FACTORY);
            qcon = qconFactory.createQueueConnection();
            qcon.setExceptionListener(async);
            qsession = qcon.createQueueSession(false, Session.AUTO_ACKNOWLEDGE);
            Queue queue = (Queue) ctx.lookup(QUEUE);
            qreceiver = qsession.createReceiver(queue);
            qreceiver.setMessageListener(async);

            qcon.start();
            System.out.println("readingMessage()..");
            // TextMessage msg = (TextMessage) qreceiver.receive();
            // System.out.println("Message read from " + QUEUE + " : "
            // + msg.getText());
            // msg.acknowledge();
        } catch (Exception ex) {
            ex.printStackTrace();
        }
        // } finally {
        // if (qreceiver != null)
        // qreceiver.close();
        // if (qsession != null)
        // qsession.close();
        // if (qcon != null)
        // qcon.close();
        // }
    }

    public static void main(String[] args) throws Exception {
        ReadJMS test = new ReadJMS();
        System.out.println("init");
        test.receiveMessage();
        while (true) {
            Thread.sleep(10000);
        }
    }

    public void onException(JMSException arg0) {
        System.err.println("Exception: " + arg0.getLocalizedMessage());
    }

    public synchronized void onMessage(Message arg0) {
        try {
            if(((TextMessage)arg0).getText() == null || ((TextMessage)arg0).getText().trim().length()==0){
                System.out.println(" " + QUEUE + " : "
                        + ((TextMessage) arg0).getText());
            }
            System.out.print(".");
            PrintWriter out = new PrintWriter(new BufferedWriter(
                    new FileWriter("Output.txt", true)));
            Date now = new Date();
            out.println("message: "+now.toString()+ " - "+((TextMessage)arg0).getText()+"");
            out.close();

        } catch (JMSException e) {
            e.printStackTrace();
        } catch (IOException e) {
            e.printStackTrace();
        }
    }
}

WriteJms

package net;

import java.io.BufferedReader;
import java.io.InputStreamReader;
import java.util.Hashtable;

import javax.jms.Queue;
import javax.jms.QueueConnection;
import javax.jms.QueueConnectionFactory;
import javax.jms.QueueSender;
import javax.jms.QueueSession;
import javax.jms.Session;
import javax.jms.TextMessage;
import javax.naming.Context;
import javax.naming.InitialContext;

public class WriteJMS {
    public final static String JNDI_FACTORY = "weblogic.jndi.WLInitialContextFactory";
    public final static String PROVIDER_URL = "t3://url:7003";
    public final static String JMS_FACTORY = "MSS.QueueConnectionFactory";
    public final static String QUEUE = "jms.queue";

    @SuppressWarnings("unchecked")
    public void sendMessage() throws Exception {

        @SuppressWarnings("rawtypes")
        Hashtable env = new Hashtable();
        env.put(Context.INITIAL_CONTEXT_FACTORY, JNDI_FACTORY);
        env.put(Context.PROVIDER_URL, PROVIDER_URL);

        // Define queue
        QueueSender qsender = null;
        QueueSession qsession = null;
        QueueConnection qcon = null;
        try {
            InitialContext ctx = new InitialContext(env);

            QueueConnectionFactory qconFactory = (QueueConnectionFactory) ctx
                    .lookup(JMS_FACTORY);
            qcon = qconFactory.createQueueConnection();

            qsession = qcon.createQueueSession(false, Session.AUTO_ACKNOWLEDGE);
            Queue queue = (Queue) ctx.lookup(QUEUE);

            TextMessage msg = qsession.createTextMessage();
            msg.setText("<eventMessage><eventId>123</eventId><eventName>123</eventName><documentNumber>123</documentNumber><customerId>123</customerId><actDDTaskDate>123</actDDTaskDate><taskStatusErrorMessage>123</taskStatusErrorMessage></eventMessage>");
            qsender = qsession.createSender(queue);
            qsender.send(msg);
            System.out.println("Message [" + msg.getText()
                    + "] sent to Queue: " + QUEUE);
        } catch (Exception ex) {
            ex.printStackTrace();
        } finally {
            if (qsender != null)
                qsender.close();
            if (qsession != null)
                qsession.close();
            if (qcon != null)
                qcon.close();
        }
    }


}

【问题讨论】:

  • 你说发送为空是什么意思?收到的某些消息的负载是否为空?
  • @AniketThakur 是的,有效载荷为空。
  • @AniketThakur 接收端的消息为空。发送时,我正在验证日志,它似乎发送没有任何问题。发送消息之前和之后
  • @jtahlborn 更新了代码
  • 嗯,发件人很有趣。

标签: java concurrency jms weblogic


【解决方案1】:

众所周知,禁止跨多个线程重用会话。但你不这样做。您为每条消息重新创建所有内容(连接会话和生产者)。这是低效的,但不是不正确的,不应该导致这些错误。你给我们的代码在我看来不错。

我有点惊讶,发送方没有出现异常。您能否提供有关 JMS 实现的更多详细信息?也许消息代理的日志中有更多信息?

您是否计算过消息数量,收到的数量是否等于发送的数量?其他人是否可以将消息发送到同一个队列?

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2015-02-11
    • 1970-01-01
    • 2021-10-18
    • 1970-01-01
    • 2015-05-19
    • 2019-05-19
    • 2015-08-16
    • 2016-01-11
    相关资源
    最近更新 更多