diff options
author | Arnaud Simon <arnaudsimon@apache.org> | 2009-01-16 12:58:42 +0000 |
---|---|---|
committer | Arnaud Simon <arnaudsimon@apache.org> | 2009-01-16 12:58:42 +0000 |
commit | 64eadcd9e4324c2b0f51acb920f09535d62a89a5 (patch) | |
tree | f65960e7dc5bf30c09c6d0ae1c404ffda09b31d3 /java/client/example/src | |
parent | c9487be5665905d8a146d038695acdd38fab6081 (diff) | |
download | qpid-python-64eadcd9e4324c2b0f51acb920f09535d62a89a5.tar.gz |
Qpid-1587: Added LVQ samples
git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk/qpid@734994 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'java/client/example/src')
3 files changed, 246 insertions, 0 deletions
diff --git a/java/client/example/src/main/java/org/apache/qpid/example/amqpexample/lvq/DeclareLVQueue.java b/java/client/example/src/main/java/org/apache/qpid/example/amqpexample/lvq/DeclareLVQueue.java new file mode 100644 index 0000000000..86a0f362ad --- /dev/null +++ b/java/client/example/src/main/java/org/apache/qpid/example/amqpexample/lvq/DeclareLVQueue.java @@ -0,0 +1,64 @@ +package org.apache.qpid.example.amqpexample.lvq; +/* + * + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + * + */ + + +import org.apache.qpid.transport.Connection; +import org.apache.qpid.transport.Session; + +import java.util.Map; +import java.util.HashMap; + +/** + * This creates a queue a LVQueue with key test and binds it to the + * amq.direct exchange + * + */ +public class DeclareLVQueue +{ + + public static void main(String[] args) + { + // Create connection + Connection con = new Connection(); + con.connect("localhost", 5672, "test", "guest", "guest",false); + + // Create session + Session session = con.createSession(0); + + // declare and bind queue + Map<String, Object> arguments = new HashMap<String, Object>(); + // We use a lvq + arguments.put("qpid.last_value_queue", true); + // We want this queue to use the key test + arguments.put("qpid.LVQ_key", "test"); + session.queueDeclare("message_queue", null, arguments); + session.exchangeBind("message_queue", "amq.direct", "routing_key", null); + + // confirm completion + session.sync(); + + //cleanup + session.close(); + con.close(); + } +} + diff --git a/java/client/example/src/main/java/org/apache/qpid/example/amqpexample/lvq/Listener.java b/java/client/example/src/main/java/org/apache/qpid/example/amqpexample/lvq/Listener.java new file mode 100644 index 0000000000..b930062813 --- /dev/null +++ b/java/client/example/src/main/java/org/apache/qpid/example/amqpexample/lvq/Listener.java @@ -0,0 +1,113 @@ +package org.apache.qpid.example.amqpexample.lvq; +/* + * + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + * + */ + + +import org.apache.qpid.transport.Connection; +import org.apache.qpid.transport.MessageAcceptMode; +import org.apache.qpid.transport.MessageAcquireMode; +import org.apache.qpid.transport.MessageCreditUnit; +import org.apache.qpid.transport.MessageTransfer; +import org.apache.qpid.transport.Session; +import org.apache.qpid.transport.SessionException; +import org.apache.qpid.transport.SessionListener; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; + +/** + * This listens to messages on a queue and terminates + * when it sees the final message + * + */ +public class Listener implements SessionListener +{ + private static CountDownLatch _countDownLatch = new CountDownLatch(1); + + public void opened(Session ssn) {} + + public void message(Session ssn, MessageTransfer xfr) + { + String body = xfr.getBodyString(); + System.out.println("Message: " + body); + if ( body.equals("That's all, folks!")) + { + System.out.println("Received final message"); + _countDownLatch.countDown(); + } + } + + public void exception(Session ssn, SessionException exc) + { + exc.printStackTrace(); + } + + public void closed(Session ssn) {} + + /** + * This sends 10 messages to the + * amq.direct exchange using the + * routing key as "routing_key" + * + * @param args + */ + public static void main(String[] args) throws InterruptedException + { + // Create connection + Connection con = new Connection(); + con.connect("localhost", 5672, "test", "guest", "guest",false); + + // Create session + Session session = con.createSession(0); + + // Create an instance of the listener + Listener listener = new Listener(); + session.setSessionListener(listener); + + // create a subscription + session.messageSubscribe("message_queue", + "listener_destination", + MessageAcceptMode.NONE, + MessageAcquireMode.PRE_ACQUIRED, + null, 0, null); + + + // issue credits + // XXX + session.messageFlow("listener_destination", MessageCreditUnit.BYTE, Session.UNLIMITED_CREDIT); + session.messageFlow("listener_destination", MessageCreditUnit.MESSAGE, 11); + + // confirm completion + session.sync(); + + // wait to receive all the messages + System.out.println("Waiting 100 seconds for messages from listener_destination"); + + _countDownLatch.await(30, TimeUnit.SECONDS); + System.out.println("Shutting down listener for listener_destination"); + session.messageCancel("listener_destination"); + + //cleanup + session.close(); + con.close(); + } + +} diff --git a/java/client/example/src/main/java/org/apache/qpid/example/amqpexample/lvq/Producer.java b/java/client/example/src/main/java/org/apache/qpid/example/amqpexample/lvq/Producer.java new file mode 100644 index 0000000000..482e6a6b11 --- /dev/null +++ b/java/client/example/src/main/java/org/apache/qpid/example/amqpexample/lvq/Producer.java @@ -0,0 +1,69 @@ +package org.apache.qpid.example.amqpexample.lvq; +/* + * + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + * + */ + +import java.util.Map; +import java.util.HashMap; +import org.apache.qpid.transport.*; + +public class Producer +{ + + public static void main(String[] args) + { + // Create connection + Connection con = new Connection(); + con.connect("localhost", 5672, "test", "guest", "guest",false); + + // Create session + Session session = con.createSession(0); + DeliveryProperties deliveryProps = new DeliveryProperties(); + deliveryProps.setRoutingKey("routing_key"); + + // set message headers + MessageProperties messageProperties = new MessageProperties(); + Map<String, Object> messageHeaders = new HashMap<String, Object>(); + // set the message key + messageHeaders.put("qpid.LVQ_key", "test"); + messageProperties.setApplicationHeaders(messageHeaders); + + Header header = new Header(deliveryProps, messageProperties); + + for (int i=0; i<10; i++) + { + session.messageTransfer("amq.direct", MessageAcceptMode.EXPLICIT,MessageAcquireMode.PRE_ACQUIRED, + header, + "Message " + i); + } + + session.messageTransfer("amq.direct", MessageAcceptMode.EXPLICIT,MessageAcquireMode.PRE_ACQUIRED, + header, + "That's all, folks!"); + + // confirm completion + session.sync(); + + //cleanup + session.close(); + con.close(); + } + +} |