Examples of registerQueue()


Examples of backtype.storm.messaging.IConnection.registerQueue()

    IContext context = workerData.getContext();
    String topologyId = workerData.getTopologyId();

    IConnection recvConnection = context.bind(topologyId,
        workerData.getPort());
    recvConnection.registerQueue(recvQueue);

    RunnableCallback recvDispather = new VirtualPortDispatch(workerData,
        recvConnection, recvQueue);

    AsyncLoopThread vthread = new AsyncLoopThread(recvDispather, false,
View Full Code Here

Examples of backtype.storm.messaging.IConnection.registerQueue()

    WaitStrategy waitStrategy = (WaitStrategy) Utils
        .newInstance((String) storm_conf
            .get(Config.TOPOLOGY_DISRUPTOR_WAIT_STRATEGY));
    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();
View Full Code Here

Examples of backtype.storm.messaging.IConnection.registerQueue()

    WaitStrategy waitStrategy = (WaitStrategy) Utils
        .newInstance((String) storm_conf
            .get(Config.TOPOLOGY_DISRUPTOR_WAIT_STRATEGY));
    DisruptorQueue recvQueue = new DisruptorQueue(
        "NettyUnitTest", ProducerType.SINGLE, 1024, waitStrategy);
    server.registerQueue(recvQueue);


    new Thread(new Runnable() {
       
        public void send() {
View Full Code Here

Examples of backtype.storm.messaging.IConnection.registerQueue()

    WaitStrategy waitStrategy = (WaitStrategy) Utils
        .newInstance((String) storm_conf
            .get(Config.TOPOLOGY_DISRUPTOR_WAIT_STRATEGY));
    DisruptorQueue recvQueue = new DisruptorQueue(
        "NettyJunitTest", ProducerType.SINGLE, 4, waitStrategy);
    server.registerQueue(recvQueue);


    new Thread(new Runnable() {

      @Override
View Full Code Here

Examples of backtype.storm.messaging.IConnection.registerQueue()

    WaitStrategy waitStrategy = (WaitStrategy) Utils
        .newInstance((String) storm_conf
            .get(Config.TOPOLOGY_DISRUPTOR_WAIT_STRATEGY));
    DisruptorQueue recvQueue = new DisruptorQueue(
        "NettyJunitTest", ProducerType.SINGLE, 4, waitStrategy);
    server.registerQueue(recvQueue);

    new Thread(new Runnable() {

      @Override
      public void run() {
View Full Code Here

Examples of backtype.storm.messaging.IConnection.registerQueue()

    WaitStrategy waitStrategy = (WaitStrategy) Utils
        .newInstance((String) storm_conf
            .get(Config.TOPOLOGY_DISRUPTOR_WAIT_STRATEGY));
    DisruptorQueue recvQueue = new DisruptorQueue(
        "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()));
View Full Code Here

Examples of backtype.storm.messaging.IConnection.registerQueue()

    System.out.println("!!shutdow server and sleep 30s, please wait!!");
    Thread.sleep(30000);

    IConnection server2 = context.bind(null, port);
    server2.registerQueue(recvQueue);
    System.out.println("!!!!!!!!!!!!!!!!!!!! restart server !!!!!!!!!!!");

    TaskMessage recv2 = server2.recv(0);
    Assert.assertEquals(req_msg, new String(recv2.message()));
View Full Code Here

Examples of org.ajax4jsf.component.QueueRegistry.registerQueue()

      if (queueRegistry.isShouldCreateDefaultGlobalQueue()) {
        String encodedGlobalQueueName = context.getExternalContext().encodeNamespace(
            UIQueue.GLOBAL_QUEUE_NAME);
       
        if (!queueRegistry.containsQueue(encodedGlobalQueueName)) {
          queueRegistry.registerQueue(context, encodedGlobalQueueName, null);
        }
      }

      if (queueRegistry.hasQueuesToEncode()) {
        super.encode(resource, context, queueRegistry.getRegisteredQueues(context), attributes);
View Full Code Here

Examples of org.apache.qpid.server.exchange.Exchange.registerQueue()

                if (defaultExchange == null)
                {
                    throw new ConfigurationException("Attempt to bind queue to unknown exchange:" + binding);
                }

                defaultExchange.registerQueue(queue.getName(), queue, null);

                if (binding.getRoutingKey() == null || binding.getRoutingKey().equals(""))
                {
                    throw new ConfigurationException("Unknown binding not specified on url:" + binding);
                }
View Full Code Here

Examples of org.apache.qpid.server.exchange.Exchange.registerQueue()

        final Exchange exch = exchangeRegistry.getExchange(body.exchange);
        if (exch == null)
        {
            throw body.getChannelException(AMQConstant.NOT_FOUND.getCode(), "Exchange " + body.exchange + " does not exist.");
        }
        exch.registerQueue(body.routingKey, queue, body.arguments);
        queue.bind(body.routingKey, exch);
        if (_log.isInfoEnabled())
        {
            _log.info("Binding queue " + queue + " to exchange " + exch + " with routing key " + body.routingKey);
        }
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.