Skip to content
Merged
Show file tree
Hide file tree
Changes from 13 commits
Commits
Show all changes
136 commits
Select commit Hold shift + click to select a range
a6875dd
add implementation
Jul 21, 2021
6b6fd00
fix pylint
Jul 22, 2021
515f886
fix pylint
Jul 22, 2021
773020a
fix pylint
Jul 22, 2021
8b241bd
fix pylint
Jul 22, 2021
7a8a62d
fix pylint
Jul 22, 2021
63fb4cb
fix pylint
Jul 22, 2021
ffa5541
fix pylint
Jul 22, 2021
2986370
fix pylint
Jul 22, 2021
1f06d3c
fix pylint
Jul 22, 2021
68e6e05
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Jul 23, 2021
be8362a
new update
Jul 23, 2021
29c7a84
Merge branch 'josue_third_branch' of https://github.com/Jg1255/azure-…
Jul 23, 2021
4fc4b19
Update shared_requirements.txt
Jg1255 Jul 23, 2021
9bca97a
environment
Jul 26, 2021
0159f6d
redo
Jul 26, 2021
eb6d52f
redo
Jul 26, 2021
fc36f66
redo
Jul 26, 2021
d698e8d
changed variable name
Jul 26, 2021
475b7a8
changed variable name
Jul 26, 2021
e5d68a6
Update sdk/eventhub/azure-eventhub-checkpointstoretable/tests/test_st…
Jg1255 Jul 26, 2021
9ff942c
update on test file
Jul 26, 2021
319bb77
update on test file
Jul 26, 2021
d585ec8
update on test file
Jul 26, 2021
e732e86
update on test file
Jul 26, 2021
22c6f2c
update based on feedack
Jul 28, 2021
9a24922
new update
Jul 28, 2021
233a993
update
Jul 28, 2021
1d82efb
update
Jul 28, 2021
5645319
update
Jul 28, 2021
a44af13
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Jul 29, 2021
eb8a9b5
new update
Aug 4, 2021
77d116f
Merge branch 'josue_third_branch' of https://github.com/Jg1255/azure-…
Aug 4, 2021
7bfb04b
new
Aug 4, 2021
786a700
new
Aug 4, 2021
252dab6
new
Aug 4, 2021
877f1d8
new
Aug 4, 2021
bc2e430
new
Aug 4, 2021
a36a1c1
update on spacing
Aug 4, 2021
a2593d4
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 5, 2021
76ccc65
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 5, 2021
66a3670
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 5, 2021
dabe61a
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 5, 2021
54a913d
new update
Aug 5, 2021
44bc8fe
update
Aug 5, 2021
ce4ea29
update on test file
Aug 6, 2021
9b7f144
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 6, 2021
0cec994
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 7, 2021
cbd01b6
update
Aug 7, 2021
148312a
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 7, 2021
037d62d
update
Aug 7, 2021
6e56a1e
Merge branch 'josue_third_branch' of https://github.com/Jg1255/azure-…
Aug 7, 2021
3d5c888
update
Aug 7, 2021
76a2d6c
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 7, 2021
7a01f35
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 7, 2021
061250c
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 7, 2021
bdda591
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 7, 2021
dca6422
Update sdk/eventhub/azure-eventhub-checkpointstoretable/tests/test_st…
Jg1255 Aug 7, 2021
ac89465
Update sdk/eventhub/azure-eventhub-checkpointstoretable/tests/test_st…
Jg1255 Aug 7, 2021
a91665e
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 7, 2021
3d6798e
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 7, 2021
5f496d3
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 7, 2021
bafce75
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 7, 2021
233d50f
update
Aug 7, 2021
578d18b
update
Aug 7, 2021
d4dd3f1
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 7, 2021
2c89a54
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 7, 2021
3bd50c7
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 7, 2021
790896b
update
Aug 7, 2021
65a9ddf
Update sdk/eventhub/azure-eventhub-checkpointstoretable/tests/test_st…
Jg1255 Aug 7, 2021
e997bc8
update
Aug 7, 2021
0bd82a9
Update sdk/eventhub/azure-eventhub-checkpointstoretable/tests/test_st…
Jg1255 Aug 7, 2021
2614ce3
update
Aug 7, 2021
676b235
update
Aug 9, 2021
a6aa61e
update
Aug 9, 2021
edbe697
update
Aug 9, 2021
5af3e05
update
Aug 9, 2021
d2dbb2a
update
Aug 9, 2021
610606e
Revert "update"
Aug 9, 2021
eee983f
update
Aug 9, 2021
6bc5e7e
update
Aug 9, 2021
affd864
update
Aug 9, 2021
aef09bc
newupdate
Aug 9, 2021
fc7fbcb
update
Aug 9, 2021
d67cabf
update
Aug 9, 2021
f97b399
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 10, 2021
695ed0b
update
Aug 10, 2021
aeeb27e
update
Aug 10, 2021
13f3ed3
update
Aug 10, 2021
8fd90ec
update
Aug 10, 2021
8a3951d
update
Aug 10, 2021
e991fed
update
Aug 10, 2021
7fb1c62
update
Aug 11, 2021
19cbf7e
update
Aug 11, 2021
0e8d95a
update
Aug 11, 2021
d6b450e
update
Aug 11, 2021
b3f333a
update
Aug 11, 2021
f3c3578
update
Aug 11, 2021
d2f9390
update
Aug 11, 2021
b148f1c
update
Aug 11, 2021
bc3b01d
update
Aug 11, 2021
5e2333e
update
Aug 11, 2021
f9af616
update
Aug 11, 2021
23e9546
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 12, 2021
ef11a52
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 12, 2021
bcb9dcc
Update sdk/eventhub/azure-eventhub-checkpointstoretable/samples/recei…
Jg1255 Aug 12, 2021
e2c9ae4
Update sdk/eventhub/azure-eventhub-checkpointstoretable/setup.py
Jg1255 Aug 12, 2021
a2bdb0f
Update sdk/eventhub/azure-eventhub-checkpointstoretable/tests/test_st…
Jg1255 Aug 12, 2021
00f1479
Update sdk/eventhub/azure-eventhub-checkpointstoretable/tests/test_st…
Jg1255 Aug 12, 2021
d02db94
update
Aug 12, 2021
0b9d28a
update
Aug 12, 2021
6712ea1
update
Aug 13, 2021
bcb95e2
update
Aug 13, 2021
879dae7
update
Aug 13, 2021
222a6c5
update
Aug 13, 2021
7353298
p
Aug 13, 2021
7c508e4
update
Aug 13, 2021
dbc12da
update
Aug 13, 2021
f93a037
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 13, 2021
a2b871c
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 13, 2021
d837f5b
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 13, 2021
43bbc11
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 13, 2021
141fbbd
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 13, 2021
dca1f0d
Update sdk/eventhub/azure-eventhub-checkpointstoretable/tests/test_st…
Jg1255 Aug 13, 2021
74d73e1
Update sdk/eventhub/azure-eventhub-checkpointstoretable/tests/test_st…
Jg1255 Aug 13, 2021
de08199
Update sdk/eventhub/azure-eventhub-checkpointstoretable/tests/test_st…
Jg1255 Aug 13, 2021
6553356
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 13, 2021
a6e9481
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 13, 2021
fb8fc57
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 13, 2021
96af725
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 13, 2021
de6ce7e
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 13, 2021
6be0ce9
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 13, 2021
0bce5e8
Update sdk/eventhub/azure-eventhub-checkpointstoretable/azure/eventhu…
Jg1255 Aug 13, 2021
b156941
update
Aug 13, 2021
a545417
update
Aug 13, 2021
003f741
update
Aug 13, 2021
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 @@ -2,6 +2,11 @@
# Copyright (c) Microsoft Corporation. All rights reserved.
# Licensed under the MIT License. See License.txt in the project root for license information.
# --------------------------------------------------------------------------------------------
from collections import defaultdict
from azure.data.tables import TableClient, UpdateMode
from azure.core import MatchConditions
from azure.data.tables._base_client import parse_connection_str
from azure.core.exceptions import ResourceModifiedError, ResourceExistsError, ResourceNotFoundError

