from __future__ import print_function import base64 import json import uuid from aws_kinesis_agg.deaggregator import iter_deaggregate_records import boto3 import config import const from connectors import elasticache from connectors import mysql from models import statement_db from models import s3_bucket rds_conn = mysql.get_connection() redis_conn = elasticache.get_redis() s3_client = boto3.client('s3') def update_report(records): if config.SAVE_TO_S3: # TODO: use some kinesis ID instead of UUID for csv file name s3_bucket.save_csv_report(s3_client, records, str(uuid.uuid1())) statement_db.update_report(rds_conn, records) def lambda_handler(event, context): print('Received event') raw_kinesis_records = event['Records'] records_to_save = [] records_processed = 0 job_id = '' # Iterate through deaggregated records for record in iter_deaggregate_records(raw_kinesis_records): # Kinesis data in Python Lambdas is base64 encoded payload = base64.b64decode(record['kinesis']['data']) json_record = json.loads(payload) if not job_id: job_id = json_record.get('job_id', '') # save report for records with non-empty tuid only if json_record['tuid']: records_to_save.append(json_record) records_processed += 1 update_report(records_to_save) # decrement counter in Redis try: items_left = redis_conn.hincrby( const.KINESIS_COUNTER, job_id, -records_processed) print('{} items left in the queue'.format(items_left)) except Exception as ex: print(ex) finish_message = 'Successfully processed {} records.'.format( records_processed) print(finish_message) return finish_message