"""Lambda for saving DynamoDB items to SnowFlake.""" from __future__ import print_function import base64 import json from aws_kinesis_agg.deaggregator import iter_deaggregate_records import const import field_const from connectors import snowflakedb from connectors import elasticache from models import snowflake as snowflake_model redis_conn = elasticache.get_redis() snowflake_context = snowflakedb.get_snowflake_context() def lambda_handler(event, context): """Consumer Lambda function handler. Args: event (dict): Lambda event information context (dict): Lambda event context Returns: str: result message """ raw_kinesis_records = event['Records'] deaggregated_records = [] 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_payload = json.loads(payload) if not job_id: job_id = json_payload.get(field_const.JOB_ID, '') deaggregated_records.append(json_payload) records_count = len(deaggregated_records) print('Deaggregated {} records'.format(records_count)) snowflake_model.update_snowflake( deaggregated_records, 'INSERT', snowflake_context) update_redis_counter(job_id, records_count) finish_message = 'Successfully processed {} records.'.format(records_count) print(finish_message) return finish_message def update_redis_counter(job_id, records_processed): """Update records counter in Redis. Args: job_id (str): import job id. records_processed (int): number of records processed. """ 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)