From de72e14741916acd2f0ab81a3ccb3788b4613964 Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Thu, 1 Aug 2019 17:42:00 -0700 Subject: [PATCH 1/3] Update code according to review --- sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py index dd4b7aa2396a..98508f790860 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py @@ -104,7 +104,7 @@ def _create_handler(self): link_properties=self._link_properties, properties=self.client._create_properties(self.client.config.user_agent)) # pylint: disable=protected-access - def _open(self, timeout_time=None, **kwargs): + def _open(self, timeout_time=None): """ Open the EventHubProducer using the supplied connection. If the handler has previously been redirected, the redirect From bb2f7cef037a44a4d566dec413f20321e5c747bf Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Thu, 1 Aug 2019 17:58:08 -0700 Subject: [PATCH 2/3] Small fix --- sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py index 98508f790860..dd4b7aa2396a 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py @@ -104,7 +104,7 @@ def _create_handler(self): link_properties=self._link_properties, properties=self.client._create_properties(self.client.config.user_agent)) # pylint: disable=protected-access - def _open(self, timeout_time=None): + def _open(self, timeout_time=None, **kwargs): """ Open the EventHubProducer using the supplied connection. If the handler has previously been redirected, the redirect From dc6d61c78d34a6b3af7716548d437f4a6e9ca5ce Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Fri, 2 Aug 2019 12:02:43 -0700 Subject: [PATCH 3/3] Update decorator implementation --- .../azure/eventhub/_consumer_producer_mixin.py | 2 +- .../eventhub/aio/_consumer_producer_mixin_async.py | 2 +- .../azure/eventhub/aio/consumer_async.py | 11 ++++++----- .../azure/eventhub/aio/producer_async.py | 12 ++++++++++-- .../azure-eventhubs/azure/eventhub/consumer.py | 10 ++++++---- .../azure-eventhubs/azure/eventhub/producer.py | 12 ++++++++++-- 6 files changed, 34 insertions(+), 15 deletions(-) diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/_consumer_producer_mixin.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/_consumer_producer_mixin.py index bebef7a51982..3659b7c680c1 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/_consumer_producer_mixin.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/_consumer_producer_mixin.py @@ -25,7 +25,7 @@ def wrapped_func(self, *args, **kwargs): kwargs.pop("timeout", None) while True: try: - return to_be_wrapped_func(timeout_time=timeout_time, last_exception=last_exception, **kwargs) + return to_be_wrapped_func(self, timeout_time=timeout_time, last_exception=last_exception, **kwargs) except Exception as exception: last_exception = self._handle_exception(exception, retry_count, max_retries, timeout_time) retry_count += 1 diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/_consumer_producer_mixin_async.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/_consumer_producer_mixin_async.py index aa539110e50a..21dc6bb8a012 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/_consumer_producer_mixin_async.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/_consumer_producer_mixin_async.py @@ -25,7 +25,7 @@ async def wrapped_func(self, *args, **kwargs): kwargs.pop("timeout", None) while True: try: - return await to_be_wrapped_func(timeout_time=timeout_time, last_exception=last_exception, **kwargs) + return await to_be_wrapped_func(self, timeout_time=timeout_time, last_exception=last_exception, **kwargs) except Exception as exception: last_exception = await self._handle_exception(exception, retry_count, max_retries, timeout_time) retry_count += 1 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 8457913abcf0..404fc23312f0 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py @@ -147,10 +147,8 @@ async def _open(self, timeout_time=None): self.source = self.redirected.address await super(EventHubConsumer, self)._open(timeout_time) - async def _receive(self, **kwargs): - timeout_time = kwargs.get("timeout_time") + async def _receive(self, timeout_time=None, max_batch_size=None, **kwargs): last_exception = kwargs.get("last_exception") - max_batch_size = kwargs.get("max_batch_size") data_batch = kwargs.get("data_batch") await self._open(timeout_time) @@ -171,6 +169,10 @@ async def _receive(self, **kwargs): data_batch.append(event_data) return data_batch + @_retry_decorator + async def _receive_with_try(self, timeout_time=None, max_batch_size=None, **kwargs): + return await self._receive(timeout_time=timeout_time, max_batch_size=max_batch_size, **kwargs) + @property def queue_size(self): # type: () -> int @@ -217,8 +219,7 @@ async def receive(self, *, max_batch_size=None, timeout=None): max_batch_size = max_batch_size or min(self.client.config.max_batch_size, self.prefetch) data_batch = [] # type: List[EventData] - return await _retry_decorator(self._receive)(self, timeout=timeout, - max_batch_size=max_batch_size, data_batch=data_batch) + return await self._receive_with_try(timeout=timeout, max_batch_size=max_batch_size, data_batch=data_batch) async def close(self, exception=None): # type: (Exception) -> None diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/producer_async.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/producer_async.py index e3cd1d9fcb09..10de6017026b 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/producer_async.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/producer_async.py @@ -110,6 +110,10 @@ async def _open(self, timeout_time=None, **kwargs): self.target = self.redirected.address await super(EventHubProducer, self)._open(timeout_time) + @_retry_decorator + async def _open_with_retry(self, timeout_time=None, **kwargs): + return await self._open(timeout_time=timeout_time, **kwargs) + async def _send_event_data(self, timeout_time=None, last_exception=None): if self.unsent_events: await self._open(timeout_time) @@ -131,6 +135,10 @@ async def _send_event_data(self, timeout_time=None, last_exception=None): _error(self._outcome, self._condition) return + @_retry_decorator + async def _send_event_data_with_retry(self, timeout_time=None, last_exception=None): + return await self._send_event_data(timeout_time=timeout_time, last_exception=last_exception) + def _on_outcome(self, outcome, condition): """ Called when the outcome is received for a delivery. @@ -158,7 +166,7 @@ async def create_batch(self, max_size=None, partition_key=None): """ if not self._max_message_size_on_link: - await _retry_decorator(self._open)(self, timeout=self.client.config.send_timeout) + await self._open_with_retry(timeout=self.client.config.send_timeout) if max_size and max_size > self._max_message_size_on_link: raise ValueError('Max message size: {} is too large, acceptable max batch size is: {} bytes.' @@ -212,7 +220,7 @@ async def send(self, event_data, *, partition_key=None, timeout=None): wrapper_event_data = EventDataBatch._from_batch(event_data, partition_key) # pylint: disable=protected-access wrapper_event_data.message.on_send_complete = self._on_outcome self.unsent_events = [wrapper_event_data.message] - await _retry_decorator(self._send_event_data)(self, timeout=timeout) + await self._send_event_data_with_retry(timeout=timeout) async def close(self, exception=None): # type: (Exception) -> None diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py index 33a0e8ed187e..85c09bf1a308 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py @@ -142,10 +142,8 @@ def _open(self, timeout_time=None): self.source = self.redirected.address super(EventHubConsumer, self)._open(timeout_time) - def _receive(self, **kwargs): - timeout_time = kwargs.get("timeout_time") + def _receive(self, timeout_time=None, max_batch_size=None, **kwargs): last_exception = kwargs.get("last_exception") - max_batch_size = kwargs.get("max_batch_size") data_batch = kwargs.get("data_batch") self._open(timeout_time) @@ -165,6 +163,10 @@ def _receive(self, **kwargs): data_batch.append(event_data) return data_batch + @_retry_decorator + def _receive_with_try(self, timeout_time=None, max_batch_size=None, **kwargs): + return self._receive(timeout_time=timeout_time, max_batch_size=max_batch_size, **kwargs) + @property def queue_size(self): # type:() -> int @@ -210,7 +212,7 @@ def receive(self, max_batch_size=None, timeout=None): max_batch_size = max_batch_size or min(self.client.config.max_batch_size, self.prefetch) data_batch = [] # type: List[EventData] - return _retry_decorator(self._receive)(self, timeout=timeout, max_batch_size=max_batch_size, data_batch=data_batch) + return self._receive_with_try(timeout=timeout, max_batch_size=max_batch_size, data_batch=data_batch) def close(self, exception=None): # type:(Exception) -> None diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py index dd4b7aa2396a..f49fc280814e 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py @@ -117,6 +117,10 @@ def _open(self, timeout_time=None, **kwargs): self.target = self.redirected.address super(EventHubProducer, self)._open(timeout_time) + @_retry_decorator + def _open_with_retry(self, timeout_time=None, **kwargs): + return self._open(timeout_time=timeout_time, **kwargs) + def _send_event_data(self, timeout_time=None, last_exception=None): if self.unsent_events: self._open(timeout_time) @@ -138,6 +142,10 @@ def _send_event_data(self, timeout_time=None, last_exception=None): _error(self._outcome, self._condition) return + @_retry_decorator + def _send_event_data_with_retry(self, timeout_time=None, last_exception=None): + return self._send_event_data(timeout_time=timeout_time, last_exception=last_exception) + def _on_outcome(self, outcome, condition): """ Called when the outcome is received for a delivery. @@ -165,7 +173,7 @@ def create_batch(self, max_size=None, partition_key=None): """ if not self._max_message_size_on_link: - _retry_decorator(self._open)(self, timeout=self.client.config.send_timeout) + self._open_with_retry(timeout=self.client.config.send_timeout) if max_size and max_size > self._max_message_size_on_link: raise ValueError('Max message size: {} is too large, acceptable max batch size is: {} bytes.' @@ -219,7 +227,7 @@ def send(self, event_data, partition_key=None, timeout=None): wrapper_event_data = EventDataBatch._from_batch(event_data, partition_key) # pylint: disable=protected-access wrapper_event_data.message.on_send_complete = self._on_outcome self.unsent_events = [wrapper_event_data.message] - _retry_decorator(self._send_event_data)(self, timeout=timeout) + self._send_event_data_with_retry(timeout=timeout) def close(self, exception=None): # type:(Exception) -> None