diff --git a/laborious/utils/models/minio_dataframe_payload.py b/laborious/utils/models/minio_dataframe_payload.py index bb3f2fd..30f5a00 100644 --- a/laborious/utils/models/minio_dataframe_payload.py +++ b/laborious/utils/models/minio_dataframe_payload.py @@ -20,6 +20,7 @@ from os import getenv from typing import Any, Literal from pandas import DataFrame, read_parquet +from sientia_do.observability.logger import Logger from sientia_do.repository.minio_repository import MinioRepository from sientia_do.temporal.constants import DATETIME_FORMAT_FILENAME, DATETIME_FORMAT_WITH_TZ, now @@ -81,6 +82,24 @@ class MinioDataFramePayload: object_prefix: str | None = None uri: str | None = None + @staticmethod + def _debug( + logger: Logger | None, + message: str, + metadata: dict[str, Any] | None = None, + ) -> None: + """ + Emit debug logs only when logger is provided + + Args: + - logger (Logger | None): Logger instance used for debug messages + - message (str): Message to be logged + - metadata (dict[str, Any] | None): Optional workflow metadata context + """ + if logger is None: + return + logger.debug(message, metadata) + @classmethod def from_dict(cls, raw: 'dict[str, Any] | MinioDataFramePayload') -> 'MinioDataFramePayload': """ @@ -114,7 +133,11 @@ class MinioDataFramePayload: ) @staticmethod - def estimate_size_bytes(df: DataFrame) -> int: + def estimate_size_bytes( + df: DataFrame, + metadata: dict[str, Any] | None = None, + logger: Logger | None = None, + ) -> int: """ Approximate serialized size of the DataFrame as the default-orient dict. @@ -125,9 +148,13 @@ class MinioDataFramePayload: int: Estimated size in bytes (pickle of dict representation). """ try: - return len(pickle.dumps(df.to_dict())) + size = len(pickle.dumps(df.to_dict())) except Exception: - return len(pickle.dumps(df)) + size = len(pickle.dumps(df)) + + if logger is not None: + logger.debug(f'DataFrame size: {size} bytes', metadata) + return size @staticmethod def parse_object_timestamp(object_key: str) -> datetime | None: @@ -174,6 +201,7 @@ class MinioDataFramePayload: status: dict[str, Any] | None = None, workflow_metadata: dict | None = None, last_timestamp: str | None = None, + logger: Logger | None = None, ) -> 'MinioDataFramePayload': """ Evaluate the DataFrame size, then either inline dict or upload parquet to MinIO. @@ -196,6 +224,11 @@ class MinioDataFramePayload: """ if dataframe is None or dataframe.empty: + cls._debug( + logger, + 'MinioDataFramePayload.from_dataframe received empty dataframe, returning empty payload', + workflow_metadata, + ) return cls( data=None, last_timestamp=now().strftime(DATETIME_FORMAT_WITH_TZ), status=status ) @@ -203,11 +236,34 @@ class MinioDataFramePayload: if last_timestamp is None: last_timestamp = max(dataframe['timestamp'].values.tolist()) - if cls.estimate_size_bytes(dataframe) <= OFFLOAD_THRESHOLD_BYTES: + dataframe_size = cls.estimate_size_bytes(dataframe) + cls._debug( + logger, + ( + f'MinioDataFramePayload.from_dataframe estimated size: {dataframe_size} bytes ' + f'(threshold: {OFFLOAD_THRESHOLD_BYTES} bytes)' + ), + workflow_metadata, + ) + + if dataframe_size <= OFFLOAD_THRESHOLD_BYTES: + cls._debug( + logger, + 'MinioDataFramePayload.from_dataframe using inline payload', + workflow_metadata, + ) return cls(data=dataframe.to_dict(), last_timestamp=last_timestamp, status=status) timestamp = now().strftime(DATETIME_FORMAT_FILENAME) object_key, object_prefix = _build_object_key(model_name, operation, timestamp) + cls._debug( + logger, + ( + 'MinioDataFramePayload.from_dataframe offloading payload to MinIO ' + f'with key {object_key}' + ), + workflow_metadata, + ) # Upload using the relative object key. The upstream repository will # prefix it internally under its MinIO namespace. @@ -224,6 +280,11 @@ class MinioDataFramePayload: bucket = minio_repo.bucket object_key_full = upload_result.get('minio_object_name', object_key) uri = f's3://{bucket}/{object_key_full}' if bucket else None + cls._debug( + logger, + f'MinioDataFramePayload.from_dataframe upload completed: {uri}', + workflow_metadata, + ) return cls( data=None, @@ -236,7 +297,10 @@ class MinioDataFramePayload: ) async def retrieve( - self, minio_repo: MinioRepository, workflow_metadata: dict[str, Any] | None = None + self, + minio_repo: MinioRepository, + workflow_metadata: dict[str, Any] | None = None, + logger: Logger | None = None, ) -> DataFrame: """ Load parquet from MinIO when object_key is set and populate inline data. @@ -249,13 +313,33 @@ class MinioDataFramePayload: dict[str, Any]: Flat dict with data filled (same keys as to_dict after load). """ if self.data is not None: + self._debug( + logger, + 'MinioDataFramePayload.retrieve using inline payload data', + workflow_metadata, + ) return DataFrame(self.data) if not self.has_data(): + self._debug( + logger, + 'MinioDataFramePayload.retrieve found no payload data, returning empty dataframe', + workflow_metadata, + ) return DataFrame() + self._debug( + logger, + f'MinioDataFramePayload.retrieve downloading object from MinIO: {self.object_key}', + workflow_metadata, + ) file_bytes = await minio_repo.download_file( object_name=self.object_key, metadata=workflow_metadata ) df = read_parquet(BytesIO(file_bytes)) + self._debug( + logger, + f'MinioDataFramePayload.retrieve loaded dataframe from MinIO with shape {df.shape}', + workflow_metadata, + ) return df