diff --git a/.vscode/settings.json b/.vscode/settings.json index bdd01e8..a3a53a3 100644 --- a/.vscode/settings.json +++ b/.vscode/settings.json @@ -7,6 +7,6 @@ "projectKey": "Aignosi_sientia-dataops-opc-ingestor_1642c8bf-a148-4911-8362-0903d7fef99a" }, "python.languageServer": "Pylance", - "python.analysis.typeCheckingMode": "off", + "python.analysis.typeCheckingMode": "basic", "editor.suggestSelection": "first" } diff --git a/ingestor/managers/data_manager.py b/ingestor/managers/data_manager.py index b16b933..bb2c1e4 100644 --- a/ingestor/managers/data_manager.py +++ b/ingestor/managers/data_manager.py @@ -6,11 +6,17 @@ from kafka.errors import NoBrokersAvailable from sientia_do.notifications.handlers import NotificationHandler from sientia_do.notifications.models import NotificationLevel import traceback +import ingestor.metrics as metrics +import os -class DataManager(): - def __init__(self, kafka_servers: str, logger: Logger, - notification_handler: NotificationHandler) -> None: +class DataManager: + def __init__( + self, + kafka_servers: str, + 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 @@ -23,31 +29,39 @@ class DataManager(): NoBrokersAvailable: If the connection to Kafka servers fails after 3 attempts. """ + self.pod_id = os.getenv("HOSTNAME", "localhost") self.kafka_producer = None for i in range(0, 3): logger.info( - f"Trying ({i}) to initializing DataManager with Kafka servers: {kafka_servers}") + 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, + "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...") + 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.") + 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.") + 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"DataManager initialized with Kafka servers: {kafka_servers}") self.logger = logger self.notification_handler = notification_handler @@ -57,12 +71,12 @@ class DataManager(): 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}") + self.logger.error(f"Error closing Kafka producer: {e}") else: - self.logger.warning( - "Kafka producer is already closed or not initialized.") + self.logger.warning("Kafka producer is already closed or not initialized.") def __del__(self): self.shutdown() @@ -70,7 +84,8 @@ class DataManager(): 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}") + 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.""" @@ -93,22 +108,22 @@ class DataManager(): 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.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.notification_handler.build_and_send_notification( 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 + attachment_content=trace, ) self.logger.error(trace)