Class QueueingConsumer
- All Implemented Interfaces:
Consumer
Consumer with
straightforward blocking semantics.
The general pattern for using QueueingConsumer is as follows:
// Create connection and channel.ConnectionFactoryfactory = new ConnectionFactory(); Connection conn = factory.newConnection();Channelch1 = conn.createChannel(); // Declare a queue and bind it to an exchange. String queueName = ch1.queueDeclare().getQueue(); ch1.queueBind(queueName, exchangeName, queueName); // Create the QueueingConsumer and have it consume from the queue QueueingConsumer consumer = newQueueingConsumer(ch1); ch1.basicConsume(queueName, false, consumer); // Process deliveries while (/* some condition * /) {QueueingConsumer.Deliverydelivery = consumer.nextDelivery(); // process delivery ch1.basicAck(delivery.getEnvelope().getDeliveryTag(), false); }
For a more complete example, see LogTail in the test/src/com/rabbitmq/examples
directory of the source distribution.
Historical Perspective
QueueingConsumer was introduced to allow
applications to overcome a limitation in the way Connection
managed threads and consumer dispatching. When QueueingConsumer
was introduced, callbacks to Consumers were made on the
Connection's thread. This had two main drawbacks. Firstly, the
Consumer could stall the processing of all
Channels on the Connection. Secondly, if a
Consumer made a recursive synchronous call into its
Channel the client would deadlock.
QueueingConsumer provided client code with an easy way to
obviate this problem by queueing incoming messages and processing them on
a separate, application-managed thread.
The threading behaviour of Connection and Channel
has been changed so that each Channel uses a distinct thread
for dispatching to Consumers. This prevents
Consumers on one Channel holding up
Consumers on another and it also prevents recursive calls from
deadlocking the client.
As such, it is now safe to implement Consumer directly or
to extend DefaultConsumer and QueueingConsumer
is a lot less relevant.
-
Nested Class Summary
Nested ClassesModifier and TypeClassDescriptionstatic classEncapsulates an arbitrary message - simple "bean" holder structure. -
Constructor Summary
ConstructorsConstructorDescription -
Method Summary
Modifier and TypeMethodDescriptionvoidhandleCancel(String consumerTag) No-op implementation ofConsumer.handleCancel(String)voidhandleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) voidhandleShutdownSignal(String consumerTag, ShutdownSignalException sig) No-op implementation ofConsumer.handleShutdownSignal(java.lang.String, com.rabbitmq.client.ShutdownSignalException).Main application-side API: wait for the next message delivery and return it.nextDelivery(long timeout) Main application-side API: wait for the next message delivery and return it.Methods inherited from class com.rabbitmq.client.DefaultConsumer
getChannel, getConsumerTag, handleCancelOk, handleConsumeOk, handleRecoverOk
-
Constructor Details
-
QueueingConsumer
-
QueueingConsumer
-
-
Method Details
-
handleShutdownSignal
Description copied from class:DefaultConsumerNo-op implementation ofConsumer.handleShutdownSignal(java.lang.String, com.rabbitmq.client.ShutdownSignalException).- Specified by:
handleShutdownSignalin interfaceConsumer- Overrides:
handleShutdownSignalin classDefaultConsumer- Parameters:
consumerTag- the consumer tag associated with the consumersig- aShutdownSignalExceptionindicating the reason for the shut down
-
handleCancel
Description copied from class:DefaultConsumerNo-op implementation ofConsumer.handleCancel(String)- Specified by:
handleCancelin interfaceConsumer- Overrides:
handleCancelin classDefaultConsumer- Parameters:
consumerTag- the defined consumer tag (client- or server-generated)- Throws:
IOException
-
handleDelivery
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException Description copied from class:DefaultConsumerNo-op implementation ofConsumer.handleDelivery(java.lang.String, com.rabbitmq.client.Envelope, com.rabbitmq.client.AMQP.BasicProperties, byte[]).- Specified by:
handleDeliveryin interfaceConsumer- Overrides:
handleDeliveryin classDefaultConsumer- Parameters:
consumerTag- the consumer tag associated with the consumerenvelope- packaging data for the messageproperties- content header data for the messagebody- the message body (opaque, client-specific byte array)- Throws:
IOException- if the consumer encounters an I/O error while processing the message- See Also:
-
nextDelivery
public QueueingConsumer.Delivery nextDelivery() throws InterruptedException, ShutdownSignalException, ConsumerCancelledExceptionMain application-side API: wait for the next message delivery and return it.- Returns:
- the next message
- Throws:
InterruptedException- if an interrupt is received while waitingShutdownSignalException- if the connection is shut down while waitingConsumerCancelledException- if this consumer is cancelled while waiting
-
nextDelivery
public QueueingConsumer.Delivery nextDelivery(long timeout) throws InterruptedException, ShutdownSignalException, ConsumerCancelledException Main application-side API: wait for the next message delivery and return it.- Parameters:
timeout- timeout in millisecond- Returns:
- the next message or null if timed out
- Throws:
InterruptedException- if an interrupt is received while waitingShutdownSignalException- if the connection is shut down while waitingConsumerCancelledException- if this consumer is cancelled while waiting
-