class TableCheckpointStore:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think this class should be extending the base CheckpointStore class. You can look at the storage-blob implementation to see how that's being done there.

"""A CheckpointStore that uses Azure Table Storage to store the partition ownership and checkpoint data.
Expand All @@ -22,17 +27,246 @@ class TableCheckpointStore:
The hostname of the secondary endpoint.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is the secondary hostname actually able to be used, or can we remove it from the docs?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I removed it

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I can't comment on the api_version line directly, but it says the default value is '2019-07-07'.

However it looks like the default value is actually '2018-03-28' if we look at the tables package:

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

fixed it

"""

def __init__(self, **kwargs):
pass
def __init__(self, table_account_url, table_name, credential=None, **kwargs):
# type(str, str, Optional[Any], Any) -> None
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
self.table_client = kwargs.pop("table_client", None)
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
if not self.table_client:
api_version = kwargs.pop("api_version", None)
if api_version:
headers = kwargs.get("headers")
if headers:
headers["x-ms-version"] = api_version
else:
kwargs["headers"] = {"x-ms-version": api_version}
self.table_client = TableClient(
table_account_url, table_name, credential=credential, **kwargs
)
Comment thread
swathipil marked this conversation as resolved.
self._cached_table_clients = defaultdict() # type: Dict[str, TableClient]
Comment thread
Jg1255 marked this conversation as resolved.
Outdated

