diff --git a/laborious/activities/gates.py b/laborious/activities/gates.py index 0b747f9..f9eeb16 100644 --- a/laborious/activities/gates.py +++ b/laborious/activities/gates.py @@ -479,7 +479,8 @@ class Gates(MinioManager): minio_repo=self.minio_repository, model_name=input_data['model_name'], operation='transform', - workflow_metadata=metadata + workflow_metadata=metadata, + last_timestamp=payload.last_timestamp, ) @activity.defn(name='format_prediction') diff --git a/laborious/activities/mlflow.py b/laborious/activities/mlflow.py index e6522f0..ea97f3b 100644 --- a/laborious/activities/mlflow.py +++ b/laborious/activities/mlflow.py @@ -178,6 +178,7 @@ class MLFlow(MinioManager): operation='transform', status=response_data, workflow_metadata=metadata, + last_timestamp=payload.last_timestamp, ) return await MinioDataFramePayload.from_dataframe( @@ -189,6 +190,7 @@ class MLFlow(MinioManager): status={ 'success': True, }, + last_timestamp=payload.last_timestamp, ) @activity.defn(name='request_predict') @@ -260,6 +262,7 @@ class MLFlow(MinioManager): operation='predict', status=response_data, workflow_metadata=metadata, + last_timestamp=payload.last_timestamp, ) return await MinioDataFramePayload.from_dataframe( @@ -271,6 +274,7 @@ class MLFlow(MinioManager): status={ 'success': True, }, + last_timestamp=payload.last_timestamp, ) @activity.defn(name='retrain_model') diff --git a/laborious/utils/models/minio_dataframe_payload.py b/laborious/utils/models/minio_dataframe_payload.py index c49218d..7e0d865 100644 --- a/laborious/utils/models/minio_dataframe_payload.py +++ b/laborious/utils/models/minio_dataframe_payload.py @@ -173,6 +173,7 @@ class MinioDataFramePayload: operation: OperationKind, status: dict[str, Any] | None = None, workflow_metadata: dict | None = None, + last_timestamp: str | None = None, ) -> 'MinioDataFramePayload': """ Evaluate the DataFrame size, then either inline dict or upload parquet to MinIO. @@ -199,7 +200,9 @@ class MinioDataFramePayload: data=None, last_timestamp=now().strftime(DATETIME_FORMAT_WITH_TZ), status=status ) - last_timestamp = max(dataframe['timestamp'].values.tolist()) + + if last_timestamp is None: + last_timestamp = max(dataframe['timestamp'].values.tolist()) if cls.estimate_size_bytes(dataframe) <= OFFLOAD_THRESHOLD_BYTES: return cls(data=dataframe.to_dict(), last_timestamp=last_timestamp)