import json from time import sleep from pymongo import MongoClient from kafka import KafkaProducer from kafka.errors import NoBrokersAvailable from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler from sientia_do.notifications.models import NotificationLevel from sientia_do.temporal.constants import now from sientia_do.temporal.activities.base import BaseActivity from sientia_do.observability.logger import Logger import traceback import ingestor.metrics as metrics import os class DataManager(BaseActivity): def __init__( self, kafka_servers: str, mongo_connection_string: str, mongo_database: str, export_to_kafka: bool, metadata: dict, logger: Logger, notification_handler: NotificationHandler, ) -> None: """ Initializes the DataManager instance with a Kafka producer. This constructor attempts to establish a connection to the specified Kafka servers and initializes a Kafka producer for sending messages. It retries the connection up to 3 times if the Kafka servers are unavailable. Args: kafka_servers (str): A comma-separated string of Kafka server addresses. logger (Logger): A logger instance for logging messages. Raises: NoBrokersAvailable: If the connection to Kafka servers fails after 3 attempts. """ self.pod_id = os.getenv("HOSTNAME", "localhost") self.kafka_producer = None self.export_to_kafka = export_to_kafka if self.export_to_kafka: for i in range(0, 3): logger.info( f"Trying ({i}) to initializing DataManager with Kafka servers: {kafka_servers}" ) try: self.kafka_producer = KafkaProducer( bootstrap_servers=kafka_servers, value_serializer=lambda v: json.dumps(v).encode( "utf-8" ), # Serialize JSON messages key_serializer=lambda k: str( k).encode("utf-8") if k else None, ) # Kafka connected metrics.KAFKA_CONNECTION_STATUS.labels( pod_id=self.pod_id).set(1) break except NoBrokersAvailable: logger.error( f"Kafka servers {kafka_servers} are not available. Retrying..." ) sleep(5) else: # Kafka not connected metrics.KAFKA_CONNECTION_STATUS.labels( pod_id=self.pod_id).set(0) logger.error( f"Failed to connect to Kafka servers {kafka_servers} after 3 attempts." ) raise NoBrokersAvailable( f"Failed to connect to Kafka servers {kafka_servers} after 3 attempts." ) logger.info( f"DataManager initialized with Kafka servers: {kafka_servers}") logger.info( f"Trying to initializing DataManager with MongoDB servers: {mongo_connection_string}" ) self.connection_string = mongo_connection_string self.database = mongo_database self.mongo_client = MongoClient(self.connection_string) self.mongo_client.server_info() self.metadata = metadata self.mongo_db = self.mongo_client[self.database] logger.info( f"DataManager initialized with MongoDB servers: {self.connection_string}" ) BaseActivity.__init__(self, logger=logger, notification_handler=notification_handler, set_error_counter=True) def shutdown(self): """Closes the Kafka producer connection.""" if self.kafka_producer: try: self.kafka_producer.flush(timeout=10) self.kafka_producer.close() # Mark as disconnected metrics.KAFKA_CONNECTION_STATUS.labels( pod_id=self.pod_id).set(0) except Exception as e: self.logger.error(f"Error closing Kafka producer: {e}") else: self.logger.warning( "Kafka producer is already closed or not initialized.") if self.mongo_client: try: self.mongo_client.close() except Exception as e: self.logger.error(f"Error closing MongoDB client: {e}") else: self.logger.warning( "MongoDB client is already closed or not initialized.") def __del__(self): self.shutdown() def delivery_report(self, msg: str): """Callback for delivery reports from Kafka.""" self.logger.debug( f"Record successfully produced to {msg.topic} [{msg.partition}] at offset {msg.offset}" ) def delivery_error(self, err: str): """Callback for delivery reports from Kafka.""" self.logger.error(f"Delivery failed for record : {err}") def publish(self, topic: str, data: dict) -> None: """ Publishes a message to a specified Kafka topic. Args: topic (str): The name of the Kafka topic to which the message will be published. data (dict): The message data to be sent to the Kafka topic. Returns: None Raises: Exception: If there is an error during message delivery, it will be handled by the `delivery_error` callback. """ if self.export_to_kafka: try: self.logger.debug( f"Publishing message to topic {topic}: {data}") self.kafka_producer.send(topic=topic, value=data).add_callback( self.delivery_report ).add_errback(self.delivery_error) self.kafka_producer.flush(timeout=10) metrics.KAFKA_MESSAGES_SENT.labels( pod_id=self.pod_id, topic=topic).inc() except Exception as e: metrics.KAFKA_MESSAGES_ERRORS.labels( pod_id=self.pod_id, topic=topic).inc() trace = traceback.format_exc() self.send_notification( metadata=self.metadata, notification_id=f"KAFKA_PRODUCER_ERROR_{topic}", message=f"Error publishing message to topic {topic}: {e}", block="kafka_producer", level=NotificationLevel.ERROR, attachment_content=trace, ) self.logger.error(trace) try: collection = self.mongo_db[topic] collection.insert_one( { **data, "inserted_at": now(), } ) self.logger.debug( f"Message inserted into MongoDB collection {topic}: {data}") metrics.TAG_WRITTEN_COUNT.labels( pod_id=self.pod_id, tag_name=data["name"], collection_name=topic ).inc() except Exception as e: trace = traceback.format_exc() self.send_notification( metadata=self.metadata, notification_id=f"MONGO_PRODUCER_ERROR_{topic}", message=f"Error inserting message to MongoDB: {e}", block="mongo_producer", level=NotificationLevel.ERROR, attachment_content=trace, ) self.logger.error(trace)