"""Kafka Producer initialization.""" import json import logging from typing import Dict from confluent_kafka.serialization import StringSerializer from kafka_utils.producer.event import EventProducer import config from utils.kafka_serializer import JSONOrNoneSerializer logger = logging.getLogger(config.APPLICATION_NAME) def get_producer(): """Get kafka producer.""" return EventProducer( bootstrap_servers=config.KAFKA_BOOTSTRAP_SERVERS, key_serializer=StringSerializer(), value_serializer=JSONOrNoneSerializer(), security_protocol=config.KAFKA_SECURITY_PROTOCOL ) def get_kafka_key_for_delete(data): """Remove prefix from id to get numeric id.""" # During delete operation, logic layer only sends id and type not fields. # However, the id is prefixed by strings like rel{}, track, proj. # For Opensearch we want numeric values as the key. key = str(data.get('id')) kafka_key = key.lstrip('rel').lstrip('track').lstrip('proj') return kafka_key def get_kafka_key_and_value(event_type, message_key, data): """Return key and value for kafka message based on event_type.""" if not event_type or event_type not in ['add', 'delete']: logger.error( 'Bad type of event. Its not add or delete.', extra={'bad_data': json.dumps(data)} ) return None, None if event_type == 'delete': # tombstone message for deleting the record from Opensearch. key = get_kafka_key_for_delete(data) return key, None # when event_type is 'add', we need both key and value. if not data.get('fields'): logger.error( 'Bad data. There are no fields to write to kafka.', extra={'bad_data': json.dumps(data)} ) return None, None kafka_key = data.get('fields').get(message_key) if not kafka_key: logger.error( 'Bad data. There are no key to write to kafka.', extra={'bad_data': json.dumps(data)} ) return None, None kafka_value = data.get('fields') return kafka_key, kafka_value def write_to_kafka(domain: str, messages: list[Dict]): """Produce a message to Kafka topic. Args: domain (str): domain of event ie: releases, tracks etc. messages (list[dict]): list of messages to write to kafka topic. Returns: None """ if not config.KAFKA_TARGET_TOPICS.get(domain): return message_key = config.KAFKA_KEY_FIELD_IN_DOC.get(domain) if not message_key: return with get_producer() as producer: logger.info(f'Writing to kafka topic: {config.KAFKA_TARGET_TOPICS.get(domain)}') for each in messages: event_type = each.get('type') kafka_key, kafka_value = get_kafka_key_and_value( event_type, message_key, each) if not kafka_key: continue producer.produce( topic=config.KAFKA_TARGET_TOPICS.get(domain), event_key=str(kafka_key), event_value=kafka_value, callback=msg_delivery_callback, auto_flush=False ) def msg_delivery_callback(err, msg): """Handle Kafka message delivery result.""" if err: logger.error( f'There was an error: {err} delivering message with key: {msg.key()} ' f'with data: {msg.value()} to Kafka')