Files
sientia-dataops-scouter_tem…/scouter/activities/mongodb.py
vitor-aignosi 97eb5bc904 SIENTIAPDE-1316
Update .gitignore and refactor metrics.py, activities.py, and gates.py for improved clarity and consistency. Added coverage.xml and cache directories to .gitignore. Standardized string formatting and parameter handling in metrics and activities classes, enhancing code readability. Removed the deprecated faker.py file and adjusted related tests accordingly.
2025-10-16 13:31:12 -03:00

196 lines
6.7 KiB
Python

from collections.abc import Hashable
from datetime import UTC
from temporalio import activity, workflow
with workflow.unsafe.imports_passed_through():
import traceback
from datetime import datetime
from typing import Any
from pandas import DataFrame
from pymongo import MongoClient
from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler
from sientia_do.notifications.models import NotificationLevel
from sientia_do.observability.logger import Logger
from sientia_do.temporal.activities.base import BaseActivity
from sientia_do.temporal.constants import DATETIME_FORMAT_MS_WITH_TZ
def clear_mongo_id(docs: list) -> list:
"""
Remove MongoDB internal `_id` fields from documents.
This utility function recursively removes the MongoDB `_id` field from
documents and nested structures. It's used to clean data before
processing or export operations.
Args:
docs (list): List of documents to clean
Returns:
list: Documents with `_id` fields removed
Note:
This function modifies the input list in-place and returns the same reference
"""
for doc in docs:
if isinstance(doc, list):
clear_mongo_id(doc)
elif isinstance(doc, dict):
if '_id' in doc:
del doc['_id']
for _key, value in doc.items():
if isinstance(value, list):
clear_mongo_id(value)
elif isinstance(value, dict):
clear_mongo_id([value])
return docs
class MongoDB(BaseActivity):
"""
MongoDB operations for data retrieval and storage.
This class provides MongoDB connectivity and operations for the Scouter system.
It handles:
- Connection management with automatic reconnection
- Data retrieval with timestamp-based filtering
- Document cleaning and preprocessing
- Error handling and notification integration
The class implements Temporal activities for MongoDB operations, enabling
distributed data processing with fault tolerance and monitoring.
"""
def __init__(
self,
connection_string: str,
database_name: str,
logger: Logger,
notification_handler: NotificationHandler,
):
"""
Initialize MongoDB connection and services.
Args:
connection_string (str): MongoDB connection URI string
database_name (str): Name of the target database
logger (Logger): Logger instance for operation logging
notification_handler (NotificationHandler): Handler for system notifications
Raises:
ConnectionError: If MongoDB connection fails
"""
self.connection_string = connection_string
self.database_name = database_name
self.client: MongoClient = MongoClient(
self.connection_string, serverSelectionTimeoutMS=5000
)
self.client.server_info() # Trigger an exception if connection fails
self.database = self.client[self.database_name]
# Initialize MongoDB client here (omitted for brevity)
logger.info('MongoDB connection initialized')
BaseActivity.__init__(
self, logger=logger, notification_handler=notification_handler, set_error_counter=True
)
def shutdown(self):
"""
Gracefully close MongoDB client connection.
This method ensures proper cleanup of MongoDB connections to prevent
connection leaks and ensure graceful application termination.
"""
try:
if self.client:
self.logger.info('Closing MongoDB connection...')
self.client.close()
self.logger.info('MongoDB connection closed successfully')
except Exception as e:
self.logger.error(f'Failed to close MongoDB connection: {e}')
def __del__(self):
"""
Destructor to ensure MongoDB client is closed.
This destructor ensures that MongoDB connections are properly closed
when the object is garbage collected, preventing resource leaks.
"""
self.shutdown()
@activity.defn(name='load_latest_data')
async def load_latest_data(self, input_data: dict[str, Any]) -> dict[Hashable, Any]:
"""
Load the latest data from MongoDB collection since a specified timestamp.
This activity retrieves data from a MongoDB collection, optionally
filtering by timestamp to enable incremental data processing. It
handles connection management and provides comprehensive error reporting.
Args:
input_data (dict[str, Any]): Activity input parameters.
Required fields:
- metadata (dict[str, Any]): Workflow execution metadata
- collection_name (str): Name of the MongoDB collection
- last_data_timestamp (str | None): Last processed timestamp for filtering
Returns:
dict[str, Any]: Retrieved data, or empty dict if no data found
Raises:
Exception: If MongoDB operation fails
"""
metadata = input_data['metadata']
collection_name = input_data['collection_name']
last_data_timestamp = input_data['last_data_timestamp']
self.info(f'Loading data from MongoDB: {input_data}', metadata=metadata)
try:
if last_data_timestamp is None:
data_filter = {}
else:
data_filter = {
'inserted_at': {
'$gt': datetime.strptime(last_data_timestamp, DATETIME_FORMAT_MS_WITH_TZ)
}
}
self.debug(f'Data filter: {data_filter}', metadata=metadata)
data = list(self.database[collection_name].find(data_filter, {'_id': 0}))
data = clear_mongo_id(data)
self.debug(f'Collected: {data}', metadata=metadata)
for item in data:
item['inserted_at'] = (
item['inserted_at'].replace(tzinfo=UTC).strftime(DATETIME_FORMAT_MS_WITH_TZ)
)
self.info(f'Loaded {len(data)} documents from MongoDB', metadata=metadata)
self.debug(f'Loaded data: {data}', metadata=metadata)
return DataFrame(data).to_dict()
except Exception as e:
trace = traceback.format_exc()
self.send_notification(
metadata=metadata,
notification_id='MONGO_LOAD_ERROR',
message=f'Error loading data from MongoDB: {e}',
block='load_latest_data',
level=NotificationLevel.ERROR,
attachment_content=trace,
)
raise e