Files
sientia-dataops-opc-ingestor/ingestor/managers/data_manager.py
vitor-aignosi 0885544bc7 SIENTIAPDE-1163
Update values.yaml for image tag and GITHUB_BRANCH

- Updated image tag to "0.2.4" for the latest version.
- Changed GITHUB_BRANCH value to "SIENTIAPDE-1163-alterar-dinamica-de-notificacoes-para-usar-o-mongodb-ao-inves-do-kafka".
- Removed .vscode/settings.json as it is no longer needed.
2025-07-16 10:16:02 -03:00

198 lines
7.2 KiB
Python

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
from sientia_do.notifications.models import NotificationLevel
from sientia_do.
import traceback
import ingestor.metrics as metrics
import os
class DataManager():
def __init__(
self,
kafka_servers: str,
mongo_connection_string: str,
mongo_database: str,
export_to_kafka: bool,
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.mongo_db = self.mongo_client[self.database]
logger.info(
f"DataManager initialized with MongoDB servers: {self.connection_string}"
)
self.logger = logger
self.notification_handler = notification_handler
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.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,
)
self.logger.error(trace)
try:
collection = self.mongo_db[topic]
collection.insert_one(
{
**data,
"inserted_at": datetime.now(timezone.utc),
}
)
self.logger.debug(
f"Message inserted into MongoDB collection {topic}: {data}")
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)