def list_ownership(self, namespace, eventhub, consumergroup, **kwargs):
pass
@classmethod
def from_connection_string(cls, conn_str, table_name, credential=None, **kwargs):
Comment thread
swathipil marked this conversation as resolved.
"""Create TableCheckpointStore from a storage connection string.
Comment thread
Jg1255 marked this conversation as resolved.
:param str conn_str:
A connection string to an Azure Storage account.
:param table_name:
The table name.
:type table_name: str
:param credential:
The credentials with which to authenticate. This is optional if the
account URL already has a SAS token, or the connection string already has shared
access key values. The value can be a SAS token string, an account shared access
key, or an instance of a TokenCredentials class from azure.identity.
Credentials provided here will take precedence over those in the connection string.
:keyword str api_version:
The Storage API version to use for requests. Default value is '2019-07-07'.
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
:keyword str secondary_hostname:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It doesn't look like the TableClient kwargs describe a secondary_hostname, so I don't think we need to support it here either.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This comment was marked as resolved but I'm not sure why. Do we need secondary_hostname or can it be removed from the docs?

It also looks like it is being accepted in the constructor as well, but I don't see how it is usable today.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes we do not need secondary_hostname

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sweet, please remove then 😄

:returns: A table checkpoint store.
:rtype: ~azure.eventhub.extensions.checkpointstoretable.TableCheckpointStore
"""
endpoint, credential = parse_connection_str(
conn_str=conn_str, credential=None, keyword_args=kwargs
)
return cls(endpoint, table_name=table_name, credential=credential, **kwargs)

Comment thread
swathipil marked this conversation as resolved.
def _create_entity_checkpoint(self, checkpoint):
my_new_entity = {
u'PartitionKey': u'',
u'RowKey': u'',
u'consumer_group': checkpoint['consumer_group'],
u'fully_qualified_namespace': checkpoint['fully_qualified_namespace'],
u'eventhub_name': checkpoint['eventhub_name'],
u'partition_id': checkpoint['partition_id'],
u'offset' : checkpoint['offset'],
u'sequence_number' : checkpoint['sequence_number'],
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
}
my_new_entity['RowKey'] = my_new_entity['partition_id']
my_new_entity['PartitionKey'] = my_new_entity['eventhub_name'] + ' ' + \
my_new_entity['fully_qualified_namespace'] + ' ' + my_new_entity['consumer_group'] + ' ' + 'Checkpoint'
self.table_client.create_entity(entity=my_new_entity)

