from __future__ import print_function from datetime import datetime import json import time import uuid import aws_kinesis_agg.aggregator import boto3 from pymysql import cursors import config from connectors import mysql from connectors import elasticache import const from models import sql kinesis_client = None def send_record(agg_record): if agg_record: print('Sending {} aggregated records'.format(agg_record.get_num_user_records())) global kinesis_client pk, ehk, data = agg_record.get_contents() kinesis_client.put_record( StreamName=config.SOURCE_KINESIS_STREAM, Data=data, PartitionKey=pk, ExplicitHashKey=ehk) def wait_for_finish(job_id): 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 __name__ == '__main__': start = datetime.utcnow() redis_conn = elasticache.get_redis() job_id = str(uuid.uuid1()) redis_conn.hset(const.KINESIS_COUNTER, job_id, 0) print('{}: Job {} started.'.format(datetime.utcnow(), job_id)) # we use non-buffering cursor because we are going to process # relatively large datasets conn = mysql.get_connection(cursorclass=cursors.SSDictCursor) kinesis_client = boto3.client('kinesis', region_name='us-east-1') kinesis_agg = aws_kinesis_agg.aggregator.RecordAggregator() rows_processed = 0 with conn.cursor() as cursor: # Get total rowcount cursor.execute(sql.GET_ROWCOUNT, { 'limit': config.MAX_RECORDS, 'offset': 0 }) rowcount = cursor.fetchone()['row_count'] # Save row count to Redis redis_conn.hincrby(const.KINESIS_COUNTER, job_id, rowcount) cursor.execute(sql.GET_BATCH, { 'limit': config.MAX_RECORDS, 'offset': 0 }) for row in cursor: row['job_id'] = job_id pk = int(uuid.uuid1()) data = json.dumps(row) result = kinesis_agg.add_user_record(pk, data) if result: send_record(result) elif (kinesis_agg.get_num_user_records() == const.DYNAMO_BATCH_SIZE): send_record(kinesis_agg.clear_and_get()) rows_processed += 1 if rows_processed >= config.MAX_RECORDS: break conn.close() send_record(kinesis_agg.clear_and_get()) 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)