"""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 from connectors import elasticache import config 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): """Wait until report processing job finishes. Args: job_id (str): id of a processing job """ start_time = datetime.utcnow() while True: count = redis_conn.hget(const.KINESIS_COUNTER, job_id) print('{}: {} records left.\n'.format( datetime.utcnow(), count), end='\r') time.sleep(1) if int(count) <= 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)) break #decrease_read_capacity(active_table) def scan_dynamo(segment, total_segments, client, job_id): global total_count scan_args = dict( TableName=config.MR_ACTIVE_TABLE_NAME, Segment=segment, TotalSegments=total_segments ) 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 result.get('LastEvaluatedKey'): key = result.get('LastEvaluatedKey') else: print('Scanning {} segment completed.'.format(segment)) break send_record(kinesis_agg.clear_and_get()) def send_to_kinesis(items, job_id): global rows_processed # record aggregator is not thread-safe aggregator_lock.acquire() for item in items: data = util.dynamo_acive_record_to_json(item) 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 if rows_processed >= config.MAX_RECORDS: break aggregator_lock.release() 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 """ 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() print('{} finished'.format(worker.name)) print('All workers finished.') def increase_read_capacity(active_table): # increase masters_active read capacity to prevent throttle current_read_cap = active_table.provisioned_throughput.get( const.READ_CAPACITY_UNITS) print('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) print('New DynamoDB read capacity is {}'.format(increased_read_cap)) def decrease_read_capacity(active_table): # Decrease read capacity if necessary current_read_cap = active_table.provisioned_throughput.get( const.READ_CAPACITY_UNITS) print('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) print('Decreased DynamoDB read capacity by {}'.format( config.DDB_ADDITIONAL_READ_CAP)) else: print('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) print('{}: 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() #increase_read_capacity(active_table) rows_processed = 0 print('Getting rowcount...') # Get total rowcount (please note, the item count can be not accurate) rowcount = active_table.item_count # Save row count to Redis print('Rowcount: {}'.format(rowcount)) redis_conn.hincrby(const.KINESIS_COUNTER, job_id, rowcount) # Get records from Dynamo and send to Kinesis do_export(config.TOTAL_SEGMENTS, dynamodb_client, job_id) print('{}: {} records sent for processing.\n'.format( datetime.utcnow(), rowcount)) # Check the redis counter and wait until all records will be processed wait_for_finish(job_id)