"""Lambda kinesis-to-neo4j function module.""" import base64 import json import logging import sentry_sdk from aws_kinesis_agg import deaggregator from lambdacommon.common_config import logger from neo4j import exceptions as neo4j_exceptions from neobolt import exceptions as neobolt_exceptions from sentry_sdk.integrations.aws_lambda import AwsLambdaIntegration from sentry_sdk.integrations.logging import LoggingIntegration from splitio.api import APIException import config from src.logic.record_handlers import get_record_handler from src.logic.retry_exception import RetriableError logging_integration = LoggingIntegration( level=logging.INFO, # Capture info and above as breadcrumbs event_level=logging.CRITICAL, # Send only critical as events ) sentry_sdk.init(config.SENTRY_DSN, integrations=[AwsLambdaIntegration(), logging_integration]) def handler(event, context): """Lambda entry point.""" logger.info(event) try: if 'Records' in event: preprocessed_records = _preprocess_records(event['Records']) _process_records(preprocessed_records) else: logger.warning(f'Unknown event! {event}') return {'status': 'OK'} except Exception as e: logger.exception(str(e)) raise e def _process_records(records): for i, record in enumerate(records, start=1): logger.info('Processing record number {} from batch'.format(i)) record_handler = get_record_handler(record) if record_handler: try: record_handler(record) except ( neobolt_exceptions.ConnectionExpired, neobolt_exceptions.NotALeaderError, neobolt_exceptions.SecurityError, neobolt_exceptions.TransientError, neo4j_exceptions.DatabaseError, neo4j_exceptions.SessionExpired, neo4j_exceptions.TransientError, neo4j_exceptions.ConstraintError, ConnectionResetError, APIException, RetriableError, ): logger.info('Pushing to Sentry retriable error context.') with sentry_sdk.new_scope() as scope: scope.set_tag('is_retried', 'yes') sentry_sdk.capture_exception() raise except Exception: logger.info('Pushing to Sentry non-retriable error context.') with sentry_sdk.new_scope() as scope: scope.set_tag('is_retried', 'no') sentry_sdk.capture_exception() else: logger.warning(f'There is no handler for: {record.get("type")} {record.get("table")} record.') def _preprocess_records(raw_records): """Prepare raw records for further processing.""" flatten_records = deaggregator.iter_deaggregate_records(raw_records) decoded_records = list(map(_decode_record, flatten_records)) try: sorted_records = sorted(decoded_records, key=lambda x: (x['xid'], x['ts'], x.get('xoffset', 0))) except KeyError: logger.error('Some records have invalid structure. Falling back to the original order.') sorted_records = decoded_records if decoded_records != sorted_records: count = sum(1 for a, b in zip(decoded_records, sorted_records) if a != b) logger.info(f'The order of {count} messages is incorrect! Batch size is {len(decoded_records)}') logger.info('Initial list of messages') logger.info(decoded_records) logger.info('Sorted list of messages') logger.info(sorted_records) return sorted_records def _decode_record(record): """Decode raw record.""" record_str = base64.b64decode(record['kinesis']['data']) return json.loads(record_str)