diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/_connection_manager.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/_connection_manager.py index 9600226df2fb..166703d698e7 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/_connection_manager.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/_connection_manager.py @@ -74,4 +74,4 @@ def reset_connection_if_broken(self): def get_connection_manager(**kwargs): - return _SharedConnectionManager(**kwargs) + return _SeparateConnectionManager(**kwargs) 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 20462a7753f4..17479119ccad 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/client_async.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/client_async.py @@ -194,8 +194,7 @@ async def get_partition_properties(self, partition): return output def create_consumer( - self, consumer_group, partition_id, event_position, owner_level=None, - operation=None, prefetch=None, loop=None): + self, consumer_group, partition_id, event_position, **kwargs): # type: (str, str, EventPosition, int, str, int, asyncio.AbstractEventLoop) -> EventHubConsumer """ Create an async consumer to the client for a particular consumer group and partition. @@ -227,8 +226,12 @@ def create_consumer( :caption: Add an async consumer to the client for a particular consumer group and partition. """ - prefetch = self.config.prefetch if prefetch is None else prefetch + owner_level = kwargs.get("owner_level", None) + operation = kwargs.get("operation", None) + prefetch = kwargs.get("prefetch", None) + loop = kwargs.get("loop", None) + prefetch = prefetch or self.config.prefetch 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) @@ -238,7 +241,7 @@ def create_consumer( return handler def create_producer( - self, partition_id=None, operation=None, send_timeout=None, loop=None): + self, **kwargs): # type: (str, str, float, asyncio.AbstractEventLoop) -> EventHubProducer """ Create an async producer to send EventData object to an EventHub. @@ -265,6 +268,11 @@ def create_producer( :caption: Add an async producer to the client to send EventData. """ + partition_id = kwargs.get("partition_id", None) + operation = kwargs.get("operation", None) + send_timeout = kwargs.get("send_timeout", None) + loop = kwargs.get("loop", None) + target = "amqps://{}{}".format(self.address.hostname, self.address.path) if operation: target = target + operation 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 5cc8df6941a0..d4d4143810af 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py @@ -22,15 +22,15 @@ class EventHubConsumer(ConsumerProducerMixin): """ A consumer responsible for reading EventData from a specific Event Hub - partition and as a member of a specific consumer group. + partition and as a member of a specific consumer group. A consumer may be exclusive, which asserts ownership over the partition for the consumer - group to ensure that only one consumer from that group is reading the from the partition. - These exclusive consumers are sometimes referred to as "Epoch Consumers." + group to ensure that only one consumer from that group is reading the from the partition. + These exclusive consumers are sometimes referred to as "Epoch Consumers." A consumer may also be non-exclusive, allowing multiple consumers from the same consumer - group to be actively reading events from the partition. These non-exclusive consumers are - sometimes referred to as "Non-Epoch Consumers." + group to be actively reading events from the partition. These non-exclusive consumers are + sometimes referred to as "Non-Epoch Consumers." """ timeout = 0 @@ -38,8 +38,7 @@ class EventHubConsumer(ConsumerProducerMixin): _timeout = b'com.microsoft:timeout' def __init__( # pylint: disable=super-init-not-called - self, client, source, event_position=None, prefetch=300, owner_level=None, - keep_alive=None, auto_reconnect=True, loop=None): + self, client, source, **kwargs): """ Instantiate an async consumer. EventHubConsumer should be instantiated by calling the `create_consumer` method in EventHubClient. @@ -58,6 +57,13 @@ def __init__( # pylint: disable=super-init-not-called :type owner_level: int :param loop: An event loop. """ + 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) + loop = kwargs.get("loop", None) + super(EventHubConsumer, self).__init__() self.loop = loop or asyncio.get_event_loop() self.running = False @@ -153,7 +159,7 @@ def queue_size(self): return self._handler._received_messages.qsize() return 0 - async def receive(self, max_batch_size=None, timeout=None): + async def receive(self, **kwargs): # type: (int, float) -> List[EventData] """ Receive events asynchronously from the EventHub. @@ -180,6 +186,9 @@ async def receive(self, max_batch_size=None, timeout=None): :caption: Receives events asynchronously """ + max_batch_size = kwargs.get("max_batch_size", None) + timeout = kwargs.get("timeout", None) + self._check_closed() max_batch_size = min(self.client.config.max_batch_size, self.prefetch) if max_batch_size is None else max_batch_size timeout = self.client.config.receive_timeout if timeout is None else timeout @@ -217,7 +226,7 @@ async def receive(self, max_batch_size=None, timeout=None): last_exception = await self._handle_exception(exception, retry_count, max_retries, timeout_time) retry_count += 1 - async def close(self, exception=None): + async def close(self, **kwargs): # type: (Exception) -> None """ Close down the handler. If the handler has already closed, @@ -237,6 +246,7 @@ async def close(self, exception=None): :caption: Close down the handler. """ + exception = kwargs.get("exception", None) self.running = False if self.error: return @@ -250,4 +260,4 @@ async def close(self, exception=None): self.error = EventHubError(str(exception)) else: self.error = EventHubError("This receive handler is now closed.") - await self._handler.close_async() \ No newline at end of file + await self._handler.close_async() 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 16ddb97a36b9..9612b4156327 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/producer_async.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/producer_async.py @@ -23,16 +23,15 @@ class EventHubProducer(ConsumerProducerMixin): """ A producer responsible for transmitting EventData to a specific Event Hub, - grouped together in batches. Depending on the options specified at creation, the producer may - be created to allow event data to be automatically routed to an available partition or specific - to a partition. + grouped together in batches. Depending on the options specified at creation, the producer may + be created to allow event data to be automatically routed to an available partition or specific + to a partition. """ _timeout = b'com.microsoft:timeout' def __init__( # pylint: disable=super-init-not-called - self, client, target, partition=None, send_timeout=60, - keep_alive=None, auto_reconnect=True, loop=None): + self, client, target, **kwargs): """ Instantiate an async EventHubProducer. EventHubProducer should be instantiated by calling the `create_producer` method in EventHubClient. @@ -55,6 +54,12 @@ def __init__( # pylint: disable=super-init-not-called :type auto_reconnect: bool :param loop: An event loop. If not specified the default event loop will be used. """ + partition = kwargs.get("partition", None) + send_timeout = kwargs.get("send_timeout", 60) + keep_alive = kwargs.get("keep_alive", None) + auto_reconnect = kwargs.get("auto_reconnect", True) + loop = kwargs.get("loop", None) + super(EventHubProducer, self).__init__() self.loop = loop or asyncio.get_event_loop() self.running = False @@ -150,7 +155,7 @@ def _on_outcome(self, outcome, condition): self._outcome = outcome self._condition = condition - async def send(self, event_data, partition_key=None, timeout=None): + async def send(self, event_data, **kwargs): # type:(Union[EventData, Iterable[EventData]], Union[str, bytes]) -> None """ Sends an event data and blocks until acknowledgement is @@ -178,6 +183,9 @@ async def send(self, event_data, partition_key=None, timeout=None): :caption: Sends an event data and blocks until acknowledgement is received or operation times out. """ + partition_key = kwargs.get("partition_key", None) + timeout = kwargs.get("timeout", None) + self._check_closed() if isinstance(event_data, EventData): if partition_key: @@ -192,7 +200,7 @@ async def send(self, event_data, partition_key=None, timeout=None): self.unsent_events = [wrapper_event_data.message] await self._send_event_data(timeout) - async def close(self, exception=None): + async def close(self, **kwargs): # type: (Exception) -> None """ Close down the handler. If the handler has already closed, @@ -212,4 +220,5 @@ async def close(self, exception=None): :caption: Close down the handler. """ + exception = kwargs.get("exception", None) await super(EventHubProducer, self).close(exception) diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/client.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/client.py index b7fad40779c9..deda0ddc01fb 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/client.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/client.py @@ -200,8 +200,7 @@ def get_partition_properties(self, partition): return output def create_consumer( - self, consumer_group, partition_id, event_position, - owner_level=None, operation=None, prefetch=None, + self, consumer_group, partition_id, event_position, **kwargs ): # type: (str, str, EventPosition, int, str, int) -> EventHubConsumer """ @@ -233,8 +232,11 @@ def create_consumer( :caption: Add a consumer to the client for a particular consumer group and partition. """ - prefetch = self.config.prefetch if prefetch is None else prefetch + owner_level = kwargs.get("owner_level", None) + operation = kwargs.get("operation", None) + prefetch = kwargs.get("prefetch", None) + prefetch = prefetch or self.config.prefetch 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) @@ -243,7 +245,7 @@ def create_consumer( prefetch=prefetch) return handler - def create_producer(self, partition_id=None, operation=None, send_timeout=None): + def create_producer(self, **kwargs): # type: (str, str, float) -> EventHubProducer """ Create an producer to send EventData object to an EventHub. @@ -269,6 +271,10 @@ def create_producer(self, partition_id=None, operation=None, send_timeout=None): :caption: Add a producer to the client to send EventData. """ + partition_id = kwargs.get("partition_id", None) + operation = kwargs.get("operation", None) + send_timeout = kwargs.get("send_timeout", None) + target = "amqps://{}{}".format(self.address.hostname, self.address.path) if operation: target = target + operation diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/client_abstract.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/client_abstract.py index 38e2afde2615..7f6afb51c7fc 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/client_abstract.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/client_abstract.py @@ -219,7 +219,7 @@ def _process_redirect_uri(self, redirect): self.mgmt_target = redirect_uri @classmethod - def from_connection_string(cls, conn_str, event_hub_path=None, **kwargs): + def from_connection_string(cls, conn_str, **kwargs): """Create an EventHubClient from an EventHub/IotHub connection string. :param conn_str: The connection string of an eventhub or IoT hub @@ -266,6 +266,7 @@ def from_connection_string(cls, conn_str, event_hub_path=None, **kwargs): :caption: Create an EventHubClient from a connection string. """ + event_hub_path = kwargs.get("event_hub_path", None) is_iot_conn_str = conn_str.lstrip().lower().startswith("hostname") if not is_iot_conn_str: address, policy, key, entity = _parse_conn_str(conn_str) @@ -281,12 +282,10 @@ def from_connection_string(cls, conn_str, event_hub_path=None, **kwargs): @abstractmethod def create_consumer( - self, consumer_group, partition_id, event_position, owner_level=None, - operation=None, - prefetch=None, + self, consumer_group, partition_id, event_position, **kwargs ): pass @abstractmethod - def create_producer(self, partition_id=None, operation=None, send_timeout=None): + def create_producer(self, **kwargs): pass diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/common.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/common.py index 5ac6258eeb0a..8faded746a74 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/common.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/common.py @@ -51,7 +51,7 @@ class EventData(object): PROP_TIMESTAMP = b"x-opt-enqueued-time" PROP_DEVICE_ID = b"iothub-connection-device-id" - def __init__(self, body=None, to_device=None, message=None): + def __init__(self, **kwargs): """ Initialize EventData. @@ -64,6 +64,10 @@ def __init__(self, body=None, to_device=None, message=None): :param message: The received message. :type message: ~uamqp.message.Message """ + body = kwargs.get("body", None) + to_device = kwargs.get("to_device", None) + message = kwargs.get("message", None) + self._partition_key = types.AMQPSymbol(EventData.PROP_PARTITION_KEY) self._annotations = {} self._app_properties = {} @@ -206,7 +210,7 @@ def body(self): except TypeError: raise ValueError("Message data empty.") - def body_as_str(self, encoding='UTF-8'): + def body_as_str(self, **kwargs): """ The body of the event data as a string if the data is of a compatible type. @@ -215,6 +219,7 @@ def body_as_str(self, encoding='UTF-8'): Default is 'UTF-8' :rtype: str or unicode """ + encoding = kwargs.get("encoding", 'UTF-8') data = self.body try: return "".join(b.decode(encoding) for b in data) @@ -227,7 +232,7 @@ def body_as_str(self, encoding='UTF-8'): except Exception as e: raise TypeError("Message data is not compatible with string type: {}".format(e)) - def body_as_json(self, encoding='UTF-8'): + def body_as_json(self, **kwargs): """ The body of the event loaded as a JSON object is the data is compatible. @@ -235,6 +240,7 @@ def body_as_json(self, encoding='UTF-8'): Default is 'UTF-8' :rtype: dict """ + encoding = kwargs.get("encoding", 'UTF-8') data_str = self.body_as_str(encoding=encoding) try: return json.loads(data_str) @@ -280,7 +286,7 @@ class EventPosition(object): >>> event_pos = EventPosition(1506968696002) """ - def __init__(self, value, inclusive=False): + def __init__(self, value, **kwargs): """ Initialize EventPosition. @@ -289,6 +295,7 @@ def __init__(self, value, inclusive=False): :param inclusive: Whether to include the supplied value as the start point. :type inclusive: bool """ + inclusive = kwargs.get("inclusive", False) self.value = value if value is not None else "-1" self.inclusive = inclusive diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py index 53368dc292b0..855e6cf97d53 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py @@ -22,23 +22,22 @@ class EventHubConsumer(ConsumerProducerMixin): """ A consumer responsible for reading EventData from a specific Event Hub - partition and as a member of a specific consumer group. + partition and as a member of a specific consumer group. A consumer may be exclusive, which asserts ownership over the partition for the consumer - group to ensure that only one consumer from that group is reading the from the partition. - These exclusive consumers are sometimes referred to as "Epoch Consumers." + group to ensure that only one consumer from that group is reading the from the partition. + These exclusive consumers are sometimes referred to as "Epoch Consumers." A consumer may also be non-exclusive, allowing multiple consumers from the same consumer - group to be actively reading events from the partition. These non-exclusive consumers are - sometimes referred to as "Non-Epoch Consumers." + group to be actively reading events from the partition. These non-exclusive consumers are + sometimes referred to as "Non-Epoch Consumers." """ timeout = 0 _epoch = b'com.microsoft:epoch' _timeout = b'com.microsoft:timeout' - def __init__(self, client, source, event_position=None, prefetch=300, owner_level=None, - keep_alive=None, auto_reconnect=True): + def __init__(self, client, source, **kwargs): """ Instantiate a consumer. EventHubConsumer should be instantiated by calling the `create_consumer` method in EventHubClient. @@ -54,6 +53,12 @@ def __init__(self, client, source, event_position=None, prefetch=300, owner_leve consumer if owner_level is set. :type owner_level: int """ + 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) + super(EventHubConsumer, self).__init__() self.running = False self.client = client @@ -147,7 +152,7 @@ def queue_size(self): return self._handler._received_messages.qsize() return 0 - def receive(self, max_batch_size=None, timeout=None): + def receive(self, **kwargs): # type:(int, float) -> List[EventData] """ Receive events from the EventHub. @@ -173,8 +178,10 @@ def receive(self, max_batch_size=None, timeout=None): :caption: Receive events from the EventHub. """ - self._check_closed() + max_batch_size = kwargs.get("max_batch_size", None) + timeout = kwargs.get("timeout", None) + self._check_closed() max_batch_size = min(self.client.config.max_batch_size, self.prefetch) if max_batch_size is None else max_batch_size timeout = self.client.config.receive_timeout if timeout is None else timeout if not timeout: @@ -210,7 +217,7 @@ def receive(self, max_batch_size=None, timeout=None): last_exception = self._handle_exception(exception, retry_count, max_retries, timeout_time) retry_count += 1 - def close(self, exception=None): + def close(self, **kwargs): # type:(Exception) -> None """ Close down the handler. If the handler has already closed, @@ -230,6 +237,7 @@ def close(self, exception=None): :caption: Close down the handler. """ + exception = kwargs.get("exception", None) if self.messages_iter: self.messages_iter.close() self.messages_iter = None diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/error.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/error.py index 31f456f84eb8..0fb6933e3015 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/error.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/error.py @@ -19,6 +19,7 @@ log = logging.getLogger(__name__) + def _error_handler(error): """ Called internally when an event has failed to send so we @@ -56,7 +57,8 @@ class EventHubError(Exception): :vartype details: dict[str, str] """ - def __init__(self, message, details=None): + def __init__(self, message, **kwargs): + details = kwargs.get("details", None) self.error = None self.message = message self.details = details diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py index d85a965381ad..465fbf45d9ef 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py @@ -35,14 +35,14 @@ def _set_partition_key(event_datas, partition_key): class EventHubProducer(ConsumerProducerMixin): """ A producer responsible for transmitting EventData to a specific Event Hub, - grouped together in batches. Depending on the options specified at creation, the producer may - be created to allow event data to be automatically routed to an available partition or specific - to a partition. + grouped together in batches. Depending on the options specified at creation, the producer may + be created to allow event data to be automatically routed to an available partition or specific + to a partition. """ _timeout = b'com.microsoft:timeout' - def __init__(self, client, target, partition=None, send_timeout=60, keep_alive=None, auto_reconnect=True): + def __init__(self, client, target, **kwargs): """ Instantiate an EventHubProducer. EventHubProducer should be instantiated by calling the `create_producer` method in EventHubClient. @@ -64,6 +64,11 @@ def __init__(self, client, target, partition=None, send_timeout=60, keep_alive=N Default value is `True`. :type auto_reconnect: bool """ + partition = kwargs.get("partition", None) + send_timeout = kwargs.get("send_timeout", 60) + keep_alive = kwargs.get("keep_alive", None) + auto_reconnect = kwargs.get("auto_reconnect", True) + super(EventHubProducer, self).__init__() self.running = False self.client = client @@ -157,7 +162,7 @@ def _on_outcome(self, outcome, condition): self._outcome = outcome self._condition = condition - def send(self, event_data, partition_key=None, timeout=None): + def send(self, event_data, **kwargs): # type:(Union[EventData, Iterable[EventData]], Union[str, bytes], float) -> None """ Sends an event data and blocks until acknowledgement is @@ -186,6 +191,9 @@ def send(self, event_data, partition_key=None, timeout=None): :caption: Sends an event data and blocks until acknowledgement is received or operation times out. """ + partition_key = kwargs.get("partition_key", None) + timeout = kwargs.get("timeout", None) + self._check_closed() if isinstance(event_data, EventData): if partition_key: @@ -200,7 +208,7 @@ def send(self, event_data, partition_key=None, timeout=None): self.unsent_events = [wrapper_event_data.message] self._send_event_data(timeout=timeout) - def close(self, exception=None): + def close(self, **kwargs): # type:(Exception) -> None """ Close down the handler. If the handler has already closed, @@ -220,4 +228,5 @@ def close(self, exception=None): :caption: Close down the handler. """ + exception = kwargs.get("exception", None) super(EventHubProducer, self).close(exception)