def _create_entity_ownership(self, ownership):
my_new_entity = {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: I think a better name for this would be ownership_entity. Then it's very clear just looking at the variable name what it is meant to hold.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

You marked this as resolved, but I'm not sure why. Do you disagree with the name change?

u'PartitionKey': u'',
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
u'RowKey': u'',
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
u'consumer_group': ownership['consumer_group'],
u'fully_qualified_namespace': ownership['fully_qualified_namespace'],
u'eventhub_name': ownership['eventhub_name'],
u'partition_id': ownership['partition_id'],
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
u'owner_id' : ownership['owner_id'],
}
my_new_entity['RowKey'] = my_new_entity['partition_id']
my_new_entity['PartitionKey'] = my_new_entity['eventhub_name'] + ' ' \
+ my_new_entity['fully_qualified_namespace'] + ' ' + my_new_entity['consumer_group'] + ' ' + 'Ownership'
new_entity = self.table_client.create_entity(entity=my_new_entity)
return new_entity
Comment thread
Jg1255 marked this conversation as resolved.
Outdated

@classmethod
def _modify_entity_ownership(cls, ownership):
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
"""
create a dictionary with the new ownership attributes so that it can be updated in tables
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
"""
my_new_entity = {
u'PartitionKey': u'',
u'RowKey': u'',
u'consumer_group': ownership['consumer_group'],
u'fully_qualified_namespace': ownership['fully_qualified_namespace'],
u'eventhub_name': ownership['eventhub_name'],
u'partition_id': ownership['partition_id'],
u'owner_id' : ownership['owner_id'],
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
}
my_new_entity['RowKey'] = my_new_entity['partition_id']
my_new_entity['PartitionKey'] = my_new_entity['eventhub_name'] + ' ' \
+ my_new_entity['fully_qualified_namespace'] + ' ' + my_new_entity['consumer_group'] + ' ' + 'Ownership'
return my_new_entity

@classmethod
def _modify_entity_checkpoint(cls, checkpoint):
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
"""
create a dictionary with the new checkpoint attributes so that it can be updated in tables
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
"""
my_new_entity = {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: Since you're creating a checkpoint entity, renaming this to checkpoint_entity is a bit more clear.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

yea I will replace checkpoint_entity that make sense

u'PartitionKey': u'',
u'RowKey': u'',
u'consumer_group': checkpoint['consumer_group'],
u'fully_qualified_namespace': checkpoint['fully_qualified_namespace'],
u'eventhub_name': checkpoint['eventhub_name'],
u'partition_id': checkpoint['partition_id'],
u'offset' : checkpoint['offset'],
u'sequence_number' : checkpoint['sequence_number'],}
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
my_new_entity['RowKey'] = my_new_entity['partition_id']
my_new_entity['PartitionKey'] = my_new_entity['eventhub_name'] + ' ' \
+ my_new_entity['fully_qualified_namespace'] + ' ' + my_new_entity['consumer_group'] + ' ' + 'Checkpoint'
return my_new_entity

def list_ownership(self, fully_qualified_namespace, eventhub_name, consumer_group):
"""Retrieves a complete ownership list from the storage table.
Comment thread
Jg1255 marked this conversation as resolved.
:param str fully_qualified_namespace: The fully qualified namespace that the Event Hub belongs to.
The format is like "<namespace>.servicebus.windows.net".
:param str eventhub_name: The name of the specific Event Hub the partition ownerships are associated with,
relative to the Event Hubs namespace that contains it.
:param str consumer_group: The name of the consumer group the ownerships are associated with.
:rtype: Iterable[Dict[str, Any]], Iterable of dictionaries containing partition ownership information:
- `fully_qualified_namespace` (str): The fully qualified namespace that the Event Hub belongs to.
The format is like "<namespace>.servicebus.windows.net".
- `eventhub_name` (str): The name of the specific Event Hub the checkpoint is associated with,
relative to the Event Hubs namespace that contains it.
- `consumer_group` (str): The name of the consumer group the ownership are associated with.
- `partition_id` (str): The partition ID which the checkpoint is created for.
- `owner_id` (str): A UUID representing the current owner of this partition.
- `last_modified_time` (UTC datetime.datetime): The last time this ownership was claimed.
- `etag` (str): The Etag value for the last time this ownership was modified. Optional depending
on storage implementation.
"""
Comment thread
swathipil marked this conversation as resolved.
thePartitionKey = eventhub_name + ' ' + fully_qualified_namespace +' ' + consumer_group + ' '+ 'Ownership'
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
my_filter = "PartitionKey eq '"+ thePartitionKey + "'"
entities = self.table_client.query_entities(my_filter)
ownershiplist = []
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
for entity in entities:
dic = {}
dic[u'fully_qualified_namespace'] = entity[u'fully_qualified_namespace']
dic[u'eventhub_name'] = entity[u'eventhub_name']
dic[u'consumer_group'] = entity[u'consumer_group']
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
dic[u'partition_id'] = entity[u'partition_id']
dic[u'owner_id'] = entity[u'owner_id']
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
dic[u'etag'] = entity.metadata.get('etag')
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
ownershiplist.append(dic)
return ownershiplist

