【问题标题】:hornetq scheduled messages are not delivered on timehornetq 预定的消息没有按时传递
【发布时间】:2014-02-15 03:42:21
【问题描述】:

我正在尝试使用 hornetq core-api 预定消息 api 以便在 30 秒后发送消息 队列不持久并在 hornetq 配置文件中定义

val message:String ={...} //some string
val clientMessage:ClientMessage =session.createMessage(false)
clientMessage.getBodyBuffer.writeString(message)
clientMessage.putLongProperty(Message.HDR_SCHEDULED_DELIVERY_TIME, System.currentTimeMillis() + 30000) 
//expecting to deliver the message after 30 seconds
    producer.send(MyQueue,clientMessage)

但是,当我查看日志时,似乎消息是在同一秒发送并到达的。我应该定义其他东西吗?我错过了什么吗?

添加代码:

class HornetQMessageTest extends FunSuite with ShouldMatchers {
test("scheduleMessage") {
    def createServerLocator: ServerLocator = {
      var map = new java.util.HashMap[String, Object]
      map.put("host",  "127.0.0.1")
      map.put("port", "5445")
      val transConf = new TransportConfiguration(classOf[NettyConnectorFactory].getName,map)
      val locator = HornetQClient.createServerLocatorWithoutHA(transConf)
      locator.setConfirmationWindowSize(1024^2)//confirmationWindowSize
      locator.setBlockOnDurableSend(false)
      locator.setBlockOnNonDurableSend(false)
      locator.setClientFailureCheckPeriod(5000) //keepAlivePing
      locator.setConnectionTTL(10000) //connection TTL
      locator
    }
    val serverLocator: ServerLocator = createServerLocator
    val sessionFactory: ClientSessionFactory = serverLocator.createSessionFactory()
    val receiverSession = sessionFactory.createSession(true, true, 0)
    val senderSession =sessionFactory.createSession(true, true, 0)
    val queue = "atestq"
    def closeHorentQClient() {
      receiverSession.close()
      senderSession.close()
      sessionFactory.close()
      serverLocator.close()
    }
    senderSession should not be null
    receiverSession should not be null
    val query: QueueQuery = senderSession.queueQuery(new SimpleString(queue))
    if (query == null || !query.isExists) senderSession.createQueue(queue,queue,false)
    val producer: ClientProducer = senderSession.createProducer(queue)
    val message = senderSession.createMessage(false)
      message.getBodyBuffer().writeString("This is my string........")

    val deliverytime = System.currentTimeMillis() + 30000
    message.putLongProperty(Message.HDR_SCHEDULED_DELIVERY_TIME, deliverytime)
    println("Message Sent "+deliverytime)
    senderSession.start()
    producer.send(message)
    receiverSession.start()
    val consumer = receiverSession.createConsumer(queue)
    val message2 = consumer.receive(50000)
    val messageBody = message2.getBodyBuffer().readString()
    println("received message: "+messageBody +" after" + ((System.currentTimeMillis()-deliverytime)/1000)+" seconds ")
    message2 should not be  null
    System.currentTimeMillis() should be >= deliverytime
    assert(messageBody == "This is my string........")

    message2.acknowledge()

    // Make sure no more messages
    closeHorentQClient()
}  }

测试结果:

Message Sent 1392198959686
received message: This is my string........ after-29 seconds 

1392198929758 was not greater than or equal to 1392198959686

【问题讨论】:

    标签: hornetq


    【解决方案1】:

    我们正在使用 Scheduled Executor,并且在某些操作系统中... scheduler.schedule(在 30 秒内)调用我们的 Runnable 的时间比您预期的要短。实际上,我已经在 Windows 上看到过很多这样的问题。

    确保您了解计划时间是基于服务器的时间(而不是客户端的时间)。将您的客户端时间与您的服务器同步。

    我们最近对 2.4.0 做了很多改进,我认为您不会再次遇到这个问题,因为我们现在验证时间并且不再信任 ScheduledExecutor。

    如果您更新您的操作系统,您将不会看到此问题。这通常是实时内核问题,并且没有尊重等待。

    或者,如果您可以使用不再信任此行为的 2.4.0。

    我已经在 MAC 上使用 2.2.eap5 尝试过这段代码,但没有在 2.4.0 上打补丁,它可以工作。这几乎是您的测试转换为 java。

       public void testSomething() throws Exception
       {
          // then we create a client as normal
          ClientSessionFactory sessionFactory = createSessionFactory(locator);
          ClientSession receiverSession = sessionFactory.createSession(true, true, 0);
          ClientSession senderSession = sessionFactory.createSession(true, true, 0);
          String queue = "atestq";
    
          ClientSession.QueueQuery query = senderSession.queueQuery(new SimpleString(queue));
          if (query == null || !query.isExists()) senderSession.createQueue(queue, queue, false);
    
          ClientProducer producer = senderSession.createProducer(queue);
          ClientMessage message = senderSession.createMessage(false);
          message.getBodyBuffer().writeString("This is my string........");
    
          long original = System.currentTimeMillis();
          long deliverytime = System.currentTimeMillis() + 30000;
          message.putLongProperty(Message.HDR_SCHEDULED_DELIVERY_TIME, deliverytime);
          System.out.println("Message Sent " + deliverytime);
          senderSession.start();
          producer.send(message);
          receiverSession.start();
          ClientConsumer consumer = receiverSession.createConsumer(queue);
          ClientMessage message2 = consumer.receive(50000);
          String messageBody = message2.getBodyBuffer().readString();
          System.out.println("received message: " + messageBody + " after" + ((System.currentTimeMillis() - original) / 1000) + " seconds ");
          assert (messageBody.equals("This is my string........"));
    
          message2.acknowledge();
       }
    

    我刚刚在 HornetQ 中提出了一个功能请求。

    【讨论】:

    • 谢谢,我已经添加了我正在使用的测试代码,仍然可以同时收到消息,没有延迟,正如您在结果中看到的那样
    • 我们使用的是 2.3.0 版本
    • 我想我现在已经回答了...您在帖子编辑中的详细信息帮助我了解了您所看到的内容..
    • 将版本升级到 2.4.0.final 但结果还是一样,你可以试试我写的测试吗?也许测试代码有问题?
    • 您是在不同的服务器上运行它吗?服务器上的时间和客户端上的时间一样吗?调度程序基于服务器的时间而不是客户端。
    猜你喜欢
    • 2014-02-02
    • 2013-03-20
    • 1970-01-01
    • 2014-04-22
    • 2012-05-08
    • 2021-10-30
    • 2019-09-25
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多