"""Find missing dynamo items. The purpose of this script is to list the monthly vendors who had AVRO files generated for a given period and to check that a corresponding DynamoDB record exists for it. When a missing record is found, the replacement for that item will be printed to stdout. It is likely you'll only ever need to adjust the PERIOD_ID value. """ import json import os import boto3 from boto3.dynamodb.conditions import Key from dotenv import load_dotenv load_dotenv() # configurable options PERIOD_ID = os.environ.get('PERIOD_ID') # '249' or '247,248,249' assert len(PERIOD_ID) == 3 or len(PERIOD_ID.split(',')) == 3 USER_TYPE = os.environ.get('USER_TYPE') assert USER_TYPE in ['label', 'subaccount'] PERIOD_TYPE = 'month' if len(PERIOD_ID) == 3 else 'quarter' S3_REPORT_BUCKET = 'prod-statement-detail-exports' S3_REPORT_PATH_TEMPLATE = 'schematized_files/{}_{}_all_all_{}/' S3_REPORT_PATH = S3_REPORT_PATH_TEMPLATE.format( PERIOD_ID, USER_TYPE, PERIOD_TYPE) S3_LABEL_REPORT_PATH = S3_REPORT_PATH + 'user_id_type={}/' DYNAMO_TABLE = 'prod_accounting_statement_export' PARTITION_KEY = 'user_id_type' SORT_KEY = 'user_params' AWS_REGION = 'us-east-1' MATCH_FILE_TYPE = 'AVRO' def main(): """Main function.""" s3_client = boto3.client('s3') dynamodb = boto3.resource('dynamodb', region_name=AWS_REGION) table = dynamodb.Table(DYNAMO_TABLE) vendor_ids = get_folder_listings(s3_client) for vendor_id in vendor_ids: missing_id = dynamo_record_is_missing(vendor_id, table) if missing_id: format_replacement_entry(vendor_id, s3_client) def dynamo_record_is_missing(vendor_id, table): """Query DynamoDb.""" filter_ex = Key(PARTITION_KEY).eq(vendor_id) \ & Key(SORT_KEY).begins_with(PERIOD_ID) result = table.query(KeyConditionExpression=filter_ex) for item in result.get('Items'): if item.get('file_type') == MATCH_FILE_TYPE: return False return True def get_folder_listings(s3_client): """Main function.""" folders = [] options = { 'Bucket': S3_REPORT_BUCKET, 'Prefix': S3_REPORT_PATH, 'Delimiter': '/' } remove_len = len(options['Prefix'] + PARTITION_KEY + '=') while True: results = s3_client.list_objects_v2(**options) options['ContinuationToken'] = results.get('NextContinuationToken') if not results.get('CommonPrefixes'): print('No S3 objects found.') print(options) break for item in results.get('CommonPrefixes'): folders.append(item.get('Prefix')[remove_len:-1]) if not results.get('NextContinuationToken'): break return folders def format_replacement_entry(vendor_id, s3_client): """Print the item that should be inserted into dynamo.""" options = { 'Bucket': S3_REPORT_BUCKET, 'Prefix': S3_LABEL_REPORT_PATH.format(vendor_id), 'Delimiter': '/' } results = s3_client.list_objects_v2(**options) if results.get('Contents'): new_item = { 'file_type': MATCH_FILE_TYPE, 'period_ids': PERIOD_ID, 's3_path': 's3://{}/{}'.format( S3_REPORT_BUCKET, results.get('Contents')[0].get('Key')), 'status': 'GENERATED', 'user_id_type': vendor_id, 'user_params': '{}__all__AVRO__en_US'.format(PERIOD_ID) } print(json.dumps(new_item)) return print('what happened with {} in {}'.format(vendor_id, options['Prefix'])) if __name__ == '__main__': main()