ServiceBusSessionReceiverClient Class

  • java.lang.Object
    • com.azure.messaging.servicebus.ServiceBusSessionReceiverClient

Implements

public final class ServiceBusSessionReceiverClient
implements AutoCloseable

This synchronous session receiver client is used to acquire session locks from a queue or topic and create ServiceBusReceiverClient instances that are tied to the locked sessions. Sessions can be used as a first in first out (FIFO) processing of messages. Queues and topics/subscriptions support Service Bus sessions, however, it must be enabled at the time of entity creation.

The examples shown in this document use a credential object named DefaultAzureCredential for authentication, which is appropriate for most scenarios, including local development and production environments. Additionally, we recommend using managed identity for authentication in production environments. You can find more information on different ways of authenticating and their corresponding credential types in the Azure Identity documentation".

Sample: Receive messages from a specific session

Use acceptSession(String sessionId) to acquire the lock of a session if you know the session id. PEEK_LOCK is strongly recommended so users have control over message settlement.

TokenCredential credential = new DefaultAzureCredentialBuilder().build();

 // 'fullyQualifiedNamespace' will look similar to "{your-namespace}.servicebus.windows.net"
 // 'disableAutoComplete' indicates that users will explicitly settle their message.
 ServiceBusSessionReceiverClient sessionReceiver = new ServiceBusClientBuilder()
     .credential(fullyQualifiedNamespace, credential)
     .sessionReceiver()
     .queueName(sessionEnabledQueueName)
     .disableAutoComplete()
     .buildClient();
 ServiceBusReceiverClient receiver = sessionReceiver.acceptSession("<<my-session-id>>");

 // Keep fetching messages from the session until there are no more messages.
 // The receiveMessage operation returns when either 10 messages have been receiver or, 30 seconds have elapsed.
 boolean hasMoreMessages = true;
 while (hasMoreMessages) {
     IterableStream<ServiceBusReceivedMessage> messages =
         receiver.receiveMessages(10, Duration.ofSeconds(30));
     Iterator<ServiceBusReceivedMessage> iterator = messages.iterator();
     hasMoreMessages = iterator.hasNext();

     while (iterator.hasNext()) {
         ServiceBusReceivedMessage message = iterator.next();
         System.out.printf("Session Id: %s. Contents: %s%n.", message.getSessionId(), message.getBody());

         // Explicitly settle the message using complete, abandon, defer, dead-letter, etc.
         if (isMessageProcessed) {
             receiver.complete(message);
         } else {
             receiver.abandon(message);
         }
     }
 }

 // Use the receiver and finally close it along with the sessionReceiver.
 receiver.close();
 sessionReceiver.close();

Sample: Receive messages from the first available session

Use acceptNextSession() to acquire the lock of the next available session without specifying the session id. PEEK_LOCK is strongly recommended so users have control over message settlement.

TokenCredential credential = new DefaultAzureCredentialBuilder().build();

 // 'fullyQualifiedNamespace' will look similar to "{your-namespace}.servicebus.windows.net"
 // 'disableAutoComplete' indicates that users will explicitly settle their message.
 ServiceBusSessionReceiverClient sessionReceiver = new ServiceBusClientBuilder()
     .credential(fullyQualifiedNamespace, credential)
     .sessionReceiver()
     .disableAutoComplete()
     .queueName(sessionEnabledQueueName)
     .buildClient();

 // Creates a client to receive messages from the first available session. It waits until
 // AmqpRetryOptions.getTryTimeout() elapses. If no session is available within that operation timeout, it
 // throws a retriable error. Otherwise, a receiver is returned when a lock on the session is acquired.
 ServiceBusReceiverClient receiver = sessionReceiver.acceptNextSession();

 // Use the receiver and finally close it along with the sessionReceiver.
 try {
     // Keep fetching messages from the session until there are no more messages.
     // The receiveMessage operation returns when either 10 messages have been receiver or, 30 seconds have elapsed.
     boolean hasMoreMessages = true;
     while (hasMoreMessages) {
         IterableStream<ServiceBusReceivedMessage> messages =
             receiver.receiveMessages(10, Duration.ofSeconds(30));
         Iterator<ServiceBusReceivedMessage> iterator = messages.iterator();
         hasMoreMessages = iterator.hasNext();

         while (iterator.hasNext()) {
             ServiceBusReceivedMessage message = iterator.next();
             System.out.printf("Session Id: %s. Message: %s%n.", message.getSessionId(), message.getBody());

             // Explicitly settle the message using complete, abandon, defer, dead-letter, etc.
             if (isMessageProcessed) {
                 receiver.complete(message);
             } else {
                 receiver.abandon(message);
             }
         }
     }
 } finally {
     receiver.close();
     sessionReceiver.close();
 }

Method Summary

Modifier and Type Method and Description
ServiceBusReceiverClient acceptNextSession()

Acquires a session lock for the next available session and creates a ServiceBusReceiverClient to receive messages from the session.

ServiceBusReceiverClient acceptSession(String sessionId)

Acquires a session lock for sessionId and create a ServiceBusReceiverClient to receive messages from the session.

void close()
PagedIterable<String> listSessions()

Lists the IDs of sessions that have active messages or stored session state in this entity.

PagedIterable<String> listSessions(OffsetDateTime sessionStateUpdatedAfter)

Lists the IDs of sessions whose state was set or updated after the specified time.

Methods inherited from java.lang.Object

Method Details

acceptNextSession

public ServiceBusReceiverClient acceptNextSession()

Acquires a session lock for the next available session and creates a ServiceBusReceiverClient to receive messages from the session. If no session is available immediately, it waits for one; if none becomes available within the operation timeout it throws a retriable timeout error rather than blocking indefinitely, so callers typically invoke this in a loop.

Returns:

A ServiceBusReceiverClient that is tied to the available session.

acceptSession

public ServiceBusReceiverClient acceptSession(String sessionId)

Acquires a session lock for sessionId and create a ServiceBusReceiverClient to receive messages from the session. If the session is already locked by another client, an AmqpException is thrown immediately.

Parameters:

sessionId - The session id.

Returns:

A ServiceBusReceiverClient that is tied to the specified session.

close

public void close()

listSessions

public PagedIterable<String> listSessions()

Lists the IDs of sessions that have active messages or stored session state in this entity.

Sessions with neither active messages nor stored session state are excluded.

The returned PagedIterable<T> fetches additional pages from the broker on demand; iterate the PagedIterable (or call PagedIterable#stream()) to receive every session ID. Pages are fetched lazily as the iterator advances. The default page size is 100; callers can request a different size via PagedIterable#iterableByPage(int) (or the equivalent on the underlying PagedFlux).

Returns:

A PagedIterable<T> of session ID strings.

listSessions

public PagedIterable<String> listSessions(OffsetDateTime sessionStateUpdatedAfter)

Lists the IDs of sessions whose state was set or updated after the specified time.

The returned PagedIterable<T> fetches additional pages from the broker on demand; iterate the PagedIterable (or call PagedIterable#stream()) to receive every session ID. Pages are fetched lazily as the iterator advances. The default page size is 100; callers can request a different size via PagedIterable#iterableByPage(int) (or the equivalent on the underlying PagedFlux).

Parameters:

sessionStateUpdatedAfter - Only sessions whose session state was set or updated after this time are returned.

Returns:

A PagedIterable<T> of session ID strings.

Applies to