SIENTIAPDE-1712
SIENTIAPDE-1712 Implement debug logging in MinioDataFramePayload class for enhanced traceability. Added a static method for conditional logging and integrated debug statements throughout methods to capture DataFrame size estimates, upload actions, and retrieval processes, improving overall observability.
This commit is contained in:
@@ -20,6 +20,7 @@ from os import getenv
|
|||||||
from typing import Any, Literal
|
from typing import Any, Literal
|
||||||
|
|
||||||
from pandas import DataFrame, read_parquet
|
from pandas import DataFrame, read_parquet
|
||||||
|
from sientia_do.observability.logger import Logger
|
||||||
from sientia_do.repository.minio_repository import MinioRepository
|
from sientia_do.repository.minio_repository import MinioRepository
|
||||||
from sientia_do.temporal.constants import DATETIME_FORMAT_FILENAME, DATETIME_FORMAT_WITH_TZ, now
|
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
|
object_prefix: str | None = None
|
||||||
uri: 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
|
@classmethod
|
||||||
def from_dict(cls, raw: 'dict[str, Any] | MinioDataFramePayload') -> 'MinioDataFramePayload':
|
def from_dict(cls, raw: 'dict[str, Any] | MinioDataFramePayload') -> 'MinioDataFramePayload':
|
||||||
"""
|
"""
|
||||||
@@ -114,7 +133,11 @@ class MinioDataFramePayload:
|
|||||||
)
|
)
|
||||||
|
|
||||||
@staticmethod
|
@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.
|
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).
|
int: Estimated size in bytes (pickle of dict representation).
|
||||||
"""
|
"""
|
||||||
try:
|
try:
|
||||||
return len(pickle.dumps(df.to_dict()))
|
size = len(pickle.dumps(df.to_dict()))
|
||||||
except Exception:
|
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
|
@staticmethod
|
||||||
def parse_object_timestamp(object_key: str) -> datetime | None:
|
def parse_object_timestamp(object_key: str) -> datetime | None:
|
||||||
@@ -174,6 +201,7 @@ class MinioDataFramePayload:
|
|||||||
status: dict[str, Any] | None = None,
|
status: dict[str, Any] | None = None,
|
||||||
workflow_metadata: dict | None = None,
|
workflow_metadata: dict | None = None,
|
||||||
last_timestamp: str | None = None,
|
last_timestamp: str | None = None,
|
||||||
|
logger: Logger | None = None,
|
||||||
) -> 'MinioDataFramePayload':
|
) -> 'MinioDataFramePayload':
|
||||||
"""
|
"""
|
||||||
Evaluate the DataFrame size, then either inline dict or upload parquet to MinIO.
|
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:
|
if dataframe is None or dataframe.empty:
|
||||||
|
cls._debug(
|
||||||
|
logger,
|
||||||
|
'MinioDataFramePayload.from_dataframe received empty dataframe, returning empty payload',
|
||||||
|
workflow_metadata,
|
||||||
|
)
|
||||||
return cls(
|
return cls(
|
||||||
data=None, last_timestamp=now().strftime(DATETIME_FORMAT_WITH_TZ), status=status
|
data=None, last_timestamp=now().strftime(DATETIME_FORMAT_WITH_TZ), status=status
|
||||||
)
|
)
|
||||||
@@ -203,11 +236,34 @@ class MinioDataFramePayload:
|
|||||||
if last_timestamp is None:
|
if last_timestamp is None:
|
||||||
last_timestamp = max(dataframe['timestamp'].values.tolist())
|
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)
|
return cls(data=dataframe.to_dict(), last_timestamp=last_timestamp, status=status)
|
||||||
|
|
||||||
timestamp = now().strftime(DATETIME_FORMAT_FILENAME)
|
timestamp = now().strftime(DATETIME_FORMAT_FILENAME)
|
||||||
object_key, object_prefix = _build_object_key(model_name, operation, timestamp)
|
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
|
# Upload using the relative object key. The upstream repository will
|
||||||
# prefix it internally under its MinIO namespace.
|
# prefix it internally under its MinIO namespace.
|
||||||
@@ -224,6 +280,11 @@ class MinioDataFramePayload:
|
|||||||
bucket = minio_repo.bucket
|
bucket = minio_repo.bucket
|
||||||
object_key_full = upload_result.get('minio_object_name', object_key)
|
object_key_full = upload_result.get('minio_object_name', object_key)
|
||||||
uri = f's3://{bucket}/{object_key_full}' if bucket else None
|
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(
|
return cls(
|
||||||
data=None,
|
data=None,
|
||||||
@@ -236,7 +297,10 @@ class MinioDataFramePayload:
|
|||||||
)
|
)
|
||||||
|
|
||||||
async def retrieve(
|
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:
|
) -> DataFrame:
|
||||||
"""
|
"""
|
||||||
Load parquet from MinIO when object_key is set and populate inline data.
|
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).
|
dict[str, Any]: Flat dict with data filled (same keys as to_dict after load).
|
||||||
"""
|
"""
|
||||||
if self.data is not None:
|
if self.data is not None:
|
||||||
|
self._debug(
|
||||||
|
logger,
|
||||||
|
'MinioDataFramePayload.retrieve using inline payload data',
|
||||||
|
workflow_metadata,
|
||||||
|
)
|
||||||
return DataFrame(self.data)
|
return DataFrame(self.data)
|
||||||
|
|
||||||
if not self.has_data():
|
if not self.has_data():
|
||||||
|
self._debug(
|
||||||
|
logger,
|
||||||
|
'MinioDataFramePayload.retrieve found no payload data, returning empty dataframe',
|
||||||
|
workflow_metadata,
|
||||||
|
)
|
||||||
return DataFrame()
|
return DataFrame()
|
||||||
|
|
||||||
|
self._debug(
|
||||||
|
logger,
|
||||||
|
f'MinioDataFramePayload.retrieve downloading object from MinIO: {self.object_key}',
|
||||||
|
workflow_metadata,
|
||||||
|
)
|
||||||
file_bytes = await minio_repo.download_file(
|
file_bytes = await minio_repo.download_file(
|
||||||
object_name=self.object_key, metadata=workflow_metadata
|
object_name=self.object_key, metadata=workflow_metadata
|
||||||
)
|
)
|
||||||
df = read_parquet(BytesIO(file_bytes))
|
df = read_parquet(BytesIO(file_bytes))
|
||||||
|
self._debug(
|
||||||
|
logger,
|
||||||
|
f'MinioDataFramePayload.retrieve loaded dataframe from MinIO with shape {df.shape}',
|
||||||
|
workflow_metadata,
|
||||||
|
)
|
||||||
return df
|
return df
|
||||||
|
|||||||
Reference in New Issue
Block a user