"""Script for parsing Dynamo DB formatted records from MR_AUDIT table.""" from functools import partial import json import sys # Some shortcuts for dynamo types j = json.loads def get_dynamo_item(item, dynamo_type, default_value, cast=None): """Get DynamoDB value.""" if not item: return default_value result = item.get(dynamo_type, default_value) if cast: result = cast(result) return result s = partial(get_dynamo_item, dynamo_type='s', default_value=None) m = partial(get_dynamo_item, dynamo_type='m', default_value={}) l = partial(get_dynamo_item, dynamo_type='l', default_value=[]) n = partial(get_dynamo_item, dynamo_type='n', default_value=None, cast=int) def dynamo_audit_record_to_json(item): """Parse JSON item received from hive. Args: item (dict): partially parsed JSON items. (all values are strings) Returns: str: parsed and transformed JSON dumped as string """ parsed = {} for k, v in item.items(): parsed[k] = json.loads(v) opcode = s(parsed.get('opcode')) timestamp = n(parsed.get('timestamp'), cast=float) result = { 'isrc': s(parsed.get('isrc')), 'data': { 'territories': {}, 'opcode': opcode, 'lock_reason': s(parsed.get('lock_reason')), 'user': s(parsed.get('user')), 'source': s(parsed.get('source')), 'correlation_id': s(parsed.get('correlation_id')), 'conflict': None, }, 'timestamp': timestamp / 1000 } # 'territories' could be either list or map if 'm' in parsed.get('territories'): territories = m(parsed.get('territories')) for terr, tuid in territories.items(): result['data']['territories'][terr] = n(tuid) else: territories = l(parsed.get('territories')) for territory in territories: terr = s(territory) result['data']['territories'][terr] = None if 'conflict' in parsed: conflict = m(parsed.get('conflict')) conflict['status'] = s(conflict['status']) conflict['resolved'] = n(conflict.get('resolved')) tuids = [] for tuid in l(conflict.get('conflicting_tuid')): tuids.append(n(tuid)) conflict['conflicting_tuid'] = tuids result['data']['conflict'] = conflict return json.dumps(result) if __name__ == '__main__': while True: line = sys.stdin.readline() if not line: break try: obj = json.loads(line) print(dynamo_audit_record_to_json(obj)) except: continue