def list_checkpoints(self, fully_qualified_namespace, eventhub_name, consumer_group):
"""List the updated checkpoints from the storage blob.
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
:param str fully_qualified_namespace: The fully qualified namespace that the Event Hub belongs to.
The format is like "<namespace>.servicebus.windows.net".
:param str eventhub_name: The name of the specific Event Hub the checkpoints are associated with, relative to
the Event Hubs namespace that contains it.
:param str consumer_group: The name of the consumer group the checkpoints are associated with.
:rtype: Iterable[Dict[str,Any]], Iterable of dictionaries containing partition checkpoint information:
- `fully_qualified_namespace` (str): The fully qualified namespace that the Event Hub belongs to.
The format is like "<namespace>.servicebus.windows.net".
- `eventhub_name` (str): The name of the specific Event Hub the checkpoints are associated with,
relative to the Event Hubs namespace that contains it.
- `consumer_group` (str): The name of the consumer group the checkpoints are associated with.
- `partition_id` (str): The partition ID which the checkpoint is created for.
- `sequence_number` (int): The sequence number of the :class:`EventData<azure.eventhub.EventData>`.
- `offset` (str): The offset of the :class:`EventData<azure.eventhub.EventData>`.
"""
thePartitionKey = eventhub_name + ' ' + fully_qualified_namespace +' ' + consumer_group + ' '+ 'Checkpoint'
my_filter = "PartitionKey eq '"+ thePartitionKey + "'"
entities = self.table_client.query_entities(my_filter)
checkpointslist = []
for entity in entities:
dic = {}
Comment thread
chradek marked this conversation as resolved.
Outdated
dic[u'fully_qualified_namespace'] = entity[u'fully_qualified_namespace']
dic[u'eventhub_name'] = entity[u'eventhub_name']
dic[u'consumer_group'] = entity[u'consumer_group']
dic[u'partition_id'] = entity[u'partition_id']
dic[u'sequence_number'] = entity[u'sequence_number']
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
dic[u'offset'] = entity[u'offset']
checkpointslist.append(dic)
return checkpointslist

def update_checkpoint(self, checkpoint):
"""Updates the checkpoint using the given information for the offset, associated partition and
consumer group in the storage table.
Comment thread
Jg1255 marked this conversation as resolved.
Note: If you plan to implement a custom checkpoint store with the intention of running between
cross-language EventHubs SDKs, it is recommended to persist the offset value as an integer.
Comment thread
Jg1255 marked this conversation as resolved.
:param Dict[str,Any] checkpoint: A dict containing checkpoint information:
- `fully_qualified_namespace` (str): The fully qualified namespace that the Event Hub belongs to.
The format is like "<namespace>.servicebus.windows.net".
- `eventhub_name` (str): The name of the specific Event Hub the checkpoint is associated with,
relative to the Event Hubs namespace that contains it.
- `consumer_group` (str): The name of the consumer group the checkpoint is associated with.
- `partition_id` (str): The partition ID which the checkpoint is created for.
- `sequence_number` (int): The sequence number of the :class:`EventData<azure.eventhub.EventData>`
the new checkpoint will be associated with.
- `offset` (str): The offset of the :class:`EventData<azure.eventhub.EventData>`
the new checkpoint will be associated with.
:rtype: None
"""
try:
theentity = self._modify_entity_checkpoint(checkpoint)
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
self.table_client.update_entity(mode=UpdateMode.REPLACE, entity=theentity)
except ResourceNotFoundError:
Comment thread
swathipil marked this conversation as resolved.
self._create_entity_checkpoint(checkpoint)

