Add MongoDB integration to DataManager, switch Python language server to Jedi, and include pymongo in requirements.txt

This commit is contained in:
vitor-aignosi
2025-07-01 12:19:06 -03:00
parent 0d03141655
commit d79f882018
4 changed files with 73 additions and 10 deletions

View File

@@ -6,7 +6,7 @@
"connectionId": "sonardev-sientia-ai", "connectionId": "sonardev-sientia-ai",
"projectKey": "Aignosi_sientia-dataops-opc-ingestor_1642c8bf-a148-4911-8362-0903d7fef99a" "projectKey": "Aignosi_sientia-dataops-opc-ingestor_1642c8bf-a148-4911-8362-0903d7fef99a"
}, },
"python.languageServer": "Pylance", "python.languageServer": "Jedi",
"python.analysis.typeCheckingMode": "standard", "python.analysis.typeCheckingMode": "standard",
"editor.suggestSelection": "first", "editor.suggestSelection": "first",
"windsurfPyright.disableLanguageServices": true "windsurfPyright.disableLanguageServices": true

View File

@@ -25,8 +25,6 @@ def main():
exit_signal.set() exit_signal.set()
ingestor.logger.info("Ingestor prepared. Starting main loop.") ingestor.logger.info("Ingestor prepared. Starting main loop.")
print(exit_signal.is_set())
while not exit_signal.is_set(): while not exit_signal.is_set():
start_time = time() # Start loop timer start_time = time() # Start loop timer
try: try:

View File

@@ -1,6 +1,8 @@
import json import json
from logging import Logger from logging import Logger
from time import sleep from time import sleep
from datetime import datetime, timezone
from pymongo import MongoClient
from kafka import KafkaProducer from kafka import KafkaProducer
from kafka.errors import NoBrokersAvailable from kafka.errors import NoBrokersAvailable
from sientia_do.notifications.handlers import NotificationHandler from sientia_do.notifications.handlers import NotificationHandler
@@ -14,6 +16,8 @@ class DataManager:
def __init__( def __init__(
self, self,
kafka_servers: str, kafka_servers: str,
mongo_connection_string: str,
mongo_database: str,
logger: Logger, logger: Logger,
notification_handler: NotificationHandler, notification_handler: NotificationHandler,
) -> None: ) -> None:
@@ -41,10 +45,12 @@ class DataManager:
value_serializer=lambda v: json.dumps(v).encode( value_serializer=lambda v: json.dumps(v).encode(
"utf-8" "utf-8"
), # Serialize JSON messages ), # 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 # 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 break
except NoBrokersAvailable: except NoBrokersAvailable:
logger.error( logger.error(
@@ -61,7 +67,34 @@ class DataManager:
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.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.logger = logger
self.notification_handler = notification_handler self.notification_handler = notification_handler
@@ -72,11 +105,19 @@ class DataManager:
self.kafka_producer.flush(timeout=10) self.kafka_producer.flush(timeout=10)
self.kafka_producer.close() self.kafka_producer.close()
# Mark as disconnected # 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: except Exception as e:
self.logger.error(f"Error closing Kafka producer: {e}") self.logger.error(f"Error closing Kafka producer: {e}")
else: 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): def __del__(self):
self.shutdown() self.shutdown()
@@ -114,10 +155,12 @@ class DataManager:
).add_errback(self.delivery_error) ).add_errback(self.delivery_error)
self.kafka_producer.flush(timeout=10) 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: 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() trace = traceback.format_exc()
self.notification_handler.build_and_send_notification( self.notification_handler.build_and_send_notification(
notification_id=f"KAFKA_PRODUCER_ERROR_{topic}", notification_id=f"KAFKA_PRODUCER_ERROR_{topic}",
@@ -127,3 +170,24 @@ class DataManager:
attachment_content=trace, attachment_content=trace,
) )
self.logger.error(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)

View File

@@ -2,3 +2,4 @@ asyncua==1.1.5
redis redis
git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.1.14 git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.1.14
prometheus_client prometheus_client
pymongo