# to run the script: # install all required dependencies; # make sure you have needed releases in FEED_BUCKET; # execute python release_files_parser.py; # Signoff.xlsx file will appear in the same directory. import os from typing import Any, Dict, Generator, List import boto3 import smart_open import xlsxwriter from lxml.etree import ElementTree, parse # usually ddex team asks us to parse a few releases for testing # if you need to parse specific releases - add them to RELEASES list below # otherwise all release from the FEED_BUCKET will be processed RELEASES: List[str] = [ # 'A10301A00042689714_20211007011408166', # 'A10301A0003729265G_20211007011324124', # 'A10301A00035091829_20211007011327069', ] class ReleaseParser: FEED_BUCKET = 'dev-ddex-feed' COLUMNS = [ 'Party ID', 'MessageCreatedDateTime', 'Catalog Number', 'ICPN (UPC)', 'GRID', 'ISRC', 'Artist', 'Title', 'ReleaseType', 'CommercialType', 'UseType', 'Territories', 'Price Value/Code', 'SalesStartDate', 'SalesEndDate', ] def __init__(self): self._s3_client = boto3.client('s3') def _get_folders(self, **base_kwargs: str) -> Generator: continuation_token = None while True: list_kwargs = dict(MaxKeys=1000, **base_kwargs) if continuation_token: list_kwargs['ContinuationToken'] = continuation_token response = self._s3_client.list_objects_v2(**list_kwargs) yield from response.get('CommonPrefixes', []) if not response.get('IsTruncated'): break continuation_token = response.get('NextContinuationToken') def _get_release_list(self) -> List[str]: if RELEASES: return RELEASES result = self._get_folders(Bucket=self.FEED_BUCKET, Delimiter='/') folders = [res.get('Prefix')[:-1] for res in result] folders.remove('acknowledgements') folders.remove('acknowledgements_test') return folders def _get_release_file_path(self, release_folder: str) -> str: responce = self._s3_client.list_objects_v2(Bucket=self.FEED_BUCKET, Prefix=release_folder) if not responce.get('Contents'): print('{release_folder} release not found on S3') raise release_file_name = responce.get('Contents')[0].get('Key') return 's3://{}'.format(os.path.join(self.FEED_BUCKET, release_file_name)) def _get_release_file_content(self, release_folder: str) -> ElementTree: release_file_path = self._get_release_file_path(release_folder) with smart_open.open(release_file_path, 'r') as fin: release_file_content = parse(fin) # pylint: disable=c-extension-no-member return release_file_content def _parse_content_by_xpaths(self, content: ElementTree, xpaths: Dict[str, Any]) -> List[str]: data = [] for (key, path) in xpaths.items(): try: if not path: data.append('') else: data.append(content.xpath(path)[0].text) except: print(f'Parsing failed, release: key: {key}, path: {path}') data.append('') return data def _parse_updated_release(self, release_folder: str, release_file_content: ElementTree) -> List[List[str]]: releases = [] # common fields for all releases xpaths_level_1 = { 'party_id': 'MessageHeader/MessageRecipient/PartyId', 'message_created_datetime': 'MessageHeader/MessageCreatedDateTime', } release_data_level_1 = self._parse_content_by_xpaths(release_file_content, xpaths_level_1) ################## # main release info for release in release_file_content.xpath('ReleaseList/Release'): xpaths_level_2 = { 'catalog_number': 'ReleaseId/CatalogNumber', 'upc': 'ReleaseId/ICPN', 'grid': 'ReleaseId/GRid', 'isrc': 'ReleaseId/ISRC', 'artist': 'ReleaseDetailsByTerritory/DisplayArtistName', 'title': 'ReferenceTitle/TitleText', 'release_type': 'ReleaseType', } release_data_level_2 = release_data_level_1 + self._parse_content_by_xpaths( release, xpaths_level_2 ) ################## # secondary release info for deal in release_file_content.xpath('DealList/ReleaseDeal/Deal/DealTerms'): xpaths_level_3 = { 'commercial_type': 'CommercialModelType', 'territories': 'TerritoryCode', 'price': None, 'sales_start_date': 'ValidityPeriod/StartDate', 'sales_end_date': 'ValidityPeriod/EndDate', } release_data_level_3 = release_data_level_2 + self._parse_content_by_xpaths( deal, xpaths_level_3 ) ################## # usage release info for usage in deal.xpath('Usage/UseType'): release_data_level_4 = [] release_data_level_4.extend(release_data_level_3) release_data_level_4.insert(9, usage.text) releases.append(release_data_level_4) return releases def _parse_purged_release(self, release_folder: str, release_file_content: ElementTree) -> List[List[str]]: xpaths = { 'party_id': 'MessageHeader/MessageRecipient/PartyId', 'message_created_datetime': 'MessageHeader/MessageCreatedDateTime', 'catalog_number': 'PurgedRelease/ReleaseId/CatalogNumber', 'upc': 'PurgedRelease/ReleaseId/ICPN', 'grid': 'PurgedRelease/ReleaseId/GRid', 'isrc': None, 'artist': None, 'title': 'PurgedRelease/Title/TitleText', 'release_type': None, 'commercial_type': None, 'use_type': None, 'territories': None, 'price': None, 'sales_start_date': None, 'sales_end_date': None, } release_data = self._parse_content_by_xpaths(release_file_content, xpaths) return [release_data] def _parse_release(self, release_folder: str, release_file_content: ElementTree) -> List[List[str]]: if release_file_content.xpath('PurgedRelease'): return self._parse_purged_release(release_folder, release_file_content) else: return self._parse_updated_release(release_folder, release_file_content) def process(self): workbook = xlsxwriter.Workbook('Signoff.xlsx') worksheet = workbook.add_worksheet() for col_num, data in enumerate(self.COLUMNS): worksheet.write(0, col_num, data) release_list = self._get_release_list() total_rows = 0 for release in release_list: release_file_content = self._get_release_file_content(release) release_data = self._parse_release(release, release_file_content) print(release_data) for release in release_data: for col_num, data in enumerate(release): worksheet.write(total_rows + 1, col_num, data) total_rows += 1 workbook.close() print(f'Done, total rows number: {total_rows}') if __name__ == '__main__': parser = ReleaseParser() parser.process()