From 9bc2418dcd56fc52bdac8bd69a359c4cc993e510 Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Wed, 4 Sep 2019 13:50:52 -0700 Subject: [PATCH] Update receive method --- .../azure-eventhubs/azure/eventhub/aio/consumer_async.py | 5 ++--- sdk/eventhub/azure-eventhubs/azure/eventhub/common.py | 5 +++++ sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py | 7 +++---- 3 files changed, 10 insertions(+), 7 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 f26853e32cac..147550c4d819 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py @@ -157,7 +157,7 @@ async def _open_with_retry(self): async def _receive(self, timeout_time=None, max_batch_size=None, **kwargs): last_exception = kwargs.get("last_exception") - data_batch = kwargs.get("data_batch") + data_batch = [] await self._open() remaining_time = timeout_time - time.time() @@ -225,9 +225,8 @@ async def receive(self, *, max_batch_size=None, timeout=None): timeout = timeout or self.client.config.receive_timeout max_batch_size = max_batch_size or min(self.client.config.max_batch_size, self.prefetch) - data_batch = [] # type: List[EventData] - return await self._receive_with_retry(timeout=timeout, max_batch_size=max_batch_size, data_batch=data_batch) + return await self._receive_with_retry(timeout=timeout, max_batch_size=max_batch_size) async def close(self, exception=None): # type: (Exception) -> None diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/common.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/common.py index 11668ba367f0..73fed892db11 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/common.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/common.py @@ -202,6 +202,11 @@ def application_properties(self, value): @property def system_properties(self): + """ + Metadata set by the Event Hubs Service associated with the EventData + + :rtype: dict + """ return self._annotations @property diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py index 82550bf3b9e5..e10f52e61b59 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py @@ -152,7 +152,7 @@ def _open_with_retry(self): def _receive(self, timeout_time=None, max_batch_size=None, **kwargs): last_exception = kwargs.get("last_exception") - data_batch = kwargs.get("data_batch") + data_batch = [] self._open() remaining_time = timeout_time - time.time() @@ -163,7 +163,7 @@ def _receive(self, timeout_time=None, max_batch_size=None, **kwargs): return data_batch remaining_time_ms = 1000 * remaining_time message_batch = self._handler.receive_message_batch( - max_batch_size=max_batch_size - (len(data_batch) if data_batch else 0), + max_batch_size=max_batch_size, timeout=remaining_time_ms) for message in message_batch: event_data = EventData._from_message(message) # pylint:disable=protected-access @@ -219,9 +219,8 @@ def receive(self, max_batch_size=None, timeout=None): timeout = timeout or self.client.config.receive_timeout max_batch_size = max_batch_size or min(self.client.config.max_batch_size, self.prefetch) - data_batch = [] # type: List[EventData] - return self._receive_with_retry(timeout=timeout, max_batch_size=max_batch_size, data_batch=data_batch) + return self._receive_with_retry(timeout=timeout, max_batch_size=max_batch_size) def close(self, exception=None): # type:(Exception) -> None