"""Lambda pp_identity function module.""" from kafka_utils.producer.event import EventProducer from kafka_utils.producer.serializer.simple_json import SimpleJSONSerializer from lambdacommon.common_config import logger import config def handler(event, context): """Lambda entry point.""" try: produce_records_to_kafka(event['Records']) return {'status': 'OK'} except Exception as e: logger.exception(str(e)) raise e def produce_records_to_kafka(records: list[dict]): """Produce records to Kafka. Args: records: All records received from DynamoDB Streams trigger. """ with EventProducer( bootstrap_servers=config.KAFKA_BOOTSTRAP_SERVERS, value_serializer=SimpleJSONSerializer(), security_protocol=config.KAFKA_SECURITY_PROTOCOL ) as producer: for record in records: transformed_record = transform(record) producer.produce( topic=config.KAFKA_TARGET_TOPIC, event_key=None, event_value=transformed_record, auto_flush=False ) def transform(item: dict): """Transform record to camel-aws-ddb-streams-source Kafka Connector format. This recursive function walks through the record and does the following: - Converts keys from Pascal Case to Camel Case - Adds a 'type' key to the DynamoDB items Args: item: Record to transform. Returns: dict: Transformed record. """ if isinstance(item, dict): # d = {_change_key_case(k): transform(v) for k, v in item.items()} d = {} for k, v in item.items(): if k == 'ApproximateCreationDateTime': v = int(v * 1000) d[_change_key_case(k)] = transform(v) if len(d) == 1: key = list(item.keys())[0] if key.upper() in config.DYNAMODB_DATA_TYPES: d['type'] = key for k, v in config.CAMEL_CONNECTOR_DEFAULT_ATTRIBUTES.items(): if k not in d: d[k] = v return d elif isinstance(item, list): return [transform(i) for i in item] else: return item def _change_key_case(key): if key in config.DYNAMODB_DATA_TYPES: return key.lower() return key[0].lower() + key[1:]