From d79f8820187c09480819400e3fee122d16e60776 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Tue, 1 Jul 2025 12:19:06 -0300 Subject: [PATCH] Add MongoDB integration to DataManager, switch Python language server to Jedi, and include pymongo in requirements.txt --- .vscode/settings.json | 2 +- ingestor/app.py | 2 - ingestor/managers/data_manager.py | 78 ++++++++++++++++++++++++++++--- requirements.txt | 1 + 4 files changed, 73 insertions(+), 10 deletions(-) diff --git a/.vscode/settings.json b/.vscode/settings.json index dc84355..e7d55f8 100644 --- a/.vscode/settings.json +++ b/.vscode/settings.json @@ -6,7 +6,7 @@ "connectionId": "sonardev-sientia-ai", "projectKey": "Aignosi_sientia-dataops-opc-ingestor_1642c8bf-a148-4911-8362-0903d7fef99a" }, - "python.languageServer": "Pylance", + "python.languageServer": "Jedi", "python.analysis.typeCheckingMode": "standard", "editor.suggestSelection": "first", "windsurfPyright.disableLanguageServices": true diff --git a/ingestor/app.py b/ingestor/app.py index 75e4bed..d742bf5 100644 --- a/ingestor/app.py +++ b/ingestor/app.py @@ -25,8 +25,6 @@ def main(): exit_signal.set() ingestor.logger.info("Ingestor prepared. Starting main loop.") - print(exit_signal.is_set()) - while not exit_signal.is_set(): start_time = time() # Start loop timer try: diff --git a/ingestor/managers/data_manager.py b/ingestor/managers/data_manager.py index bb2c1e4..466542a 100644 --- a/ingestor/managers/data_manager.py +++ b/ingestor/managers/data_manager.py @@ -1,6 +1,8 @@ import json from logging import Logger from time import sleep +from datetime import datetime, timezone +from pymongo import MongoClient from kafka import KafkaProducer from kafka.errors import NoBrokersAvailable from sientia_do.notifications.handlers import NotificationHandler @@ -14,6 +16,8 @@ class DataManager: def __init__( self, kafka_servers: str, + mongo_connection_string: str, + mongo_database: str, logger: Logger, notification_handler: NotificationHandler, ) -> None: @@ -41,10 +45,12 @@ class DataManager: 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, + 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) + metrics.KAFKA_CONNECTION_STATUS.labels( + pod_id=self.pod_id).set(1) break except NoBrokersAvailable: logger.error( @@ -61,7 +67,34 @@ class DataManager: 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.connection_string = mongo_connection_string + self.database = mongo_database + + for i in range(0, 3): + logger.info( + f"Trying ({i}) to initializing DataManager with MongoDB servers: {self.connection_string}" + ) + try: + self.mongo_client = MongoClient(self.connection_string) + self.mongo_client.server_info() + + self.mongo_db = self.mongo_client[self.database] + break + except Exception as e: + logger.error(f"Error connecting to MongoDB: {e}") + sleep(5) + else: + raise ValueError( + f"Failed to connect to MongoDB servers {self.connection_string} after 3 attempts." + ) + + logger.info( + f"DataManager initialized with MongoDB servers: {self.connection_string}" + ) + self.logger = logger self.notification_handler = notification_handler @@ -72,11 +105,19 @@ class DataManager: 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) + 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.") + 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}") def __del__(self): self.shutdown() @@ -114,10 +155,12 @@ class DataManager: ).add_errback(self.delivery_error) self.kafka_producer.flush(timeout=10) - metrics.KAFKA_MESSAGES_SENT.labels(pod_id=self.pod_id, topic=topic).inc() + 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() + 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}", @@ -127,3 +170,24 @@ class DataManager: attachment_content=trace, ) self.logger.error(trace) + + try: + collection = self.mongo_db[topic] + + collection.insert_one( + { + **data, + "inserted_at": datetime.now(timezone.utc), + } + ) + + except Exception as e: + trace = traceback.format_exc() + self.notification_handler.build_and_send_notification( + 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) diff --git a/requirements.txt b/requirements.txt index 8d5d26c..f67d2f4 100644 --- a/requirements.txt +++ b/requirements.txt @@ -2,3 +2,4 @@ asyncua==1.1.5 redis git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.1.14 prometheus_client +pymongo \ No newline at end of file