diff --git a/requirements.txt b/requirements.txt index 8589717..ce4f826 100644 --- a/requirements.txt +++ b/requirements.txt @@ -5,6 +5,6 @@ asyncua redis aiokafka pymongo -git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.4.7 +/home/grezewave/Documents/projects/sientia/sientia-dataops-library pydruid[pandas] prometheus-client \ No newline at end of file diff --git a/scouter/activities/gates.py b/scouter/activities/gates.py index ca9e2ce..cb9ebf6 100644 --- a/scouter/activities/gates.py +++ b/scouter/activities/gates.py @@ -10,7 +10,7 @@ with workflow.unsafe.imports_passed_through(): from sientia_do.notifications.handlers import 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.observability.sientia_monitoring import SientiaMonitoring from scouter import metrics from scouter.utils.quality.filters import null_values_filter, out_of_bounds_filter @@ -21,7 +21,7 @@ quality_gate_filters = { } -class Gates(BaseActivity): +class Gates(SientiaMonitoring): """ Data quality gates and filtering operations. @@ -44,7 +44,13 @@ class Gates(BaseActivity): logger (Logger): Logger instance for operation logging notification_handler (NotificationHandler): Handler for system notifications """ - BaseActivity.__init__(self, logger, notification_handler, set_error_counter=True) + SientiaMonitoring.__init__(self, logger, notification_handler) + + def close(self): + """ + Close the Gates class. + """ + SientiaMonitoring.close(self) def apply_aggregation( self, values: DataFrame, aggr_function: str, metadata: dict[str, Any] diff --git a/scouter/activities/mongodb.py b/scouter/activities/mongodb.py index a20e356..586d1a6 100644 --- a/scouter/activities/mongodb.py +++ b/scouter/activities/mongodb.py @@ -13,45 +13,13 @@ with workflow.unsafe.imports_passed_through(): 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.repository.mongodb_repository import MongoDBRepository + from sientia_do.observability.sientia_monitoring import SientiaMonitoring 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): +class MongoDB(SientiaMonitoring): """ MongoDB operations for data retrieval and storage. @@ -85,37 +53,22 @@ class MongoDB(BaseActivity): 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 + + self.mongodb_repository = MongoDBRepository( + connection_string=connection_string, + database_name=database_name, + logger=logger, + notification_handler=notification_handler, ) - def shutdown(self): - """ - Gracefully close MongoDB client connection. + SientiaMonitoring.__init__(self, logger=logger, notification_handler=notification_handler) - This method ensures proper cleanup of MongoDB connections to prevent - connection leaks and ensure graceful application termination. + def close(self): """ - 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}') + Close the MongoDB connection. + """ + self.mongodb_repository.close() + SientiaMonitoring.shutdown(self) def __del__(self): """ @@ -124,7 +77,7 @@ class MongoDB(BaseActivity): This destructor ensures that MongoDB connections are properly closed when the object is garbage collected, preventing resource leaks. """ - self.shutdown() + self.close() @activity.defn(name='load_latest_data') async def load_latest_data(self, input_data: dict[str, Any]) -> dict[Hashable, Any]: @@ -166,9 +119,11 @@ class MongoDB(BaseActivity): 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) + data = self.mongodb_repository.find( + collection_name=collection_name, + filters=data_filter, + metadata=metadata, + ) self.debug(f'Collected: {data}', metadata=metadata) diff --git a/scouter/activities/redis.py b/scouter/activities/redis.py index 9c76fe4..e4513a2 100644 --- a/scouter/activities/redis.py +++ b/scouter/activities/redis.py @@ -10,11 +10,12 @@ with workflow.unsafe.imports_passed_through(): 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.redis_base import Redis as RedisBase + from sientia_do.observability.sientia_monitoring import SientiaMonitoring + from sientia_do.repository.redis_repository import RedisRepository from sientia_do.temporal.constants import DATETIME_FORMAT, now -class Redis(RedisBase): +class Redis(SientiaMonitoring): """ Redis operations for data caching and temporary storage. @@ -49,7 +50,22 @@ class Redis(RedisBase): logger (Logger): Logger instance for operation logging notification_handler (NotificationHandler): Handler for system notifications """ - RedisBase.__init__(self, host, port, username, password, logger, notification_handler) + SientiaMonitoring.__init__(self, logger, notification_handler) + self.redis_repository = RedisRepository( + host=host, + port=port, + username=username, + password=password, + logger=logger, + notification_handler=notification_handler, + ) + + def close(self): + """ + Close the Redis connection. + """ + self.redis_repository.close() + SientiaMonitoring.shutdown(self) @activity.defn(name='get_last_data_timestamp') async def get_last_data_timestamp(self, input_data: dict[str, Any]) -> str | None: @@ -79,7 +95,7 @@ class Redis(RedisBase): self.info(f'Getting last data timestamp for {key}') try: - data_hold = self.get(key) + data_hold = self.redis_repository.get(key) except Exception as e: self.send_notification( metadata=metadata, @@ -137,7 +153,7 @@ class Redis(RedisBase): self.info(f'Last collected timestamp to insert: {last_data_timestamp}', metadata=metadata) try: - self.set(key, last_data_timestamp, ttl=60 * 60 * 5) + self.redis_repository.set(key, last_data_timestamp, ttl=60 * 60 * 5) except Exception as e: self.send_notification( metadata=metadata, @@ -190,7 +206,7 @@ class Redis(RedisBase): self.info(f'Getting held data for {key}') try: - data_hold = self.get(key) + data_hold = self.redis_repository.get(key) except Exception as e: self.send_notification( metadata=metadata, @@ -234,7 +250,7 @@ class Redis(RedisBase): data['timestamp'].max() if not data.empty else data_hold['timestamp'] ) - self.set(key, data_hold, ttl=retention_time) + self.redis_repository.set(key, data_hold, ttl=retention_time) data_hold_df = DataFrame(data_hold, index=[0]) data_hold_melted = data_hold_df.melt( @@ -280,7 +296,7 @@ class Redis(RedisBase): cache = {'data': data.to_dict(), 'held_data': held_data.to_dict()} try: - self.set(key, cache, ttl=120) + self.redis_repository.set(key, cache, ttl=120) except Exception as e: self.send_notification( metadata=metadata,