-
Notifications
You must be signed in to change notification settings - Fork 3.3k
Eventhubs new EventProcessor #6550
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
Merged
Changes from 17 commits
Commits
Show all changes
47 commits
Select commit
Hold shift + click to select a range
e62825f
Shared connection (sync) draft
84ae397
Shared connection (sync) draft 2
6e6c827
Shared connection (sync) test update
a86146e
Shared connection
c5a23c8
Fix an issue
b83264d
add retry exponential delay and timeout to exception handling
4f17781
put module method before class def
f0d98d1
fixed Client.get_properties error
ed7d414
new eph (draft)
ee228b4
new eph (draft2)
1d65719
remove in memory partition manager
4895385
EventProcessor draft 3
6415ac2
small format change
a685f87
Fix logging
6388f4a
Add EventProcessor example
1cd6275
Merge branch 'eventhubs_preview2' into eventhubs_eph
6394f19
use decorator to implement retry logic and update some tests (#6544)
yunhaoling 56fdd1e
Update livetest (#6547)
yunhaoling 10f0be6
Remove legacy code and update livetest (#6549)
yunhaoling fbb66bd
make sync longrunning multi-threaded
00ff723
small changes on async long running test
0f5180c
reset retry_count for iterator
6e0c238
Don't return early when open a ReceiveClient or SendClient
0efc95f
type annotation change
90fbafb
Update kwargs and remove unused import
yunhaoling e06bad8
Misc changes from EventProcessor PR review
dd1d7ae
raise asyncio.CancelledError out instead of supressing it.
9d18dd9
Merge branch 'eventhubs_dev' into eventhubs_eph
a31ee67
Update livetest and small fixed (#6594)
yunhaoling 19a5539
Merge branch 'eventhubs_dev' into eventhubs_eph
997dacf
Fix feedback from PR (1)
d688090
Revert "Merge branch 'eventhubs_dev' into eventhubs_eph"
2399dcb
Fix feedback from PR (2)
5ad0255
Update code according to the review (#6623)
yunhaoling 5679065
Fix feedback from PR (3)
35245a1
Merge branch 'eventhubs_dev' into eventhubs_eph
83d0ec2
small bug fixing
a2195ba
Remove old EPH
f8acc8b
Update decorator implementation (#6642)
yunhaoling 6fe4533
Remove old EPH pytest
598245d
Revert "Revert "Merge branch 'eventhubs_dev' into eventhubs_eph""
97dfce5
Update sample codes and docstring (#6643)
yunhaoling 64c5c7d
Check tablename to prevent sql injection
13014dd
PR review update
3e82e90
Removed old EPH stuffs.
0f670d8
Small fix (#6650)
yunhaoling 2a8b34c
Merge branch 'eventhubs_dev' into eventhubs_eph
File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
17 changes: 17 additions & 0 deletions
17
sdk/eventhub/azure-eventhubs/azure/eventhub/eventprocessor/__init__.py
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,17 @@ | ||
| # -------------------------------------------------------------------------------------------- | ||
| # Copyright (c) Microsoft Corporation. All rights reserved. | ||
| # Licensed under the MIT License. See License.txt in the project root for license information. | ||
| # ----------------------------------------------------------------------------------- | ||
|
|
||
| from .event_processor import EventProcessor | ||
| from .partition_processor import PartitionProcessor, CloseReason | ||
| from .partition_manager import PartitionManager | ||
| from .sqlite3_partition_manager import Sqlite3PartitionManager | ||
|
|
||
| __all__ = [ | ||
| 'CloseReason', | ||
| 'EventProcessor', | ||
| 'PartitionProcessor', | ||
| 'PartitionManager', | ||
| 'Sqlite3PartitionManager', | ||
| ] | ||
30 changes: 30 additions & 0 deletions
30
sdk/eventhub/azure-eventhubs/azure/eventhub/eventprocessor/checkpoint_manager.py
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,30 @@ | ||
| # -------------------------------------------------------------------------------------------- | ||
| # Copyright (c) Microsoft Corporation. All rights reserved. | ||
| # Licensed under the MIT License. See License.txt in the project root for license information. | ||
| # ----------------------------------------------------------------------------------- | ||
|
|
||
|
|
||
| from .partition_manager import PartitionManager | ||
|
|
||
|
|
||
| class CheckpointManager(object): | ||
| """Every PartitionProcessor has a CheckpointManager to save the partition's checkpoint. | ||
|
|
||
| """ | ||
| def __init__(self, partition_id, eventhub_name, consumer_group_name, instance_id, partition_manager: PartitionManager): | ||
| self._partition_id = partition_id | ||
| self._eventhub_name = eventhub_name | ||
| self._consumer_group_name = consumer_group_name | ||
| self._instance_id = instance_id | ||
| self._partition_manager = partition_manager | ||
|
|
||
| async def update_checkpoint(self, offset, sequence_number): | ||
|
YijunXieMS marked this conversation as resolved.
Outdated
|
||
| """Users call this method in PartitionProcessor.process_events() to save checkpoints | ||
|
|
||
| :param offset: offset of the processed EventData | ||
| :param sequence_number: sequence_number of the processed EventData | ||
| :return: None | ||
| """ | ||
| await self._partition_manager.\ | ||
| update_checkpoint(self._eventhub_name, self._consumer_group_name, self._partition_id, self._instance_id, | ||
| offset, sequence_number) | ||
|
YijunXieMS marked this conversation as resolved.
Outdated
|
||
170 changes: 170 additions & 0 deletions
170
sdk/eventhub/azure-eventhubs/azure/eventhub/eventprocessor/event_processor.py
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,170 @@ | ||
| # -------------------------------------------------------------------------------------------- | ||
| # Copyright (c) Microsoft Corporation. All rights reserved. | ||
| # Licensed under the MIT License. See License.txt in the project root for license information. | ||
| # ----------------------------------------------------------------------------------- | ||
|
|
||
| from typing import Callable, List | ||
| import uuid | ||
| import asyncio | ||
| import logging | ||
|
|
||
| from azure.eventhub import EventPosition, EventHubError | ||
| from azure.eventhub.aio import EventHubClient | ||
| from .checkpoint_manager import CheckpointManager | ||
| from .partition_manager import PartitionManager | ||
| from .partition_processor import PartitionProcessor, CloseReason | ||
| from .utils import get_running_loop | ||
|
|
||
| log = logging.getLogger(__name__) | ||
|
|
||
|
|
||
| class EventProcessor(object): | ||
| def __init__(self, eventhub_client: EventHubClient, consumer_group_name: str, | ||
| partition_processor_factory: Callable[..., PartitionProcessor], | ||
|
YijunXieMS marked this conversation as resolved.
Outdated
|
||
| partition_manager: PartitionManager, **kwargs): | ||
| """An EventProcessor automatically creates and runs consumers for all partitions of the eventhub. | ||
|
|
||
| It provides the user a convenient way to receive events from multiple partitions and save checkpoints. | ||
| If multiple EventProcessors are running for an event hub, they will automatically balance loading. This feature | ||
| won't be availabe until preview 3. | ||
|
|
||
| :param consumer_group_name: the consumer group that is used to receive events | ||
| from the event hub that the eventhub_client is going to receive events from | ||
| :param eventhub_client: an instance of azure.eventhub.aio.EventClient object | ||
| :param partition_processor_callable: a callable that is called to return a PartitionProcessor | ||
| :param partition_manager: an instance of a PartitionManager implementation | ||
| :param initial_event_position: the offset to start a partition consumer if the partition has no checkpoint yet | ||
|
YijunXieMS marked this conversation as resolved.
|
||
| """ | ||
| self._consumer_group_name = consumer_group_name | ||
| self._eventhub_client = eventhub_client | ||
| self._eventhub_name = eventhub_client.eh_name | ||
| self._partition_processor_factory = partition_processor_factory | ||
| self._partition_manager = partition_manager | ||
| self._initial_event_position = kwargs.get("initial_event_position", "-1") | ||
| self._max_batch_size = eventhub_client.config.max_batch_size | ||
| self._receive_timeout = eventhub_client.config.receive_timeout | ||
| self._tasks: List[asyncio.Task] = [] | ||
| self._instance_id = str(uuid.uuid4()) | ||
| self._partition_ids = None | ||
|
YijunXieMS marked this conversation as resolved.
Outdated
|
||
|
|
||
| async def start(self): | ||
| """Start the EventProcessor. | ||
| 1. retrieve the partition ids from eventhubs | ||
| 2. claim partition ownership of these partitions. | ||
| 3. repeatedly call EvenHubConsumer.receive() to retrieve events and | ||
| call user defined PartitionProcessor.process_events() | ||
| """ | ||
| log.info("EventProcessor %r is being started", self._instance_id) | ||
| partition_ids = await self._eventhub_client.get_partition_ids() | ||
| self.partition_ids = partition_ids | ||
|
YijunXieMS marked this conversation as resolved.
Outdated
|
||
| claimed_list = await self._claim_partitions() | ||
| await self._start_claimed_partitions(claimed_list) | ||
|
|
||
| async def stop(self): | ||
| """Stop all the partition consumer | ||
|
|
||
| It sends out a cancellation token to stop all partitions' EventHubConsumer will stop receiving events. | ||
|
YijunXieMS marked this conversation as resolved.
Outdated
|
||
|
|
||
| """ | ||
| for task in self._tasks: | ||
| task.cancel() | ||
| # It's not agreed whether a partition manager has method close(). | ||
| log.info("EventProcessor %r has been cancelled", self._instance_id) | ||
|
|
||
| async def _claim_partitions(self): | ||
|
YijunXieMS marked this conversation as resolved.
Outdated
|
||
| partitions_ownership = await self._partition_manager.list_ownership(self._eventhub_name, self._consumer_group_name) | ||
| partitions_ownership_dict = dict() | ||
| for ownership in partitions_ownership: | ||
| partitions_ownership_dict[ownership["partition_id"]] = ownership | ||
|
|
||
| to_claim_list = [] | ||
| for pid in self.partition_ids: | ||
| p_ownership = partitions_ownership_dict.get(pid) | ||
| if p_ownership: | ||
| to_claim_list.append(p_ownership) | ||
|
YijunXieMS marked this conversation as resolved.
|
||
| else: | ||
| new_ownership = dict() | ||
| new_ownership["eventhub_name"] = self._eventhub_name | ||
| new_ownership["consumer_group_name"] = self._consumer_group_name | ||
| new_ownership["instance_id"] = self._instance_id | ||
| new_ownership["partition_id"] = pid | ||
| new_ownership["owner_level"] = 1 # will increment in preview 3 | ||
|
YijunXieMS marked this conversation as resolved.
Outdated
|
||
| to_claim_list.append(new_ownership) | ||
| claimed_list = await self._partition_manager.claim_ownership(to_claim_list) | ||
| return claimed_list | ||
|
|
||
| async def _start_claimed_partitions(self, claimed_partitions): | ||
| consumers = [] | ||
| for partition in claimed_partitions: | ||
| partition_id = partition["partition_id"] | ||
| offset = partition.get("offset") | ||
| offset = offset or self._initial_event_position | ||
| consumer = self._eventhub_client.create_consumer(self._consumer_group_name, partition_id, | ||
|
YijunXieMS marked this conversation as resolved.
|
||
| EventPosition(str(offset))) | ||
|
YijunXieMS marked this conversation as resolved.
YijunXieMS marked this conversation as resolved.
|
||
| consumers.append(consumer) | ||
|
YijunXieMS marked this conversation as resolved.
Outdated
|
||
|
|
||
| partition_processor = self._partition_processor_factory( | ||
| eventhub_name=self._eventhub_name, | ||
| consumer_group_name=self._consumer_group_name, | ||
| partition_id=partition_id, | ||
| checkpoint_manager=CheckpointManager(partition_id, self._eventhub_name, self._consumer_group_name, | ||
| self._instance_id, self._partition_manager) | ||
| ) | ||
| loop = get_running_loop() | ||
| task = loop.create_task( | ||
| _receive(consumer, partition_processor, self._receive_timeout)) | ||
| self._tasks.append(task) | ||
|
YijunXieMS marked this conversation as resolved.
|
||
|
|
||
| await asyncio.gather(*self._tasks) | ||
| await self._partition_manager.close() | ||
| log.info("EventProcessor %r partition manager is closed", self._instance_id) | ||
| log.info("EventProcessor %r has stopped", self._instance_id) | ||
|
|
||
|
|
||
| async def _receive(partition_consumer, partition_processor, receive_timeout): | ||
| try: | ||
| while True: | ||
| try: | ||
| events = await partition_consumer.receive(timeout=receive_timeout) | ||
| except asyncio.CancelledError: | ||
| await partition_processor.close(reason=CloseReason.SHUTDOWN) | ||
| log.info( | ||
| "PartitionProcessor of EventProcessor instance %r of eventhub %r partition %r consumer group %r " | ||
| "has been shutdown", | ||
| partition_processor._checkpoint_manager._instance_id, | ||
| partition_processor._eventhub_name, | ||
| partition_processor._partition_id, | ||
| partition_processor._consumer_group_name | ||
| ) | ||
| break | ||
| except EventHubError as eh_err: | ||
| reason = CloseReason.LEASE_LOST if eh_err.error == "link:stolen" else CloseReason.EVENTHUB_EXCEPTION | ||
| log.info( | ||
|
YijunXieMS marked this conversation as resolved.
Outdated
|
||
| "PartitionProcessor of EventProcessor instance %r of eventhub %r partition %r consumer group %r " | ||
| "has met an exception receiving events. It's being closed. The exception is %r.", | ||
| partition_processor._checkpoint_manager._instance_id, | ||
| partition_processor._eventhub_name, | ||
| partition_processor._partition_id, | ||
| partition_processor._consumer_group_name, | ||
| eh_err | ||
| ) | ||
| await partition_processor.process_error(eh_err) | ||
| await partition_processor.close(reason=reason) | ||
| break | ||
| try: | ||
| await partition_processor.process_events(events) | ||
|
YijunXieMS marked this conversation as resolved.
|
||
| except Exception as exp: # user code has caused an error | ||
| log.info( | ||
| "PartitionProcessor of EventProcessor instance %r of eventhub %r partition %r consumer group %r " | ||
| "has met an exception from user code process_events. It's being closed. The exception is %r.", | ||
| partition_processor._checkpoint_manager._instance_id, | ||
| partition_processor._eventhub_name, | ||
| partition_processor._partition_id, | ||
| partition_processor._consumer_group_name, | ||
| exp | ||
| ) | ||
| await partition_processor.process_error(exp) | ||
| # TODO: will review whether to break and close partition processor after user's code has an exception | ||
| # TODO: try to inform other EventProcessors to take the partition when this partition is closed in preview 3? | ||
| finally: | ||
| await partition_consumer.close() | ||
44 changes: 44 additions & 0 deletions
44
sdk/eventhub/azure-eventhubs/azure/eventhub/eventprocessor/partition_manager.py
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,44 @@ | ||
| # -------------------------------------------------------------------------------------------- | ||
| # Copyright (c) Microsoft Corporation. All rights reserved. | ||
| # Licensed under the MIT License. See License.txt in the project root for license information. | ||
| # ----------------------------------------------------------------------------------- | ||
|
|
||
| from typing import Iterable, Dict, Any | ||
| from abc import ABC, abstractmethod | ||
|
|
||
|
|
||
| class PartitionManager(ABC): | ||
| """Subclass PartitionManager to implement the read/write access to storage service to list/claim ownership and save checkpoint. | ||
|
|
||
| """ | ||
|
|
||
| @abstractmethod | ||
| async def list_ownership(self, eventhub_name: str, consumer_group_name: str) -> Iterable[Dict[str, Any]]: | ||
|
YijunXieMS marked this conversation as resolved.
|
||
| """ | ||
|
|
||
| :param eventhub_name: | ||
| :param consumer_group_name: | ||
| :return: Iterable of dictionaries containing the following partition ownership information: | ||
| eventhub_name | ||
| consumer_group_name | ||
| instance_id | ||
| partition_id | ||
| owner_level | ||
| offset | ||
| sequence_number | ||
| last_modified_time | ||
| etag | ||
| """ | ||
| pass | ||
|
|
||
| @abstractmethod | ||
| async def claim_ownership(self, partitions: Iterable[Dict[str, Any]]) -> Iterable[Dict[str, Any]]: | ||
| pass | ||
|
|
||
| @abstractmethod | ||
| async def update_checkpoint(self, eventhub_name, consumer_group_name, partition_id, instance_id, | ||
| offset, sequence_number) -> None: | ||
| pass | ||
|
|
||
| async def close(self): | ||
| pass | ||
47 changes: 47 additions & 0 deletions
47
sdk/eventhub/azure-eventhubs/azure/eventhub/eventprocessor/partition_processor.py
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,47 @@ | ||
| # -------------------------------------------------------------------------------------------- | ||
| # Copyright (c) Microsoft Corporation. All rights reserved. | ||
| # Licensed under the MIT License. See License.txt in the project root for license information. | ||
| # ----------------------------------------------------------------------------------- | ||
|
|
||
| from typing import List | ||
| from abc import ABC, abstractmethod | ||
| from enum import Enum | ||
| from .checkpoint_manager import CheckpointManager | ||
|
|
||
| from azure.eventhub import EventData | ||
|
|
||
|
|
||
| class CloseReason(Enum): | ||
| SHUTDOWN = 0 # user call EventProcessor.stop() | ||
| OWNERSHIP_LOST = 1 # lose the ownership of a partition. | ||
| EVENTHUB_EXCEPTION = 2 # Exception happens during receiving events | ||
|
|
||
|
|
||
| class PartitionProcessor(ABC): | ||
| def __init__(self, eventhub_name, consumer_group_name, partition_id, checkpoint_manager: CheckpointManager): | ||
| self._partition_id = partition_id | ||
|
YijunXieMS marked this conversation as resolved.
Outdated
|
||
| self._eventhub_name = eventhub_name | ||
| self._consumer_group_name = consumer_group_name | ||
| self._checkpoint_manager = checkpoint_manager | ||
|
|
||
| async def close(self, reason): | ||
| """Called when EventProcessor stops processing this PartitionProcessor. | ||
|
|
||
| There are different reasons to trigger the PartitionProcessor to close. | ||
| Refer to enum class CloseReason | ||
|
|
||
| """ | ||
| pass | ||
|
|
||
| @abstractmethod | ||
| async def process_events(self, events: List[EventData]): | ||
|
YijunXieMS marked this conversation as resolved.
|
||
| """Called when a batch of events have been received. | ||
|
|
||
| """ | ||
| pass | ||
|
|
||
| async def process_error(self, error): | ||
| """Called when an error happens | ||
|
|
||
| """ | ||
| pass | ||
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.