diff --git a/laborious/activities/gates.py b/laborious/activities/gates.py index 34c26f1..cbebebe 100644 --- a/laborious/activities/gates.py +++ b/laborious/activities/gates.py @@ -8,6 +8,7 @@ with workflow.unsafe.imports_passed_through(): from sientia_do.temporal.activities.base import BaseActivity from sientia_do.observability.logger import Logger from sientia_do.temporal.constants import DATETIME_FORMAT_WITH_TZ, now + from sientia_do.formatters import create_sample_dict from laborious.utils.filters.mlflow_filters import nan_values_filter, api_error_filter from typing import Any from laborious.utils.filters.conditional_filters import ( @@ -122,15 +123,13 @@ class Gates(BaseActivity): self.info("Performing input gate...", metadata) - self.debug(f"Input data: {input_data}", metadata) - filters = input_data['filters'] data = DataFrame(input_data['data']) path_priority = input_data['path_priority'] filter_output = [] - self.debug(f"Input data:\n {data}", metadata) + self.debug(f"Input data: {data.head(5).to_string()}", metadata) self.debug(f"Filters: {filters}", metadata) # Apply each configured filter @@ -207,7 +206,8 @@ class Gates(BaseActivity): filter_output = [] - self.debug(f"Input data:\n {data}", metadata) + self.debug(create_sample_dict( + data, max_items=5, max_depth=2), metadata) self.debug(f"Filters: {filters}", metadata) comments = [] @@ -292,8 +292,8 @@ class Gates(BaseActivity): filter_output = [] - self.debug(f"Input data:\n {data}", metadata) - self.debug(f"Filters: {filters}", metadata) + self.debug(f"Input data:\n {data.head(5).to_string()}", metadata) + self.debug(create_sample_dict(filters), metadata) for fil, config in filters.items(): if fil not in mlflow_content_filter_functions: @@ -410,7 +410,7 @@ class Gates(BaseActivity): self.debug( f"Prediction store policy: {prediction_store_policy}", metadata) - self.debug(f"Prediction data: {data.to_string()}", metadata) + self.debug(f"Prediction data: {data.head(5).to_string()}", metadata) policy_type, policy_value = self.get_prediction_store_policy( prediction_store_policy, metadata) @@ -445,7 +445,7 @@ class Gates(BaseActivity): data = data.reset_index(drop=True) self.info(f"Prediction formatted: {len(data)} rows", metadata) - self.debug(f"Prediction data: {data.to_string()}", metadata) + self.debug(f"Prediction data: {data.head(5).to_string()}", metadata) return data.to_dict() @@ -521,7 +521,7 @@ class Gates(BaseActivity): data = DataFrame(input_data['data']) - self.debug(f"Input data: {data.to_string()}", metadata) + self.debug(f"Input data: {data.head(5).to_string()}", metadata) if data.empty: return now().strftime(DATETIME_FORMAT_WITH_TZ) diff --git a/laborious/activities/mlflow.py b/laborious/activities/mlflow.py index c11bbaa..5b449e0 100644 --- a/laborious/activities/mlflow.py +++ b/laborious/activities/mlflow.py @@ -3,13 +3,13 @@ from temporalio import activity, workflow with workflow.unsafe.imports_passed_through(): from datetime import datetime - import json from pandas import Timestamp, to_datetime from sientia_do.temporal.constants import DATETIME_FORMAT, DATETIME_FORMAT_WITH_TZ from sientia_do.temporal.activities.base import BaseActivity 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.formatters import create_sample_dict from laborious.utils.repository.model_repository import MLFlowRepository from typing import Any import numpy as np @@ -17,39 +17,6 @@ with workflow.unsafe.imports_passed_through(): import traceback -def create_sample_dict(data: dict, max_items: int = 3, max_depth: int = 2) -> dict: - """ - Create a sample of a dictionary for debugging purposes. - - Args: - data: Dictionary to sample - max_items: Maximum number of items to show per level - max_depth: Maximum depth to traverse nested structures - - Returns: - Dictionary with sampled content - """ - if max_depth <= 0: - return {"...": "max_depth_reached"} - - sample = {} - items = list(data.items())[:max_items] - - for key, value in items: - if isinstance(value, dict): - sample[key] = create_sample_dict(value, max_items, max_depth - 1) - elif isinstance(value, list): - sample[key] = value[:max_items] if len( - value) > max_items else value - else: - sample[key] = value - - if len(data) > max_items: - sample["..."] = f"({len(data) - max_items} more items)" - - return sample - - class MLFlow(BaseActivity): """ MLFlow integration activities for model inference operations. diff --git a/requirements.txt b/requirements.txt index 10995c8..1ce9f9c 100644 --- a/requirements.txt +++ b/requirements.txt @@ -3,7 +3,6 @@ psycopg2-binary sqlalchemy asyncua redis -git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.4.5 -# git+ssh://git@github.com/Aignosi/sientia-mlops-library.git@0.38.13 -/home/grezewave/Documents/projects/sientia/sientia-mlops-library +git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.4.6 +git+ssh://git@github.com/Aignosi/sientia-mlops-library.git@0.39.0 prometheus-client