from __future__ import print_function import base64 import json from aws_kinesis_agg.deaggregator import iter_deaggregate_records from aws_kinesis_agg import aggregator import boto3 import config import const from util import retry dynamodb_client = boto3.client('dynamodb') kinesis_client = boto3.client('kinesis', region_name=config.AWS_REGION) kinesis_agg = aggregator.RecordAggregator() def lambda_handler(event, context): print('Received event') raw_kinesis_records = event['Records'] deaggregated_records = [] # 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) deaggregated_records.append(json_payload) print('Deaggregated {} records'.format(len(deaggregated_records))) # Split records into batches by 100 isrcs to effectively query # DynamoDB. It allows to query up to 100 items at once with batch_get_item. batches = chunks(deaggregated_records) for batch in batches: isrcs = set([i['isrc'] for i in batch]) ownership_records = get_records_from_dynamo(isrcs) enriched_records = enrich_records(ownership_records, batch) send_records_to_output_stream(enriched_records) finish_message = 'Successfully processed {} records.'.format( len(deaggregated_records)) print(finish_message) return finish_message def chunks(long_list, length=const.DYNAMO_BATCH_SIZE): """Split a list into chunks with desired length""" length = max(1, length) return (long_list[i:i + length] for i in xrange(0, len(long_list), length)) def enrich_records(ownership_records, batch): enriched = [] for item in batch: isrc = item['isrc'] job_id = item['job_id'] territories = ownership_records.get(isrc, {}) for terr in item['territories'].split(','): tuid = territories.get(terr, '') enriched.append({ 'isrc': isrc, 'tuid': tuid, 'territory': terr, 'job_id': job_id }) return enriched @retry(retry_count=3, error_condition=lambda err: True) def get_records_from_dynamo(isrcs): record_keys = map(lambda isrc: {'isrc': {'S': isrc}}, isrcs) response = dynamodb_client.batch_get_item( RequestItems={ config.MR_ACTIVE_TABLE_NAME: { 'Keys': record_keys, 'AttributesToGet': [ 'isrc', 'territories', ], 'ConsistentRead': True, } }, ReturnConsumedCapacity='TOTAL') isrc_territories = {} for item in response['Responses'][config.MR_ACTIVE_TABLE_NAME]: isrc = item['isrc']['S'] territories = item['territories']['M'] isrc_territories[isrc] = { k: int(v['M']['tuid']['N']) for k, v in territories.iteritems()} return isrc_territories def send_record(agg_record): if agg_record: pk, ehk, data = agg_record.get_contents() kinesis_client.put_record( StreamName=config.DESTINATION_KINESIS_STREAM, Data=data, PartitionKey=pk, ExplicitHashKey=ehk) print('Aggregated record sent to kinesis. pk: ' + pk) def send_records_to_output_stream(records): print('Sending {} to output stream'.format(len(records))) for rec in records: pk = rec['isrc'] data = json.dumps(rec) result = kinesis_agg.add_user_record(pk, data) send_record(result) send_record(kinesis_agg.clear_and_get())