import boto3 import json import time import config import sys import os import logging import argparse import aws_kinesis_agg.aggregator from util import formatter from datetime import datetime from connector import mysql from sql import releases, projects, tracks, labels sys.path.append(os.path.join(os.path.dirname(__file__), '..')) log = logging.getLogger('main') log.setLevel(logging.DEBUG) fmt = logging.Formatter('%(asctime)s - %(levelname)s - %(message)s') sh = logging.StreamHandler(sys.stdout) sh.setFormatter(fmt) log.addHandler(sh) OFFSET = 0 LIMIT = 10000 # A releases.last_updated date START_DATE = '2019-10-01 00:00:00' END_DATE = datetime.now().strftime('%Y-%m-%d %H:%M:%S') ITERATION_TABLE = 'releases' def send_record(agg_record): if agg_record is None: return log.info('No records found for provided criteria.') kinesis_client = boto3.client('kinesis') pk, ehk, data = agg_record.get_contents() kinesis_client.put_record(StreamName=config.KINESIS_STREAM_NAME,Data=data,PartitionKey=pk) def main(): parser = argparse.ArgumentParser() parser.add_argument('--offset', type=int, default=OFFSET) parser.add_argument('--limit', type=int, default=LIMIT) parser.add_argument('--start_date', type=str, default=START_DATE) parser.add_argument('--end_date', type=str, default=END_DATE) parser.add_argument('--iteration_table', type=str, default=ITERATION_TABLE) args = parser.parse_args() log.info('STARTED producing kinesis events for %s', args.iteration_table) log.info( 'start_date: %s end_date: %s limit: %s offset: %s table %s', args.start_date, args.end_date, args.limit, args.offset, args.iteration_table) ar_db = None cursor = None ar_db = mysql.get_ar_mysql_connection() cursor = ar_db.cursor() if args.iteration_table == 'releases': cursor.execute(releases.SELECTSQL.format(limit=args.limit, offset=args.offset), { 'start_date': args.start_date, 'end_date': args.end_date } ) elif args.iteration_table == 'project': cursor.execute(projects.SELECTSQL.format(limit=args.limit, offset=args.offset), { 'start_date': args.start_date, 'end_date': args.end_date } ) elif args.iteration_table == 'track': cursor.execute(tracks.SELECTSQL.format(limit=args.limit, offset=args.offset), { 'start_date': args.start_date, 'end_date': args.end_date } ) elif args.iteration_table == 'vendor': cursor.execute(labels.SELECTSQL.format(limit=args.limit, offset=args.offset), { 'start_date': args.start_date, 'end_date': args.end_date } ) items = cursor.fetchall() kinesis_agg = aws_kinesis_agg.aggregator.RecordAggregator() kinesis_agg.on_record_complete(send_record) for item in items: record_event = formatter.get_kinesis_event(item, args.iteration_table) log.info('Kinesis event : %s', json.dumps(record_event)) kinesis_agg.add_user_record('partition-key', bytes(json.dumps(record_event), 'utf-8'), 'ehk') #Clear out any remaining records that didn't trigger a callback yet send_record(kinesis_agg.clear_and_get()) log.info('COMPLETE kinesis events for %s', args.iteration_table) if __name__ == "__main__": main()