Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -481,6 +481,62 @@ def send_message( # type: ignore
except StorageErrorException as error:
process_storage_error(error)

@distributed_trace
def receive_message(self, **kwargs):
# type: (Optional[Any]) -> QueueMessage
"""Removes one message from the front of the queue.

When the message is retrieved from the queue, the response includes the message
content and a pop_receipt value, which is required to delete the message.
The message is not automatically deleted from the queue, but after it has
been retrieved, it is not visible to other clients for the time interval
specified by the visibility_timeout parameter.

If the key-encryption-key or resolver field is set on the local service object, the message will be
decrypted before being returned.

:keyword int visibility_timeout:
If not specified, the default value is 0. Specifies the
new visibility timeout value, in seconds, relative to server time.
The value must be larger than or equal to 0, and cannot be
larger than 7 days. The visibility timeout of a message cannot be
set to a value later than the expiry time. visibility_timeout
should be set to a value smaller than the time-to-live value.
:keyword int timeout:
The server timeout, expressed in seconds.
:return:
Returns a message from the Queue.
:rtype: ~azure.storage.queue.QueueMessage

.. admonition:: Example:

.. literalinclude:: ../samples/queue_samples_message.py
:start-after: [START receive_one_message]
:end-before: [END receive_one_message]
:language: python
:dedent: 12
:caption: Receive one message from the queue.
"""
visibility_timeout = kwargs.pop('visibility_timeout', None)
timeout = kwargs.pop('timeout', None)
self._config.message_decode_policy.configure(
require_encryption=self.require_encryption,
key_encryption_key=self.key_encryption_key,
resolver=self.key_resolver_function)
try:
message = self._client.messages.dequeue(
number_of_messages=1,
visibilitytimeout=visibility_timeout,
timeout=timeout,
cls=self._config.message_decode_policy,
**kwargs
)
wrapped_message = QueueMessage._from_generated( # pylint: disable=protected-access
message[0]) if message != [] else None
return wrapped_message
except StorageErrorException as error:
process_storage_error(error)

@distributed_trace
def receive_messages(self, **kwargs):
# type: (Optional[Any]) -> ItemPaged[QueueMessage]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -400,6 +400,62 @@ async def send_message( # type: ignore
except StorageErrorException as error:
process_storage_error(error)

@distributed_trace_async
async def receive_message(self, **kwargs):
# type: (Optional[Any]) -> QueueMessage
"""Removes one message from the front of the queue.

When the message is retrieved from the queue, the response includes the message
content and a pop_receipt value, which is required to delete the message.
The message is not automatically deleted from the queue, but after it has
been retrieved, it is not visible to other clients for the time interval
specified by the visibility_timeout parameter.

If the key-encryption-key or resolver field is set on the local service object, the message will be
decrypted before being returned.

:keyword int visibility_timeout:
If not specified, the default value is 0. Specifies the
new visibility timeout value, in seconds, relative to server time.
The value must be larger than or equal to 0, and cannot be
larger than 7 days. The visibility timeout of a message cannot be
set to a value later than the expiry time. visibility_timeout
should be set to a value smaller than the time-to-live value.
:keyword int timeout:
The server timeout, expressed in seconds.
:return:
Returns a message from the Queue.
:rtype: ~azure.storage.queue.QueueMessage

.. admonition:: Example:

.. literalinclude:: ../samples/queue_samples_message_async.py
:start-after: [START receive_one_message]
:end-before: [END receive_one_message]
:language: python
:dedent: 12
:caption: Receive one message from the queue.
"""
visibility_timeout = kwargs.pop('visibility_timeout', None)
timeout = kwargs.pop('timeout', None)
self._config.message_decode_policy.configure(
require_encryption=self.require_encryption,
key_encryption_key=self.key_encryption_key,
resolver=self.key_resolver_function)
try:
message = await self._client.messages.dequeue(
number_of_messages=1,
visibilitytimeout=visibility_timeout,
timeout=timeout,
cls=self._config.message_decode_policy,
**kwargs
)
wrapped_message = QueueMessage._from_generated( # pylint: disable=protected-access
message[0]) if message != [] else None
return wrapped_message
except StorageErrorException as error:
process_storage_error(error)

@distributed_trace
def receive_messages(self, **kwargs):
# type: (Optional[Any]) -> AsyncItemPaged[QueueMessage]
Expand Down
35 changes: 32 additions & 3 deletions sdk/storage/azure-storage-queue/samples/queue_samples_message.py
Original file line number Diff line number Diff line change
Expand Up @@ -179,14 +179,42 @@ def list_message_pages(self):
finally:
queue.delete_queue()

def delete_and_clear_messages(self):
def receive_one_message_from_queue(self):
# Instantiate a queue client
from azure.storage.queue import QueueClient
queue = QueueClient.from_connection_string(self.connection_string, "myqueue5")

# Create the queue
queue.create_queue()

try:
queue.send_message(u"message1")
queue.send_message(u"message2")
queue.send_message(u"message3")

# [START receive_one_message]
# Pop two messages from the front of the queue
message1 = queue.receive_message()
message2 = queue.receive_message()
# We should see message 3 if we peek
message3 = queue.peek_messages()[0]

print(message1.content)
print(message2.content)
print(message3.content)
# [END receive_one_message]

finally:
queue.delete_queue()

def delete_and_clear_messages(self):
# Instantiate a queue client
from azure.storage.queue import QueueClient
queue = QueueClient.from_connection_string(self.connection_string, "myqueue6")

# Create the queue
queue.create_queue()

try:
# Send messages
queue.send_message(u"message1")
Expand Down Expand Up @@ -214,7 +242,7 @@ def delete_and_clear_messages(self):
def peek_messages(self):
# Instantiate a queue client
from azure.storage.queue import QueueClient
queue = QueueClient.from_connection_string(self.connection_string, "myqueue6")
queue = QueueClient.from_connection_string(self.connection_string, "myqueue7")

# Create the queue
queue.create_queue()
Expand Down Expand Up @@ -246,7 +274,7 @@ def peek_messages(self):
def update_message(self):
# Instantiate a queue client
from azure.storage.queue import QueueClient
queue = QueueClient.from_connection_string(self.connection_string, "myqueue7")
queue = QueueClient.from_connection_string(self.connection_string, "myqueue8")

# Create the queue
queue.create_queue()
Expand Down Expand Up @@ -279,6 +307,7 @@ def update_message(self):
sample.queue_metadata()
sample.send_and_receive_messages()
sample.list_message_pages()
sample.receive_one_message_from_queue()
sample.delete_and_clear_messages()
sample.peek_messages()
sample.update_message()
Original file line number Diff line number Diff line change
Expand Up @@ -149,6 +149,36 @@ async def send_and_receive_messages_async(self):
# Delete the queue
await queue.delete_queue()

async def receive_one_message_from_queue(self):
# Instantiate a queue client
from azure.storage.queue.aio import QueueClient
queue = QueueClient.from_connection_string(self.connection_string, "myqueue3")

# Create the queue
async with queue:
await queue.create_queue()

try:
await asyncio.gather(
queue.send_message(u"message1"),
queue.send_message(u"message2"),
queue.send_message(u"message3"))

# [START receive_one_message]
# Pop two messages from the front of the queue
message1 = await queue.receive_message()
message2 = await queue.receive_message()
# We should see message 3 if we peek
message3 = await queue.peek_messages()

print(message1.content)
print(message2.content)
print(message3[0].content)
# [END receive_one_message]

finally:
await queue.delete_queue()

async def delete_and_clear_messages_async(self):
# Instantiate a queue client
from azure.storage.queue.aio import QueueClient
Expand Down Expand Up @@ -256,6 +286,7 @@ async def main():
await sample.set_access_policy_async()
await sample.queue_metadata_async()
await sample.send_and_receive_messages_async()
await sample.receive_one_message_from_queue()
await sample.delete_and_clear_messages_async()
await sample.peek_messages_async()
await sample.update_message_async()
Expand Down
Loading