def _upload_ownership(self, ownership):
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
try:
theentity = self._modify_entity_ownership(ownership)
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
entity = self.table_client.update_entity(mode=UpdateMode.REPLACE, entity=theentity,
etag=ownership['etag'], match_condition=MatchConditions.IfNotModified)
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
ownership['etag'] = entity['etag']
ownership['last_modified_time'] = entity['date']
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
except (ResourceNotFoundError, ValueError):
try:
entity = self._create_entity_ownership(ownership)
ownership['etag'] = entity['etag']
ownership['last_modified_time'] = entity['date']
except ResourceExistsError:
raise'your etag does not match the last time this ownership was modified.'
Comment thread
Jg1255 marked this conversation as resolved.
Outdated
except ResourceModifiedError:
raise'your etag does not match the last time this ownership was modified.'

def list_checkpoints(self, namespace, eventhub, consumergroup, **kwargs):
pass

def update_checkpoint(self, checkpoint, **kwargs):
pass
def _claim_one_partition(self, ownership):
self._upload_ownership(ownership)
return ownership

def claim_ownership(self, ownershiplist, **kwargs):
pass
def claim_ownership(self, ownershiplist):
# type: (Iterable[Dict[str, Any]], Any) -> Iterable[Dict[str, Any]]
"""Tries to claim ownership for a list of specified partitions.
Comment thread
swathipil marked this conversation as resolved.
:param Iterable[Dict[str,Any]] ownership_list: Iterable of dictionaries containing all the ownerships to claim.
:rtype: Iterable[Dict[str,Any]], Iterable of dictionaries containing partition ownership information:
- `fully_qualified_namespace` (str): The fully qualified namespace that the Event Hub belongs to.
The format is like "<namespace>.servicebus.windows.net".
- `eventhub_name` (str): The name of the specific Event Hub the checkpoint is associated with,
relative to the Event Hubs namespace that contains it.
- `consumer_group` (str): The name of the consumer group the ownership are associated with.
- `partition_id` (str): The partition ID which the checkpoint is created for.
- `owner_id` (str): A UUID representing the owner attempting to claim this partition.
- `last_modified_time` (UTC datetime.datetime): The last time this ownership was claimed.
- `etag` (str): The Etag value for the last time this ownership was modified. Optional depending
on storage implementation.
"""
newlist = []

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

snake case + consistency with blob implementation:

Suggested change
newlist = []
gathered_results = []

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

ok

for x in ownershiplist:
newlist.append(self._claim_one_partition(x))
return newlist
Comment thread
chradek marked this conversation as resolved.
Outdated
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
-e ../../../tools/azure-sdk-tools
../../core/azure-core
../../tables/azure-data-tables

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If we copy the the tables code into a _vendor folder now, we can remove this line.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@Jg1255 I think you can remove this line now

-e ../../../tools/azure-devtools
2 changes: 1 addition & 1 deletion sdk/eventhub/azure-eventhub-checkpointstoretable/setup.py
Original file line number Diff line number Diff line change
Expand Up @@ -64,7 +64,7 @@
]
),
install_requires=[
"azure-core<2.0.0,>=1.2.2",

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

can we also add the dependencies required for vendored tables?

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

and also the tables dependencies to extras_require below?

":python_version<'3.0'": ['futures', 'azure-data-nspkg<2.0.0,>=1.0.0'],

(should be these 3:
":python_version<'3.0'": ['futures', 'azure-data-nspkg<2.0.0,>=1.0.0'],
":python_version<'3.4'": ['enum34>=1.0.4'],
":python_version<'3.5'": ["typing"]
)

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

ok

"azure-core<2.0.0,>=1.14.0",
],
extras_require={
":python_version<'3.0'": ["azure-nspkg"],
Expand Down
Loading