Package backtype.storm.messaging

Examples of backtype.storm.messaging.TaskMessage.message()


      LOG.warn("Received invalid message directed at port " + task
          + ". Dropping...");
      return;
    }

    queue.publish(message.message());

  }

}
View Full Code Here


    list.add(message);

    client.send(message);

    TaskMessage recv = server.recv(0);
    Assert.assertEquals(req_msg, new String(recv.message()));

    System.out.println("!!!!!!!!!!!!!!!!!!Test one time!!!!!!!!!!!!!!!!!");

    server.close();
    client.close();
View Full Code Here

    LOG.info("Client send data");
    client.send(message);

    TaskMessage recv = server.recv(0);
    Assert.assertEquals(req_msg, new String(recv.message()));

    client.close();
    server.close();
    context.term();
    System.out.println("!!!!!!!!!!End larget message test!!!!!!!!");
View Full Code Here

    LOG.info("Client send data");
    client.send(message);
    Thread.sleep(1000);

    TaskMessage recv = server.recv(0);
    Assert.assertEquals(req_msg, new String(recv.message()));

    server.close();
    client.close();
    context.term();
    System.out.println("!!!!!!!!!!End delay message test!!!!!!!!");
View Full Code Here

    DisruptorQueue recvQueue = new DisruptorQueue(
        "NettyUnitTest", ProducerType.SINGLE, 1024, waitStrategy);
    server.registerQueue(recvQueue);

    TaskMessage recv = server.recv(0);
    Assert.assertEquals(req_msg, new String(recv.message()));

    lock.lock();
    clientClose.signal();
    server.close();
    contextClose.await();
View Full Code Here

    for (int i = 1; i < Short.MAX_VALUE; i++) {
      TaskMessage message = server.recv(0);

      Assert.assertEquals(String.valueOf(i + base),
          new String(message.message()));

      if (i % 1000 == 0) {
        System.out.println("Receive " + message.task());
      }
    }
View Full Code Here

    for (int i = 1; i < Short.MAX_VALUE; i++) {
      TaskMessage message = server.recv(0);
      JStormUtils.sleepMs(1);

      Assert.assertEquals(String.valueOf(i + base),
          new String(message.message()));

      if (i % 1000 == 0) {
        System.out.println("Receive " + message.task());
      }
    }
View Full Code Here

    for (int i = 1; i < base; i++) {
      TaskMessage message = server.recv(0);
      JStormUtils.sleepMs(100);

      Assert.assertEquals(req_msg, new String(message.message()));
      System.out.println("receive " + message.task());

    }

    System.out.println("Finish Receive ");
View Full Code Here

        "NettyUnitTest", ProducerType.SINGLE, 1024, waitStrategy);
    server.registerQueue(recvQueue);

    TaskMessage recv = server.recv(0);
    System.out.println("Sever receive first");
    Assert.assertEquals(req_msg, new String(recv.message()));

    Thread.sleep(1000);

    TaskMessage recv2 = server.recv(0);
    System.out.println("Sever receive second");
View Full Code Here

    Thread.sleep(1000);

    TaskMessage recv2 = server.recv(0);
    System.out.println("Sever receive second");
    Assert.assertEquals(req_msg, new String(recv2.message()));

    lock.lock();
        clientClose.signal();
        server.close();
        contextClose.await();
View Full Code Here

TOP
Copyright © 2018 www.massapi.com. All rights reserved.
All source code are property of their respective owners. Java is a trademark of Sun Microsystems, Inc and owned by ORACLE Inc. Contact coftware#gmail.com.