From 0cdb291495337be3e606d332df729b48e724ee56 Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Tue, 17 Sep 2019 21:49:16 -0700 Subject: [PATCH 1/6] runtime metric init commit --- .../azure/eventhub/aio/consumer_async.py | 30 +++++++++++++++++ .../azure-eventhubs/azure/eventhub/common.py | 29 +++++++++++++++-- .../azure/eventhub/consumer.py | 32 ++++++++++++++++++- sdk/eventhub/azure-eventhubs/setup.py | 2 +- 4 files changed, 88 insertions(+), 5 deletions(-) diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py index efad6a3cb7db..4ec53371fa8b 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py @@ -36,6 +36,7 @@ class EventHubConsumer(ConsumerProducerMixin): # pylint:disable=too-many-instan _timeout = 0 _epoch_symbol = b'com.microsoft:epoch' _timeout_symbol = b'com.microsoft:timeout' + _receiver_runtime_metric_symbol = b'com.microsoft:enable-receiver-runtime-metric' def __init__( # pylint: disable=super-init-not-called self, client, source, **kwargs): @@ -62,6 +63,7 @@ def __init__( # pylint: disable=super-init-not-called owner_level = kwargs.get("owner_level", None) keep_alive = kwargs.get("keep_alive", None) auto_reconnect = kwargs.get("auto_reconnect", True) + track_last_enqueued_event_info = kwargs.get("track_last_enqueued_event_info", False) loop = kwargs.get("loop", None) super(EventHubConsumer, self).__init__() @@ -88,6 +90,8 @@ def __init__( # pylint: disable=super-init-not-called link_property_timeout_ms = (self._client._config.receive_timeout or self._timeout) * 1000 # pylint:disable=protected-access self._link_properties[types.AMQPSymbol(self._timeout_symbol)] = types.AMQPLong(int(link_property_timeout_ms)) self._handler = None + self._track_last_enqueued_event_info = track_last_enqueued_event_info + self._runtime_info = {} def __aiter__(self): return self @@ -104,6 +108,8 @@ async def __anext__(self): event_data = EventData._from_message(message) # pylint:disable=protected-access self._offset = EventPosition(event_data.offset, inclusive=False) retried_times = 0 + if self._track_last_enqueued_event_info: + self._runtime_info = event_data._runtime_info # pylint:disable=protected-access return event_data except Exception as exception: # pylint:disable=broad-except last_exception = await self._handle_exception(exception) @@ -122,6 +128,10 @@ def _create_handler(self): source = Source(self._source) if self._offset is not None: source.set_filter(self._offset._selector()) # pylint:disable=protected-access + + desired_capabilities = types.AMQPArray([types.AMQPSymbol(self._receiver_runtime_metric_symbol)])\ + if self._track_last_enqueued_event_info else None + self._handler = ReceiveClientAsync( source, auth=self._client._get_auth(**alt_creds), # pylint:disable=protected-access @@ -134,6 +144,7 @@ def _create_handler(self): client_name=self._name, properties=self._client._create_properties( # pylint:disable=protected-access self._client._config.user_agent), # pylint:disable=protected-access + desired_capabilities=desired_capabilities, # pylint:disable=protected-access loop=self._loop) self._messages_iter = None @@ -179,12 +190,31 @@ async def _receive(self, timeout_time=None, max_batch_size=None, **kwargs): event_data = EventData._from_message(message) # pylint:disable=protected-access self._offset = EventPosition(event_data.offset) data_batch.append(event_data) + + if self._track_last_enqueued_event_info and len(data_batch) > 0: + self._runtime_info = data_batch[-1]._runtime_info # pylint:disable=protected-access return data_batch async def _receive_with_retry(self, timeout=None, max_batch_size=None, **kwargs): return await self._do_retryable_operation(self._receive, timeout=timeout, max_batch_size=max_batch_size, **kwargs) + @property + def runtime_info(self): + """ + The latest enqueued event information. This property will be updated each time an event is received when + the receiver is created with `track_last_enqueued_event_info` being `True`. + The dict includes following information of the partition: + + - `last_enqueued_sequence_number` + - `last_enqueued_offset` + - `last_enqueued_time_utc` + - `runtime_info_retrieval_time_utc` + + :rtype: dict or None + """ + return self._runtime_info if self._track_last_enqueued_event_info else None + @property def queue_size(self): # type: () -> int diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/common.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/common.py index 5923d7f57972..d1675ddcc9b4 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/common.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/common.py @@ -56,6 +56,14 @@ class EventData(object): PROP_PARTITION_KEY_AMQP_SYMBOL = types.AMQPSymbol(PROP_PARTITION_KEY) PROP_TIMESTAMP = b"x-opt-enqueued-time" PROP_DEVICE_ID = b"iothub-connection-device-id" + PROP_LAST_ENQUEUED_SEQUENCE_NUMBER = b"last_enqueued_sequence_number" + PROP_LAST_ENQUEUED_SEQUENCE_NUMBER_AMQP_SYMBOL = types.AMQPSymbol(PROP_LAST_ENQUEUED_SEQUENCE_NUMBER) + PROP_LAST_ENQUEUED_OFFSET = b"last_enqueued_offset" + PROP_LAST_ENQUEUED_OFFSET_AMQP_SYMBOL = types.AMQPSymbol(PROP_LAST_ENQUEUED_OFFSET) + PROP_LAST_ENQUEUED_TIME_UTC = b"last_enqueued_time_utc" + PROP_LAST_ENQUEUED_TIME_UTC_AMQP_SYMBOL = types.AMQPSymbol(PROP_LAST_ENQUEUED_TIME_UTC) + PROP_RUNTIME_INFO_RETRIEVAL_TIME_UTC = b"runtime_info_retrieval_time_utc" + PROP_RUNTIME_INFO_RETRIEVAL_TIME_UTC_AMQP_SYMBOL = types.AMQPSymbol(PROP_RUNTIME_INFO_RETRIEVAL_TIME_UTC) def __init__(self, body=None, to_device=None): """ @@ -68,8 +76,10 @@ def __init__(self, body=None, to_device=None): """ self._annotations = {} + self._delivery_annotations = {} self._app_properties = {} self._msg_properties = MessageProperties() + self._runtime_info = {} if to_device: self._msg_properties.to = '/devices/{}/messages/devicebound'.format(to_device) if body and isinstance(body, list): @@ -116,11 +126,24 @@ def _set_partition_key(self, value): @staticmethod def _from_message(message): + # pylint:disable=protected-access event_data = EventData(body='') event_data.message = message - event_data._msg_properties = message.properties # pylint:disable=protected-access - event_data._annotations = message.annotations # pylint:disable=protected-access - event_data._app_properties = message.application_properties # pylint:disable=protected-access + event_data._msg_properties = message.properties + event_data._annotations = message.annotations + event_data._app_properties = message.application_properties + event_data._delivery_annotations = message.delivery_annotations + if event_data._delivery_annotations: + event_data._runtime_info = { + "last_enqueued_sequence_number": + event_data._delivery_annotations.get(EventData.PROP_LAST_ENQUEUED_SEQUENCE_NUMBER_AMQP_SYMBOL, None), + "last_enqueued_offset": + event_data._delivery_annotations.get(EventData.PROP_LAST_ENQUEUED_OFFSET, None), + "last_enqueued_time_utc": + event_data._delivery_annotations.get(EventData.PROP_LAST_ENQUEUED_TIME_UTC_AMQP_SYMBOL, None), + "runtime_info_retrieval_time_utc": + event_data._delivery_annotations.get(EventData.PROP_RUNTIME_INFO_RETRIEVAL_TIME_UTC_AMQP_SYMBOL, None) + } return event_data @property diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py index 604d9c7d7b82..cad633b84d9c 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py @@ -38,6 +38,7 @@ class EventHubConsumer(ConsumerProducerMixin): # pylint:disable=too-many-instan _timeout = 0 _epoch_symbol = b'com.microsoft:epoch' _timeout_symbol = b'com.microsoft:timeout' + _receiver_runtime_metric_symbol = b'com.microsoft:enable-receiver-runtime-metric' def __init__(self, client, source, **kwargs): """ @@ -60,6 +61,7 @@ def __init__(self, client, source, **kwargs): owner_level = kwargs.get("owner_level", None) keep_alive = kwargs.get("keep_alive", None) auto_reconnect = kwargs.get("auto_reconnect", True) + track_last_enqueued_event_info = kwargs.get("track_last_enqueued_event_info", False) super(EventHubConsumer, self).__init__() self._running = False @@ -84,6 +86,8 @@ def __init__(self, client, source, **kwargs): link_property_timeout_ms = (self._client._config.receive_timeout or self._timeout) * 1000 # pylint:disable=protected-access self._link_properties[types.AMQPSymbol(self._timeout_symbol)] = types.AMQPLong(int(link_property_timeout_ms)) self._handler = None + self._track_last_enqueued_event_info = track_last_enqueued_event_info + self._runtime_info = {} def __iter__(self): return self @@ -100,6 +104,8 @@ def __next__(self): event_data = EventData._from_message(message) # pylint:disable=protected-access self._offset = EventPosition(event_data.offset, inclusive=False) retried_times = 0 + if self._track_last_enqueued_event_info: + self._runtime_info = event_data._runtime_info # pylint:disable=protected-access return event_data except Exception as exception: # pylint:disable=broad-except last_exception = self._handle_exception(exception) @@ -118,6 +124,10 @@ def _create_handler(self): source = Source(self._source) if self._offset is not None: source.set_filter(self._offset._selector()) # pylint:disable=protected-access + + desired_capabilities = types.AMQPArray([types.AMQPSymbol(self._receiver_runtime_metric_symbol)])\ + if self._track_last_enqueued_event_info else None + self._handler = ReceiveClient( source, auth=self._client._get_auth(**alt_creds), # pylint:disable=protected-access @@ -129,7 +139,8 @@ def _create_handler(self): keep_alive_interval=self._keep_alive, client_name=self._name, properties=self._client._create_properties( # pylint:disable=protected-access - self._client._config.user_agent)) # pylint:disable=protected-access + self._client._config.user_agent), # pylint:disable=protected-access + desired_capabilities=desired_capabilities) # pylint:disable=protected-access self._messages_iter = None def _redirect(self, redirect): @@ -173,12 +184,31 @@ def _receive(self, timeout_time=None, max_batch_size=None, **kwargs): event_data = EventData._from_message(message) # pylint:disable=protected-access self._offset = EventPosition(event_data.offset) data_batch.append(event_data) + + if self._track_last_enqueued_event_info and len(data_batch) > 0: + self._runtime_info = data_batch[-1]._runtime_info # pylint:disable=protected-access return data_batch def _receive_with_retry(self, timeout=None, max_batch_size=None, **kwargs): return self._do_retryable_operation(self._receive, timeout=timeout, max_batch_size=max_batch_size, **kwargs) + @property + def runtime_info(self): + """ + The latest enqueued event information. This property will be updated each time an event is received when + the receiver is created with `track_last_enqueued_event_info` being `True`. + The dict includes following information of the partition: + + - `last_enqueued_sequence_number` + - `last_enqueued_offset` + - `last_enqueued_time_utc` + - `runtime_info_retrieval_time_utc` + + :rtype: dict or None + """ + return self._runtime_info if self._track_last_enqueued_event_info else None + @property def queue_size(self): # type:() -> int diff --git a/sdk/eventhub/azure-eventhubs/setup.py b/sdk/eventhub/azure-eventhubs/setup.py index 1ffa5c93005f..25ab4f45d657 100644 --- a/sdk/eventhub/azure-eventhubs/setup.py +++ b/sdk/eventhub/azure-eventhubs/setup.py @@ -67,7 +67,7 @@ zip_safe=False, packages=find_packages(exclude=exclude_packages), install_requires=[ - 'uamqp~=1.2.0', + 'uamqp>=1.2.3', 'azure-common~=1.1', ], extras_require={ From be5273f2608e257f5580639b3a91b468f6749905 Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Wed, 18 Sep 2019 01:22:43 -0700 Subject: [PATCH 2/6] evenhubts-runtime-metric implementation --- .../azure-eventhubs/azure/eventhub/aio/client_async.py | 4 +++- .../azure/eventhub/aio/consumer_async.py | 5 +++-- sdk/eventhub/azure-eventhubs/azure/eventhub/client.py | 4 +++- sdk/eventhub/azure-eventhubs/azure/eventhub/common.py | 10 +++------- .../azure-eventhubs/azure/eventhub/consumer.py | 5 +++-- 5 files changed, 15 insertions(+), 13 deletions(-) diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/client_async.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/client_async.py index 88b693d157ec..0ba3d94dc455 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/client_async.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/client_async.py @@ -266,6 +266,7 @@ def create_consumer( owner_level = kwargs.get("owner_level") operation = kwargs.get("operation") prefetch = kwargs.get("prefetch") or self._config.prefetch + track_last_enqueued_event_info = kwargs.get("track_last_enqueued_event_info", False) loop = kwargs.get("loop") path = self._address.path + operation if operation else self._address.path @@ -273,7 +274,8 @@ def create_consumer( self._address.hostname, path, consumer_group, partition_id) handler = EventHubConsumer( self, source_url, event_position=event_position, owner_level=owner_level, - prefetch=prefetch, loop=loop) + prefetch=prefetch, + track_last_enqueued_event_info=track_last_enqueued_event_info, loop=loop) return handler def create_producer( diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py index 4ec53371fa8b..08b53bb67718 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py @@ -8,7 +8,7 @@ from typing import List import time -from uamqp import errors, types # type: ignore +from uamqp import errors, types, utils # type: ignore from uamqp import ReceiveClientAsync, Source # type: ignore from azure.eventhub import EventData, EventPosition @@ -129,7 +129,8 @@ def _create_handler(self): if self._offset is not None: source.set_filter(self._offset._selector()) # pylint:disable=protected-access - desired_capabilities = types.AMQPArray([types.AMQPSymbol(self._receiver_runtime_metric_symbol)])\ + desired_capabilities = utils.data_factory( + types.AMQPArray([types.AMQPSymbol(self._receiver_runtime_metric_symbol)]))\ if self._track_last_enqueued_event_info else None self._handler = ReceiveClientAsync( diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/client.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/client.py index 06d264b5b9ac..6229848e2974 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/client.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/client.py @@ -263,13 +263,15 @@ def create_consumer(self, consumer_group, partition_id, event_position, **kwargs owner_level = kwargs.get("owner_level") operation = kwargs.get("operation") prefetch = kwargs.get("prefetch") or self._config.prefetch + track_last_enqueued_event_info = kwargs.get("track_last_enqueued_event_info", False) path = self._address.path + operation if operation else self._address.path source_url = "amqps://{}{}/ConsumerGroups/{}/Partitions/{}".format( self._address.hostname, path, consumer_group, partition_id) handler = EventHubConsumer( self, source_url, event_position=event_position, owner_level=owner_level, - prefetch=prefetch) + prefetch=prefetch, + track_last_enqueued_event_info=track_last_enqueued_event_info) return handler def create_producer(self, partition_id=None, operation=None, send_timeout=None): diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/common.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/common.py index d1675ddcc9b4..8b187fae3ad0 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/common.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/common.py @@ -57,13 +57,9 @@ class EventData(object): PROP_TIMESTAMP = b"x-opt-enqueued-time" PROP_DEVICE_ID = b"iothub-connection-device-id" PROP_LAST_ENQUEUED_SEQUENCE_NUMBER = b"last_enqueued_sequence_number" - PROP_LAST_ENQUEUED_SEQUENCE_NUMBER_AMQP_SYMBOL = types.AMQPSymbol(PROP_LAST_ENQUEUED_SEQUENCE_NUMBER) PROP_LAST_ENQUEUED_OFFSET = b"last_enqueued_offset" - PROP_LAST_ENQUEUED_OFFSET_AMQP_SYMBOL = types.AMQPSymbol(PROP_LAST_ENQUEUED_OFFSET) PROP_LAST_ENQUEUED_TIME_UTC = b"last_enqueued_time_utc" - PROP_LAST_ENQUEUED_TIME_UTC_AMQP_SYMBOL = types.AMQPSymbol(PROP_LAST_ENQUEUED_TIME_UTC) PROP_RUNTIME_INFO_RETRIEVAL_TIME_UTC = b"runtime_info_retrieval_time_utc" - PROP_RUNTIME_INFO_RETRIEVAL_TIME_UTC_AMQP_SYMBOL = types.AMQPSymbol(PROP_RUNTIME_INFO_RETRIEVAL_TIME_UTC) def __init__(self, body=None, to_device=None): """ @@ -136,13 +132,13 @@ def _from_message(message): if event_data._delivery_annotations: event_data._runtime_info = { "last_enqueued_sequence_number": - event_data._delivery_annotations.get(EventData.PROP_LAST_ENQUEUED_SEQUENCE_NUMBER_AMQP_SYMBOL, None), + event_data._delivery_annotations.get(EventData.PROP_LAST_ENQUEUED_SEQUENCE_NUMBER, None), "last_enqueued_offset": event_data._delivery_annotations.get(EventData.PROP_LAST_ENQUEUED_OFFSET, None), "last_enqueued_time_utc": - event_data._delivery_annotations.get(EventData.PROP_LAST_ENQUEUED_TIME_UTC_AMQP_SYMBOL, None), + event_data._delivery_annotations.get(EventData.PROP_LAST_ENQUEUED_TIME_UTC, None), "runtime_info_retrieval_time_utc": - event_data._delivery_annotations.get(EventData.PROP_RUNTIME_INFO_RETRIEVAL_TIME_UTC_AMQP_SYMBOL, None) + event_data._delivery_annotations.get(EventData.PROP_RUNTIME_INFO_RETRIEVAL_TIME_UTC, None) } return event_data diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py index cad633b84d9c..1d40d55c5a6b 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py @@ -9,7 +9,7 @@ import time from typing import List -from uamqp import types, errors # type: ignore +from uamqp import types, errors, utils # type: ignore from uamqp import ReceiveClient, Source # type: ignore from azure.eventhub.common import EventData, EventPosition @@ -125,7 +125,8 @@ def _create_handler(self): if self._offset is not None: source.set_filter(self._offset._selector()) # pylint:disable=protected-access - desired_capabilities = types.AMQPArray([types.AMQPSymbol(self._receiver_runtime_metric_symbol)])\ + desired_capabilities = utils.data_factory( + types.AMQPArray([types.AMQPSymbol(self._receiver_runtime_metric_symbol)]))\ if self._track_last_enqueued_event_info else None self._handler = ReceiveClient( From f471fff46a8ee95673fdb0e0dca8388d8e0e0059 Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Thu, 26 Sep 2019 12:23:17 -0700 Subject: [PATCH 3/6] Update code, and test and docstring --- .../azure/eventhub/aio/client_async.py | 8 +++++ .../azure/eventhub/aio/consumer_async.py | 9 +++--- .../azure-eventhubs/azure/eventhub/client.py | 8 +++++ .../azure/eventhub/consumer.py | 9 +++--- .../tests/asynctests/test_receive_async.py | 31 +++++++++++++++++++ .../azure-eventhubs/tests/test_receive.py | 30 ++++++++++++++++++ 6 files changed, 87 insertions(+), 8 deletions(-) diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/client_async.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/client_async.py index 0ba3d94dc455..4bfbee465f8b 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/client_async.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/client_async.py @@ -251,6 +251,14 @@ def create_consumer( :type operation: str :param prefetch: The message prefetch count of the consumer. Default is 300. :type prefetch: int + :type track_last_enqueued_event_info: bool + :param track_last_enqueued_event_info: Indicates whether or not the consumer should request information on the + last enqueued event on its associated partition, and track that information as events are received. + When information about the partition's last enqueued event is being tracked, each event received from the + Event Hubs service will carry metadata about the partition. This results in a small amount of additional + network bandwidth consumption that is generally a favorable trade-off when considered against periodically + making requests for partition properties using the Event Hub client. + It is set to `False` by default. :param loop: An event loop. If not specified the default event loop will be used. :rtype: ~azure.eventhub.aio.consumer_async.EventHubConsumer diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py index 08b53bb67718..a63891e0b50e 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py @@ -129,9 +129,10 @@ def _create_handler(self): if self._offset is not None: source.set_filter(self._offset._selector()) # pylint:disable=protected-access - desired_capabilities = utils.data_factory( - types.AMQPArray([types.AMQPSymbol(self._receiver_runtime_metric_symbol)]))\ - if self._track_last_enqueued_event_info else None + desired_capabilities = None + if self._track_last_enqueued_event_info: + symbol_array = [types.AMQPSymbol(self._receiver_runtime_metric_symbol)] + desired_capabilities = utils.data_factory(types.AMQPArray(symbol_array)) self._handler = ReceiveClientAsync( source, @@ -192,7 +193,7 @@ async def _receive(self, timeout_time=None, max_batch_size=None, **kwargs): self._offset = EventPosition(event_data.offset) data_batch.append(event_data) - if self._track_last_enqueued_event_info and len(data_batch) > 0: + if self._track_last_enqueued_event_info and len(data_batch): self._runtime_info = data_batch[-1]._runtime_info # pylint:disable=protected-access return data_batch diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/client.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/client.py index 6229848e2974..3c7381838d9e 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/client.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/client.py @@ -249,6 +249,14 @@ def create_consumer(self, consumer_group, partition_id, event_position, **kwargs :type operation: str :param prefetch: The message prefetch count of the consumer. Default is 300. :type prefetch: int + :type track_last_enqueued_event_info: bool + :param track_last_enqueued_event_info: Indicates whether or not the consumer should request information on the + last enqueued event on its associated partition, and track that information as events are received. + When information about the partition's last enqueued event is being tracked, each event received from the + Event Hubs service will carry metadata about the partition. This results in a small amount of additional + network bandwidth consumption that is generally a favorable trade-off when considered against periodically + making requests for partition properties using the Event Hub client. + It is set to `False` by default. :rtype: ~azure.eventhub.consumer.EventHubConsumer Example: diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py index 1d40d55c5a6b..89673bbfe6f6 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py @@ -125,9 +125,10 @@ def _create_handler(self): if self._offset is not None: source.set_filter(self._offset._selector()) # pylint:disable=protected-access - desired_capabilities = utils.data_factory( - types.AMQPArray([types.AMQPSymbol(self._receiver_runtime_metric_symbol)]))\ - if self._track_last_enqueued_event_info else None + desired_capabilities = None + if self._track_last_enqueued_event_info: + symbol_array = [types.AMQPSymbol(self._receiver_runtime_metric_symbol)] + desired_capabilities = utils.data_factory(types.AMQPArray(symbol_array)) self._handler = ReceiveClient( source, @@ -186,7 +187,7 @@ def _receive(self, timeout_time=None, max_batch_size=None, **kwargs): self._offset = EventPosition(event_data.offset) data_batch.append(event_data) - if self._track_last_enqueued_event_info and len(data_batch) > 0: + if self._track_last_enqueued_event_info and len(data_batch): self._runtime_info = data_batch[-1]._runtime_info # pylint:disable=protected-access return data_batch diff --git a/sdk/eventhub/azure-eventhubs/tests/asynctests/test_receive_async.py b/sdk/eventhub/azure-eventhubs/tests/asynctests/test_receive_async.py index 2a2e4836c2d5..97f3ac5e3faf 100644 --- a/sdk/eventhub/azure-eventhubs/tests/asynctests/test_receive_async.py +++ b/sdk/eventhub/azure-eventhubs/tests/asynctests/test_receive_async.py @@ -314,3 +314,34 @@ async def test_receive_over_websocket_async(connstr_senders): received = await receiver.receive(max_batch_size=50, timeout=5) assert len(received) == 20 + + +@pytest.mark.asyncio +@pytest.mark.liveTest +async def test_receive_run_time_metric_async(connstr_senders): + connection_str, senders = connstr_senders + client = EventHubClient.from_connection_string(connection_str, transport_type=TransportType.AmqpOverWebsocket, + network_tracing=False) + receiver = client.create_consumer(consumer_group="$default", partition_id="0", + event_position=EventPosition('@latest'), prefetch=500, + track_last_enqueued_event_info=True) + + event_list = [] + for i in range(20): + event_list.append(EventData("Event Number {}".format(i))) + + async with receiver: + received = await receiver.receive(timeout=5) + assert len(received) == 0 + + senders[0].send(event_list) + + await asyncio.sleep(1) + + received = await receiver.receive(max_batch_size=50, timeout=5) + assert len(received) == 20 + assert receiver.runtime_info + assert receiver.runtime_info.get('last_enqueued_sequence_number', None) + assert receiver.runtime_info.get('last_enqueued_offset', None) + assert receiver.runtime_info.get('last_enqueued_time_utc', None) + assert receiver.runtime_info.get('runtime_info_retrieval_time_utc', None) diff --git a/sdk/eventhub/azure-eventhubs/tests/test_receive.py b/sdk/eventhub/azure-eventhubs/tests/test_receive.py index d241a8e6e585..b0039c13b8f2 100644 --- a/sdk/eventhub/azure-eventhubs/tests/test_receive.py +++ b/sdk/eventhub/azure-eventhubs/tests/test_receive.py @@ -272,3 +272,33 @@ def test_receive_over_websocket_sync(connstr_senders): received = receiver.receive(max_batch_size=50, timeout=5) assert len(received) == 20 + + +@pytest.mark.liveTest +def test_receive_run_time_metric(connstr_senders): + connection_str, senders = connstr_senders + client = EventHubClient.from_connection_string(connection_str, transport_type=TransportType.AmqpOverWebsocket, + network_tracing=False) + receiver = client.create_consumer(consumer_group="$default", partition_id="0", + event_position=EventPosition('@latest'), prefetch=500, + track_last_enqueued_event_info=True) + + event_list = [] + for i in range(20): + event_list.append(EventData("Event Number {}".format(i))) + + with receiver: + received = receiver.receive(timeout=5) + assert len(received) == 0 + + senders[0].send(event_list) + + time.sleep(1) + + received = receiver.receive(max_batch_size=50, timeout=5) + assert len(received) == 20 + assert receiver.runtime_info + assert receiver.runtime_info.get('last_enqueued_sequence_number', None) + assert receiver.runtime_info.get('last_enqueued_offset', None) + assert receiver.runtime_info.get('last_enqueued_time_utc', None) + assert receiver.runtime_info.get('runtime_info_retrieval_time_utc', None) From 34097bce0908234bc0ea33e0dbac837f628f5a87 Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Thu, 26 Sep 2019 12:25:08 -0700 Subject: [PATCH 4/6] Add example code --- ...cv_track_last_enqueued_event_info_async.py | 48 +++++++++++++++++++ .../recv_track_last_enqueued_event_info.py | 45 +++++++++++++++++ 2 files changed, 93 insertions(+) create mode 100644 sdk/eventhub/azure-eventhubs/examples/async_examples/recv_track_last_enqueued_event_info_async.py create mode 100644 sdk/eventhub/azure-eventhubs/examples/recv_track_last_enqueued_event_info.py diff --git a/sdk/eventhub/azure-eventhubs/examples/async_examples/recv_track_last_enqueued_event_info_async.py b/sdk/eventhub/azure-eventhubs/examples/async_examples/recv_track_last_enqueued_event_info_async.py new file mode 100644 index 000000000000..eb9088497fef --- /dev/null +++ b/sdk/eventhub/azure-eventhubs/examples/async_examples/recv_track_last_enqueued_event_info_async.py @@ -0,0 +1,48 @@ +#!/usr/bin/env python + +# -------------------------------------------------------------------------------------------- +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the MIT License. See License.txt in the project root for license information. +# -------------------------------------------------------------------------------------------- + +""" +An example to show running concurrent consumers. +""" + +import os +import time +import asyncio + +from azure.eventhub.aio import EventHubClient +from azure.eventhub import EventPosition, EventHubSharedKeyCredential + +HOSTNAME = os.environ['EVENT_HUB_HOSTNAME'] # .servicebus.windows.net +EVENT_HUB = os.environ['EVENT_HUB_NAME'] +USER = os.environ['EVENT_HUB_SAS_POLICY'] +KEY = os.environ['EVENT_HUB_SAS_KEY'] + +EVENT_POSITION = EventPosition("-1") + + +async def pump(client, partition): + consumer = client.create_consumer(consumer_group="$default", partition_id=partition, event_position=EVENT_POSITION, + prefetch=5, track_last_enqueued_event_info=True) + async with consumer: + total = 0 + start_time = time.time() + for event_data in await consumer.receive(timeout=10): + last_offset = event_data.offset + last_sn = event_data.sequence_number + print("Received: {}, {}".format(last_offset, last_sn)) + total += 1 + end_time = time.time() + run_time = end_time - start_time + print("Consumer runtime information: {}.".format(consumer.runtime_info)) + print("Received {} messages in {} seconds".format(total, run_time)) + + +loop = asyncio.get_event_loop() +client = EventHubClient(host=HOSTNAME, event_hub_path=EVENT_HUB, credential=EventHubSharedKeyCredential(USER, KEY), + network_tracing=False) +tasks = [asyncio.ensure_future(pump(client, "0"))] +loop.run_until_complete(asyncio.wait(tasks)) diff --git a/sdk/eventhub/azure-eventhubs/examples/recv_track_last_enqueued_event_info.py b/sdk/eventhub/azure-eventhubs/examples/recv_track_last_enqueued_event_info.py new file mode 100644 index 000000000000..6cfc9a6f7315 --- /dev/null +++ b/sdk/eventhub/azure-eventhubs/examples/recv_track_last_enqueued_event_info.py @@ -0,0 +1,45 @@ +#!/usr/bin/env python + +# -------------------------------------------------------------------------------------------- +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the MIT License. See License.txt in the project root for license information. +# -------------------------------------------------------------------------------------------- + +""" +An example to show receiving events from an Event Hub partition. +""" +import os +import time +from azure.eventhub import EventHubClient, EventPosition, EventHubSharedKeyCredential + +HOSTNAME = os.environ['EVENT_HUB_HOSTNAME'] # .servicebus.windows.net +EVENT_HUB = os.environ['EVENT_HUB_NAME'] + +USER = os.environ['EVENT_HUB_SAS_POLICY'] +KEY = os.environ['EVENT_HUB_SAS_KEY'] + +EVENT_POSITION = EventPosition("-1") +PARTITION = "0" + + +total = 0 +last_sn = -1 +last_offset = "-1" +client = EventHubClient(host=HOSTNAME, event_hub_path=EVENT_HUB, credential=EventHubSharedKeyCredential(USER, KEY), + network_tracing=False) + +consumer = client.create_consumer(consumer_group="$default", partition_id=PARTITION, + event_position=EVENT_POSITION, prefetch=5000, + track_last_enqueued_event_info=True) +with consumer: + start_time = time.time() + batch = consumer.receive(timeout=5) + for event_data in batch: + last_offset = event_data.offset + last_sn = event_data.sequence_number + print("Received: {}, {}".format(last_offset, last_sn)) + print(event_data.body_as_str()) + total += 1 + batch = consumer.receive(timeout=5) + print("Consumer runtime information: {}.".format(consumer.runtime_info)) + print("Received {} messages in {} seconds".format(total, time.time() - start_time)) From 9928d9cd31b4e880246fae8f5225c73061c8e607 Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Thu, 26 Sep 2019 13:25:43 -0700 Subject: [PATCH 5/6] Update property name --- .../azure/eventhub/aio/consumer_async.py | 10 +++++----- .../azure-eventhubs/azure/eventhub/consumer.py | 10 +++++----- .../recv_track_last_enqueued_event_info_async.py | 2 +- .../examples/recv_track_last_enqueued_event_info.py | 2 +- .../tests/asynctests/test_receive_async.py | 10 +++++----- sdk/eventhub/azure-eventhubs/tests/test_receive.py | 10 +++++----- 6 files changed, 22 insertions(+), 22 deletions(-) diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py index a63891e0b50e..3fbe5b89e52c 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py @@ -91,7 +91,7 @@ def __init__( # pylint: disable=super-init-not-called self._link_properties[types.AMQPSymbol(self._timeout_symbol)] = types.AMQPLong(int(link_property_timeout_ms)) self._handler = None self._track_last_enqueued_event_info = track_last_enqueued_event_info - self._runtime_info = {} + self._last_enqueued_event_info = {} def __aiter__(self): return self @@ -109,7 +109,7 @@ async def __anext__(self): self._offset = EventPosition(event_data.offset, inclusive=False) retried_times = 0 if self._track_last_enqueued_event_info: - self._runtime_info = event_data._runtime_info # pylint:disable=protected-access + self._last_enqueued_event_info = event_data._runtime_info # pylint:disable=protected-access return event_data except Exception as exception: # pylint:disable=broad-except last_exception = await self._handle_exception(exception) @@ -194,7 +194,7 @@ async def _receive(self, timeout_time=None, max_batch_size=None, **kwargs): data_batch.append(event_data) if self._track_last_enqueued_event_info and len(data_batch): - self._runtime_info = data_batch[-1]._runtime_info # pylint:disable=protected-access + self._last_enqueued_event_info = data_batch[-1]._runtime_info # pylint:disable=protected-access return data_batch async def _receive_with_retry(self, timeout=None, max_batch_size=None, **kwargs): @@ -202,7 +202,7 @@ async def _receive_with_retry(self, timeout=None, max_batch_size=None, **kwargs) max_batch_size=max_batch_size, **kwargs) @property - def runtime_info(self): + def last_enqueued_event_info(self): """ The latest enqueued event information. This property will be updated each time an event is received when the receiver is created with `track_last_enqueued_event_info` being `True`. @@ -215,7 +215,7 @@ def runtime_info(self): :rtype: dict or None """ - return self._runtime_info if self._track_last_enqueued_event_info else None + return self._last_enqueued_event_info if self._track_last_enqueued_event_info else None @property def queue_size(self): diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py index 89673bbfe6f6..fe8baf267508 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py @@ -87,7 +87,7 @@ def __init__(self, client, source, **kwargs): self._link_properties[types.AMQPSymbol(self._timeout_symbol)] = types.AMQPLong(int(link_property_timeout_ms)) self._handler = None self._track_last_enqueued_event_info = track_last_enqueued_event_info - self._runtime_info = {} + self._last_enqueued_event_info = {} def __iter__(self): return self @@ -105,7 +105,7 @@ def __next__(self): self._offset = EventPosition(event_data.offset, inclusive=False) retried_times = 0 if self._track_last_enqueued_event_info: - self._runtime_info = event_data._runtime_info # pylint:disable=protected-access + self._last_enqueued_event_info_info = event_data._runtime_info # pylint:disable=protected-access return event_data except Exception as exception: # pylint:disable=broad-except last_exception = self._handle_exception(exception) @@ -188,7 +188,7 @@ def _receive(self, timeout_time=None, max_batch_size=None, **kwargs): data_batch.append(event_data) if self._track_last_enqueued_event_info and len(data_batch): - self._runtime_info = data_batch[-1]._runtime_info # pylint:disable=protected-access + self._last_enqueued_event_info = data_batch[-1]._runtime_info # pylint:disable=protected-access return data_batch def _receive_with_retry(self, timeout=None, max_batch_size=None, **kwargs): @@ -196,7 +196,7 @@ def _receive_with_retry(self, timeout=None, max_batch_size=None, **kwargs): max_batch_size=max_batch_size, **kwargs) @property - def runtime_info(self): + def last_enqueued_event_info(self): """ The latest enqueued event information. This property will be updated each time an event is received when the receiver is created with `track_last_enqueued_event_info` being `True`. @@ -209,7 +209,7 @@ def runtime_info(self): :rtype: dict or None """ - return self._runtime_info if self._track_last_enqueued_event_info else None + return self._last_enqueued_event_info if self._track_last_enqueued_event_info else None @property def queue_size(self): diff --git a/sdk/eventhub/azure-eventhubs/examples/async_examples/recv_track_last_enqueued_event_info_async.py b/sdk/eventhub/azure-eventhubs/examples/async_examples/recv_track_last_enqueued_event_info_async.py index eb9088497fef..2cd520fa51b5 100644 --- a/sdk/eventhub/azure-eventhubs/examples/async_examples/recv_track_last_enqueued_event_info_async.py +++ b/sdk/eventhub/azure-eventhubs/examples/async_examples/recv_track_last_enqueued_event_info_async.py @@ -37,7 +37,7 @@ async def pump(client, partition): total += 1 end_time = time.time() run_time = end_time - start_time - print("Consumer runtime information: {}.".format(consumer.runtime_info)) + print("Consumer runtime information: {}.".format(consumer.last_enqueued_event_info)) print("Received {} messages in {} seconds".format(total, run_time)) diff --git a/sdk/eventhub/azure-eventhubs/examples/recv_track_last_enqueued_event_info.py b/sdk/eventhub/azure-eventhubs/examples/recv_track_last_enqueued_event_info.py index 6cfc9a6f7315..e5f7db609a0d 100644 --- a/sdk/eventhub/azure-eventhubs/examples/recv_track_last_enqueued_event_info.py +++ b/sdk/eventhub/azure-eventhubs/examples/recv_track_last_enqueued_event_info.py @@ -41,5 +41,5 @@ print(event_data.body_as_str()) total += 1 batch = consumer.receive(timeout=5) - print("Consumer runtime information: {}.".format(consumer.runtime_info)) + print("Consumer runtime information: {}.".format(consumer.last_enqueued_event_info)) print("Received {} messages in {} seconds".format(total, time.time() - start_time)) diff --git a/sdk/eventhub/azure-eventhubs/tests/asynctests/test_receive_async.py b/sdk/eventhub/azure-eventhubs/tests/asynctests/test_receive_async.py index 97f3ac5e3faf..56bd197dee97 100644 --- a/sdk/eventhub/azure-eventhubs/tests/asynctests/test_receive_async.py +++ b/sdk/eventhub/azure-eventhubs/tests/asynctests/test_receive_async.py @@ -340,8 +340,8 @@ async def test_receive_run_time_metric_async(connstr_senders): received = await receiver.receive(max_batch_size=50, timeout=5) assert len(received) == 20 - assert receiver.runtime_info - assert receiver.runtime_info.get('last_enqueued_sequence_number', None) - assert receiver.runtime_info.get('last_enqueued_offset', None) - assert receiver.runtime_info.get('last_enqueued_time_utc', None) - assert receiver.runtime_info.get('runtime_info_retrieval_time_utc', None) + assert receiver.last_enqueued_event_info + assert receiver.last_enqueued_event_info.get('last_enqueued_sequence_number', None) + assert receiver.last_enqueued_event_info.get('last_enqueued_offset', None) + assert receiver.last_enqueued_event_info.get('last_enqueued_time_utc', None) + assert receiver.last_enqueued_event_info.get('runtime_info_retrieval_time_utc', None) diff --git a/sdk/eventhub/azure-eventhubs/tests/test_receive.py b/sdk/eventhub/azure-eventhubs/tests/test_receive.py index b0039c13b8f2..d49ef15dc2c8 100644 --- a/sdk/eventhub/azure-eventhubs/tests/test_receive.py +++ b/sdk/eventhub/azure-eventhubs/tests/test_receive.py @@ -297,8 +297,8 @@ def test_receive_run_time_metric(connstr_senders): received = receiver.receive(max_batch_size=50, timeout=5) assert len(received) == 20 - assert receiver.runtime_info - assert receiver.runtime_info.get('last_enqueued_sequence_number', None) - assert receiver.runtime_info.get('last_enqueued_offset', None) - assert receiver.runtime_info.get('last_enqueued_time_utc', None) - assert receiver.runtime_info.get('runtime_info_retrieval_time_utc', None) + assert receiver.last_enqueued_event_info + assert receiver.last_enqueued_event_info.get('last_enqueued_sequence_number', None) + assert receiver.last_enqueued_event_info.get('last_enqueued_offset', None) + assert receiver.last_enqueued_event_info.get('last_enqueued_time_utc', None) + assert receiver.last_enqueued_event_info.get('runtime_info_retrieval_time_utc', None) From bae543fc388e0794b17302d935c0fad4167cbf25 Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Fri, 27 Sep 2019 15:40:46 -0700 Subject: [PATCH 6/6] Update name to last enqueued event properties --- .../azure/eventhub/aio/client_async.py | 10 ++--- .../azure/eventhub/aio/consumer_async.py | 38 +++++++++++------- .../azure-eventhubs/azure/eventhub/client.py | 10 ++--- .../azure-eventhubs/azure/eventhub/common.py | 8 ++-- .../azure/eventhub/consumer.py | 40 +++++++++++-------- ...cv_track_last_enqueued_event_info_async.py | 4 +- .../recv_track_last_enqueued_event_info.py | 4 +- .../tests/asynctests/test_receive_async.py | 12 +++--- .../azure-eventhubs/tests/test_receive.py | 12 +++--- 9 files changed, 77 insertions(+), 61 deletions(-) diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/client_async.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/client_async.py index 4bfbee465f8b..578d0f26c059 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/client_async.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/client_async.py @@ -251,9 +251,9 @@ def create_consumer( :type operation: str :param prefetch: The message prefetch count of the consumer. Default is 300. :type prefetch: int - :type track_last_enqueued_event_info: bool - :param track_last_enqueued_event_info: Indicates whether or not the consumer should request information on the - last enqueued event on its associated partition, and track that information as events are received. + :type track_last_enqueued_event_properties: bool + :param track_last_enqueued_event_properties: Indicates whether or not the consumer should request information + on the last enqueued event on its associated partition, and track that information as events are received. When information about the partition's last enqueued event is being tracked, each event received from the Event Hubs service will carry metadata about the partition. This results in a small amount of additional network bandwidth consumption that is generally a favorable trade-off when considered against periodically @@ -274,7 +274,7 @@ def create_consumer( owner_level = kwargs.get("owner_level") operation = kwargs.get("operation") prefetch = kwargs.get("prefetch") or self._config.prefetch - track_last_enqueued_event_info = kwargs.get("track_last_enqueued_event_info", False) + track_last_enqueued_event_properties = kwargs.get("track_last_enqueued_event_properties", False) loop = kwargs.get("loop") path = self._address.path + operation if operation else self._address.path @@ -283,7 +283,7 @@ def create_consumer( handler = EventHubConsumer( self, source_url, event_position=event_position, owner_level=owner_level, prefetch=prefetch, - track_last_enqueued_event_info=track_last_enqueued_event_info, loop=loop) + track_last_enqueued_event_properties=track_last_enqueued_event_properties, loop=loop) return handler def create_producer( diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py index 3fbe5b89e52c..02b1bee781ec 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py @@ -56,6 +56,14 @@ def __init__( # pylint: disable=super-init-not-called :param owner_level: The priority of the exclusive consumer. An exclusive consumer will be created if owner_level is set. :type owner_level: int + :type track_last_enqueued_event_properties: bool + :param track_last_enqueued_event_properties: Indicates whether or not the consumer should request information + on the last enqueued event on its associated partition, and track that information as events are received. + When information about the partition's last enqueued event is being tracked, each event received from the + Event Hubs service will carry metadata about the partition. This results in a small amount of additional + network bandwidth consumption that is generally a favorable trade-off when considered against periodically + making requests for partition properties using the Event Hub client. + It is set to `False` by default. :param loop: An event loop. """ event_position = kwargs.get("event_position", None) @@ -63,7 +71,7 @@ def __init__( # pylint: disable=super-init-not-called owner_level = kwargs.get("owner_level", None) keep_alive = kwargs.get("keep_alive", None) auto_reconnect = kwargs.get("auto_reconnect", True) - track_last_enqueued_event_info = kwargs.get("track_last_enqueued_event_info", False) + track_last_enqueued_event_properties = kwargs.get("track_last_enqueued_event_properties", False) loop = kwargs.get("loop", None) super(EventHubConsumer, self).__init__() @@ -90,8 +98,8 @@ def __init__( # pylint: disable=super-init-not-called link_property_timeout_ms = (self._client._config.receive_timeout or self._timeout) * 1000 # pylint:disable=protected-access self._link_properties[types.AMQPSymbol(self._timeout_symbol)] = types.AMQPLong(int(link_property_timeout_ms)) self._handler = None - self._track_last_enqueued_event_info = track_last_enqueued_event_info - self._last_enqueued_event_info = {} + self._track_last_enqueued_event_properties = track_last_enqueued_event_properties + self._last_enqueued_event_properties = {} def __aiter__(self): return self @@ -108,8 +116,8 @@ async def __anext__(self): event_data = EventData._from_message(message) # pylint:disable=protected-access self._offset = EventPosition(event_data.offset, inclusive=False) retried_times = 0 - if self._track_last_enqueued_event_info: - self._last_enqueued_event_info = event_data._runtime_info # pylint:disable=protected-access + if self._track_last_enqueued_event_properties: + self._last_enqueued_event_properties = event_data._runtime_info # pylint:disable=protected-access return event_data except Exception as exception: # pylint:disable=broad-except last_exception = await self._handle_exception(exception) @@ -130,7 +138,7 @@ def _create_handler(self): source.set_filter(self._offset._selector()) # pylint:disable=protected-access desired_capabilities = None - if self._track_last_enqueued_event_info: + if self._track_last_enqueued_event_properties: symbol_array = [types.AMQPSymbol(self._receiver_runtime_metric_symbol)] desired_capabilities = utils.data_factory(types.AMQPArray(symbol_array)) @@ -193,8 +201,8 @@ async def _receive(self, timeout_time=None, max_batch_size=None, **kwargs): self._offset = EventPosition(event_data.offset) data_batch.append(event_data) - if self._track_last_enqueued_event_info and len(data_batch): - self._last_enqueued_event_info = data_batch[-1]._runtime_info # pylint:disable=protected-access + if self._track_last_enqueued_event_properties and len(data_batch): + self._last_enqueued_event_properties = data_batch[-1]._runtime_info # pylint:disable=protected-access return data_batch async def _receive_with_retry(self, timeout=None, max_batch_size=None, **kwargs): @@ -202,20 +210,20 @@ async def _receive_with_retry(self, timeout=None, max_batch_size=None, **kwargs) max_batch_size=max_batch_size, **kwargs) @property - def last_enqueued_event_info(self): + def last_enqueued_event_properties(self): """ The latest enqueued event information. This property will be updated each time an event is received when - the receiver is created with `track_last_enqueued_event_info` being `True`. + the receiver is created with `track_last_enqueued_event_properties` being `True`. The dict includes following information of the partition: - - `last_enqueued_sequence_number` - - `last_enqueued_offset` - - `last_enqueued_time_utc` - - `runtime_info_retrieval_time_utc` + - `sequence_number` + - `offset` + - `enqueued_time` + - `retrieval_time` :rtype: dict or None """ - return self._last_enqueued_event_info if self._track_last_enqueued_event_info else None + return self._last_enqueued_event_properties if self._track_last_enqueued_event_properties else None @property def queue_size(self): diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/client.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/client.py index 3c7381838d9e..b3ae901c3b55 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/client.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/client.py @@ -249,9 +249,9 @@ def create_consumer(self, consumer_group, partition_id, event_position, **kwargs :type operation: str :param prefetch: The message prefetch count of the consumer. Default is 300. :type prefetch: int - :type track_last_enqueued_event_info: bool - :param track_last_enqueued_event_info: Indicates whether or not the consumer should request information on the - last enqueued event on its associated partition, and track that information as events are received. + :type track_last_enqueued_event_properties: bool + :param track_last_enqueued_event_properties: Indicates whether or not the consumer should request information + on the last enqueued event on its associated partition, and track that information as events are received. When information about the partition's last enqueued event is being tracked, each event received from the Event Hubs service will carry metadata about the partition. This results in a small amount of additional network bandwidth consumption that is generally a favorable trade-off when considered against periodically @@ -271,7 +271,7 @@ def create_consumer(self, consumer_group, partition_id, event_position, **kwargs owner_level = kwargs.get("owner_level") operation = kwargs.get("operation") prefetch = kwargs.get("prefetch") or self._config.prefetch - track_last_enqueued_event_info = kwargs.get("track_last_enqueued_event_info", False) + track_last_enqueued_event_properties = kwargs.get("track_last_enqueued_event_properties", False) path = self._address.path + operation if operation else self._address.path source_url = "amqps://{}{}/ConsumerGroups/{}/Partitions/{}".format( @@ -279,7 +279,7 @@ def create_consumer(self, consumer_group, partition_id, event_position, **kwargs handler = EventHubConsumer( self, source_url, event_position=event_position, owner_level=owner_level, prefetch=prefetch, - track_last_enqueued_event_info=track_last_enqueued_event_info) + track_last_enqueued_event_properties=track_last_enqueued_event_properties) return handler def create_producer(self, partition_id=None, operation=None, send_timeout=None): diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/common.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/common.py index 8b187fae3ad0..fd262e77f833 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/common.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/common.py @@ -131,13 +131,13 @@ def _from_message(message): event_data._delivery_annotations = message.delivery_annotations if event_data._delivery_annotations: event_data._runtime_info = { - "last_enqueued_sequence_number": + "sequence_number": event_data._delivery_annotations.get(EventData.PROP_LAST_ENQUEUED_SEQUENCE_NUMBER, None), - "last_enqueued_offset": + "offset": event_data._delivery_annotations.get(EventData.PROP_LAST_ENQUEUED_OFFSET, None), - "last_enqueued_time_utc": + "enqueued_time": event_data._delivery_annotations.get(EventData.PROP_LAST_ENQUEUED_TIME_UTC, None), - "runtime_info_retrieval_time_utc": + "retrieval_time": event_data._delivery_annotations.get(EventData.PROP_RUNTIME_INFO_RETRIEVAL_TIME_UTC, None) } return event_data diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py index fe8baf267508..b983746c6d39 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py @@ -55,13 +55,21 @@ def __init__(self, client, source, **kwargs): :param owner_level: The priority of the exclusive consumer. An exclusive consumer will be created if owner_level is set. :type owner_level: int + :type track_last_enqueued_event_properties: bool + :param track_last_enqueued_event_properties: Indicates whether or not the consumer should request information + on the last enqueued event on its associated partition, and track that information as events are received. + When information about the partition's last enqueued event is being tracked, each event received from the + Event Hubs service will carry metadata about the partition. This results in a small amount of additional + network bandwidth consumption that is generally a favorable trade-off when considered against periodically + making requests for partition properties using the Event Hub client. + It is set to `False` by default. """ event_position = kwargs.get("event_position", None) prefetch = kwargs.get("prefetch", 300) owner_level = kwargs.get("owner_level", None) keep_alive = kwargs.get("keep_alive", None) auto_reconnect = kwargs.get("auto_reconnect", True) - track_last_enqueued_event_info = kwargs.get("track_last_enqueued_event_info", False) + track_last_enqueued_event_properties = kwargs.get("track_last_enqueued_event_properties", False) super(EventHubConsumer, self).__init__() self._running = False @@ -73,7 +81,7 @@ def __init__(self, client, source, **kwargs): self._owner_level = owner_level self._keep_alive = keep_alive self._auto_reconnect = auto_reconnect - self._retry_policy = errors.ErrorPolicy(max_retries=self._client._config.max_retries, on_error=_error_handler) # pylint:disable=protected-access + self._retry_policy = errors.ErrorPolicy(max_retries=self._client._config.max_retries, on_error=_error_handler) # pylint:disable=protected-access self._reconnect_backoff = 1 self._link_properties = {} self._redirected = None @@ -86,8 +94,8 @@ def __init__(self, client, source, **kwargs): link_property_timeout_ms = (self._client._config.receive_timeout or self._timeout) * 1000 # pylint:disable=protected-access self._link_properties[types.AMQPSymbol(self._timeout_symbol)] = types.AMQPLong(int(link_property_timeout_ms)) self._handler = None - self._track_last_enqueued_event_info = track_last_enqueued_event_info - self._last_enqueued_event_info = {} + self._track_last_enqueued_event_properties = track_last_enqueued_event_properties + self._last_enqueued_event_properties = {} def __iter__(self): return self @@ -104,8 +112,8 @@ def __next__(self): event_data = EventData._from_message(message) # pylint:disable=protected-access self._offset = EventPosition(event_data.offset, inclusive=False) retried_times = 0 - if self._track_last_enqueued_event_info: - self._last_enqueued_event_info_info = event_data._runtime_info # pylint:disable=protected-access + if self._track_last_enqueued_event_properties: + self._last_enqueued_event_properties = event_data._runtime_info # pylint:disable=protected-access return event_data except Exception as exception: # pylint:disable=broad-except last_exception = self._handle_exception(exception) @@ -126,7 +134,7 @@ def _create_handler(self): source.set_filter(self._offset._selector()) # pylint:disable=protected-access desired_capabilities = None - if self._track_last_enqueued_event_info: + if self._track_last_enqueued_event_properties: symbol_array = [types.AMQPSymbol(self._receiver_runtime_metric_symbol)] desired_capabilities = utils.data_factory(types.AMQPArray(symbol_array)) @@ -187,8 +195,8 @@ def _receive(self, timeout_time=None, max_batch_size=None, **kwargs): self._offset = EventPosition(event_data.offset) data_batch.append(event_data) - if self._track_last_enqueued_event_info and len(data_batch): - self._last_enqueued_event_info = data_batch[-1]._runtime_info # pylint:disable=protected-access + if self._track_last_enqueued_event_properties and len(data_batch): + self._last_enqueued_event_properties = data_batch[-1]._runtime_info # pylint:disable=protected-access return data_batch def _receive_with_retry(self, timeout=None, max_batch_size=None, **kwargs): @@ -196,20 +204,20 @@ def _receive_with_retry(self, timeout=None, max_batch_size=None, **kwargs): max_batch_size=max_batch_size, **kwargs) @property - def last_enqueued_event_info(self): + def last_enqueued_event_properties(self): """ The latest enqueued event information. This property will be updated each time an event is received when - the receiver is created with `track_last_enqueued_event_info` being `True`. + the receiver is created with `track_last_enqueued_event_properties` being `True`. The dict includes following information of the partition: - - `last_enqueued_sequence_number` - - `last_enqueued_offset` - - `last_enqueued_time_utc` - - `runtime_info_retrieval_time_utc` + - `sequence_number` + - `offset` + - `enqueued_time` + - `retrieval_time` :rtype: dict or None """ - return self._last_enqueued_event_info if self._track_last_enqueued_event_info else None + return self._last_enqueued_event_properties if self._track_last_enqueued_event_properties else None @property def queue_size(self): diff --git a/sdk/eventhub/azure-eventhubs/examples/async_examples/recv_track_last_enqueued_event_info_async.py b/sdk/eventhub/azure-eventhubs/examples/async_examples/recv_track_last_enqueued_event_info_async.py index 2cd520fa51b5..53d2e626a7f5 100644 --- a/sdk/eventhub/azure-eventhubs/examples/async_examples/recv_track_last_enqueued_event_info_async.py +++ b/sdk/eventhub/azure-eventhubs/examples/async_examples/recv_track_last_enqueued_event_info_async.py @@ -26,7 +26,7 @@ async def pump(client, partition): consumer = client.create_consumer(consumer_group="$default", partition_id=partition, event_position=EVENT_POSITION, - prefetch=5, track_last_enqueued_event_info=True) + prefetch=5, track_last_enqueued_event_properties=True) async with consumer: total = 0 start_time = time.time() @@ -37,7 +37,7 @@ async def pump(client, partition): total += 1 end_time = time.time() run_time = end_time - start_time - print("Consumer runtime information: {}.".format(consumer.last_enqueued_event_info)) + print("Consumer last enqueued event properties: {}.".format(consumer.last_enqueued_event_properties)) print("Received {} messages in {} seconds".format(total, run_time)) diff --git a/sdk/eventhub/azure-eventhubs/examples/recv_track_last_enqueued_event_info.py b/sdk/eventhub/azure-eventhubs/examples/recv_track_last_enqueued_event_info.py index e5f7db609a0d..576ef19089e6 100644 --- a/sdk/eventhub/azure-eventhubs/examples/recv_track_last_enqueued_event_info.py +++ b/sdk/eventhub/azure-eventhubs/examples/recv_track_last_enqueued_event_info.py @@ -30,7 +30,7 @@ consumer = client.create_consumer(consumer_group="$default", partition_id=PARTITION, event_position=EVENT_POSITION, prefetch=5000, - track_last_enqueued_event_info=True) + track_last_enqueued_event_properties=True) with consumer: start_time = time.time() batch = consumer.receive(timeout=5) @@ -41,5 +41,5 @@ print(event_data.body_as_str()) total += 1 batch = consumer.receive(timeout=5) - print("Consumer runtime information: {}.".format(consumer.last_enqueued_event_info)) + print("Consumer last enqueued event properties: {}.".format(consumer.last_enqueued_event_properties)) print("Received {} messages in {} seconds".format(total, time.time() - start_time)) diff --git a/sdk/eventhub/azure-eventhubs/tests/asynctests/test_receive_async.py b/sdk/eventhub/azure-eventhubs/tests/asynctests/test_receive_async.py index 56bd197dee97..a9744f259c8c 100644 --- a/sdk/eventhub/azure-eventhubs/tests/asynctests/test_receive_async.py +++ b/sdk/eventhub/azure-eventhubs/tests/asynctests/test_receive_async.py @@ -324,7 +324,7 @@ async def test_receive_run_time_metric_async(connstr_senders): network_tracing=False) receiver = client.create_consumer(consumer_group="$default", partition_id="0", event_position=EventPosition('@latest'), prefetch=500, - track_last_enqueued_event_info=True) + track_last_enqueued_event_properties=True) event_list = [] for i in range(20): @@ -340,8 +340,8 @@ async def test_receive_run_time_metric_async(connstr_senders): received = await receiver.receive(max_batch_size=50, timeout=5) assert len(received) == 20 - assert receiver.last_enqueued_event_info - assert receiver.last_enqueued_event_info.get('last_enqueued_sequence_number', None) - assert receiver.last_enqueued_event_info.get('last_enqueued_offset', None) - assert receiver.last_enqueued_event_info.get('last_enqueued_time_utc', None) - assert receiver.last_enqueued_event_info.get('runtime_info_retrieval_time_utc', None) + assert receiver.last_enqueued_event_properties + assert receiver.last_enqueued_event_properties.get('sequence_number', None) + assert receiver.last_enqueued_event_properties.get('offset', None) + assert receiver.last_enqueued_event_properties.get('enqueued_time', None) + assert receiver.last_enqueued_event_properties.get('retrieval_time', None) diff --git a/sdk/eventhub/azure-eventhubs/tests/test_receive.py b/sdk/eventhub/azure-eventhubs/tests/test_receive.py index d49ef15dc2c8..c9da841f71b2 100644 --- a/sdk/eventhub/azure-eventhubs/tests/test_receive.py +++ b/sdk/eventhub/azure-eventhubs/tests/test_receive.py @@ -281,7 +281,7 @@ def test_receive_run_time_metric(connstr_senders): network_tracing=False) receiver = client.create_consumer(consumer_group="$default", partition_id="0", event_position=EventPosition('@latest'), prefetch=500, - track_last_enqueued_event_info=True) + track_last_enqueued_event_properties=True) event_list = [] for i in range(20): @@ -297,8 +297,8 @@ def test_receive_run_time_metric(connstr_senders): received = receiver.receive(max_batch_size=50, timeout=5) assert len(received) == 20 - assert receiver.last_enqueued_event_info - assert receiver.last_enqueued_event_info.get('last_enqueued_sequence_number', None) - assert receiver.last_enqueued_event_info.get('last_enqueued_offset', None) - assert receiver.last_enqueued_event_info.get('last_enqueued_time_utc', None) - assert receiver.last_enqueued_event_info.get('runtime_info_retrieval_time_utc', None) + assert receiver.last_enqueued_event_properties + assert receiver.last_enqueued_event_properties.get('sequence_number', None) + assert receiver.last_enqueued_event_properties.get('offset', None) + assert receiver.last_enqueued_event_properties.get('enqueued_time', None) + assert receiver.last_enqueued_event_properties.get('retrieval_time', None)