Files
sientia-dataops-opc-ingestor/ingestor/managers/data_manager.py
vitor-aignosi 28866ac9c8 Update logging in DataManager to use debug level for message skipping and insertion confirmation
- Changed log level from info to debug for skipping messages to Kafka topic.
- Added debug log for successful message insertion into MongoDB collection.
2025-07-03 14:13:00 -03:00

205 lines
7.5 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
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}")
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
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}")
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)
else:
self.logger.debug(f"Skipping message to topic {topic}: {data}")
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)