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 @@ -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()
Expand Down Expand Up @@ -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
Expand Down
5 changes: 5 additions & 0 deletions sdk/eventhub/azure-eventhubs/azure/eventhub/common.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
7 changes: 3 additions & 4 deletions sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down