"""Producer script for filling the Kinesis stream with MR data.""" from __future__ import print_function from datetime import datetime import json import sys import threading import time import uuid import aws_kinesis_agg.aggregator import boto3 from botocore import exceptions import config from config import logger as log from connectors import elasticache import const import util kinesis_client = None active_table = None total_count = 0 aggregator_lock = threading.Lock() def change_dynamo_read_cap(cap): """Change dynamoDB table's read capacity. Args: cap (int): read capacity """ global active_table current_write_cap = active_table.provisioned_throughput.get( const.WRITE_CAPACITY_UNITS) print(current_write_cap) print(cap) try: active_table.update( ProvisionedThroughput={ const.READ_CAPACITY_UNITS: cap, const.WRITE_CAPACITY_UNITS: current_write_cap } ) except exceptions.ClientError: print('Could not decrease read capacity. Please contact Systems.') def send_record(agg_record): """Send aggregated record to Kinesis stream. Args: agg_record (AggRecord): aggregated Kinesis record """ if agg_record: print('Sending {} aggregated records'.format( agg_record.get_num_user_records())) global kinesis_client try: pk, ehk, data = agg_record.get_contents() kinesis_client.put_record( StreamName=config.BACKFILL_KINESIS_STREAM, Data=data, PartitionKey=pk, ExplicitHashKey=ehk) except Exception as ex: print('Failed to send kinesis record: {}'.format(ex)) def wait_for_finish(job_id, redis_conn, start, rowcount): """Wait until report processing job finishes. Args: job_id (str): id of a processing job redis_conn (StrictRedis): redis client start (datetime): producer start datetime rowcount (int): total rowcount read from DynamoDB Return: bool: True if finished on time, False if timeout exceeded """ start_time = datetime.utcnow() finished_on_time = True while True: count = int(redis_conn.hget(const.KINESIS_COUNTER, job_id)) records_left = rowcount - count print('{}: {} records left.\n'.format( datetime.utcnow(), records_left), end='\r') time.sleep(1) if records_left <= 0: print('\n\n{}: finished.'.format(datetime.utcnow())) delta = datetime.utcnow() - start print('{} records processed within {}sec'.format( rowcount, delta.total_seconds() )) break if (datetime.utcnow() - start_time).total_seconds() > config.TIMEOUT: print('Timeout for job {} exceeded.'.format(job_id)) finished_on_time = False break redis_conn.hdel(const.KINESIS_COUNTER, job_id) if config.DDB_INCREASE_CAPACITY: decrease_read_capacity(active_table) return finished_on_time def scan_dynamo(segment, total_segments, client, job_id): """Scan dynamo DB table for records. Args: segment (int): number of segment total_segments (int): total number of segments client (BaseClient): boto DynamoDB client job_id (str): import job id """ global total_count scan_args = dict( TableName=config.MR_ACTIVE_TABLE_NAME, Segment=segment, TotalSegments=total_segments, Select='SPECIFIC_ATTRIBUTES', ProjectionExpression='isrc, territories, locked_territories' ) key = None while True: print('Scanning {}, LastEvaluatedKey={}'.format(segment, key)) if key: scan_args['ExclusiveStartKey'] = key try: result = client.scan(**scan_args) except Exception as ex: print('Failed to scan DynamoDB: {}'.format(ex)) continue if result.get('Items'): cnt = len(result.get('Items')) print('Received {} isrcs'.format(cnt)) send_to_kinesis(result.get('Items'), job_id) total_count += cnt if total_count >= config.MAX_RECORDS: print( 'Scanning {} segment completed, max records read.'.format( segment)) break if result.get('LastEvaluatedKey'): key = result.get('LastEvaluatedKey') else: print('Scanning {} segment completed.'.format(segment)) break try: # record aggregator is not thread-safe aggregator_lock.acquire() send_record(kinesis_agg.clear_and_get()) except Exception as ex: print(ex) finally: aggregator_lock.release() def send_to_kinesis(items, job_id): """Send records to kinesis stream. Args: items (list): list of items from DynamoDB job_id (str): import job id """ global rows_processed try: # record aggregator is not thread-safe aggregator_lock.acquire() print('Aggregator lock acquired;') for item in items: data = util.dynamo_active_record_to_json(item, include_locked=True) data['job_id'] = job_id pk = int(uuid.uuid1()) data = json.dumps(data) result = kinesis_agg.add_user_record(pk, data) if result: send_record(result) rows_processed += 1 except Exception as ex: print(ex) finally: aggregator_lock.release() print('Aggregator lock released;') def do_export(total_segments, dynamo_client, job_id): """Read data from DynamoDB and send to Kinesis. Args: total_segments (int): number of threads to scan the table dynamo_client (BaseClient): boto3 dynamoDB client job_id (str): import job id """ workers = [] for segment in range(total_segments): worker = threading.Thread( target=scan_dynamo, name='Worker {}'.format(segment), args=(segment, total_segments, dynamo_client, job_id)) workers.append(worker) for worker in workers: worker.start() for worker in workers: worker.join() log.info('{} finished.'.format(worker.name)) log.info('All workers finished.') def increase_read_capacity(active_table): """Increase read capacity of DynamoDB table. Args: active_table (Table): boto3 dynamoDB table resource """ # increase masters_active read capacity to prevent throttle current_read_cap = active_table.provisioned_throughput.get( const.READ_CAPACITY_UNITS) log.info('Current DynamoDB read capacity is {}'.format(current_read_cap)) increased_read_cap = current_read_cap + config.DDB_ADDITIONAL_READ_CAP change_dynamo_read_cap(increased_read_cap) log.info('New DynamoDB read capacity is {}'.format(increased_read_cap)) def decrease_read_capacity(active_table): """Decrease read capacity of DynamoDB table. Args: active_table (Table): boto3 dynamoDB table resource """ # Decrease read capacity if necessary current_read_cap = active_table.provisioned_throughput.get( const.READ_CAPACITY_UNITS) log.info('Current read capacity is {}'.format(current_read_cap)) if current_read_cap > config.DDB_ADDITIONAL_READ_CAP: decreased_read_cap = current_read_cap - config.DDB_ADDITIONAL_READ_CAP change_dynamo_read_cap(decreased_read_cap) log.info('Decreased DynamoDB read capacity by {}'.format( config.DDB_ADDITIONAL_READ_CAP)) else: log.info('Will not change read capacity') if __name__ == '__main__': start = datetime.utcnow() redis_conn = elasticache.get_redis() if not redis_conn: sys.exit('Failed to connect to Redis.') job_id = str(uuid.uuid1()) redis_conn.hset(const.KINESIS_COUNTER, job_id, 0) log.info('{}: Job {} started.'.format(datetime.utcnow(), job_id)) dynamodb_resource = boto3.resource('dynamodb', region_name='us-east-1') dynamodb_client = boto3.client('dynamodb', region_name='us-east-1') kinesis_client = boto3.client('kinesis', region_name='us-east-1') active_table = dynamodb_resource.Table(config.MR_ACTIVE_TABLE_NAME) kinesis_agg = aws_kinesis_agg.aggregator.RecordAggregator() if config.DDB_INCREASE_CAPACITY: increase_read_capacity(active_table) rows_processed = 0 log.info('Getting rowcount...') # Get total rowcount (please note, the item count can be not accurate) rowcount = active_table.item_count log.info('Number of records in dynamo: {}'.format(rowcount)) # Get records from Dynamo and send to Kinesis do_export(config.TOTAL_SEGMENTS, dynamodb_client, job_id) log.info('{}: {} records sent for processing.\n'.format( datetime.utcnow(), rows_processed)) # Check the redis counter and wait until all records will be processed wait_for_finish(job_id, redis_conn, start, rows_processed)