From 8789e6693f9e2900557bfe568b1e87a5aab1bba7 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Fri, 20 Mar 2026 15:27:55 -0300 Subject: [PATCH] SIENTIAPDE-1712 Remove `query_to_minio` method from Storage class and update worker activities to eliminate its usage. This change streamlines the codebase by removing unused functionality related to MinIO queries. --- laborious/activities/storage.py | 69 --------------------------------- laborious/worker/worker.py | 3 -- 2 files changed, 72 deletions(-) diff --git a/laborious/activities/storage.py b/laborious/activities/storage.py index 9fda091..223001b 100644 --- a/laborious/activities/storage.py +++ b/laborious/activities/storage.py @@ -206,75 +206,6 @@ class Storage(Postgres, MinioManager): return report - @activity.defn(name='query_to_minio') - async def query_to_minio(self, input_data: dict[str, Any]) -> dict[str, Any]: - """ - Execute SQL query, write result as Parquet to MinIO, and return object name. - - Args (input_data): - metadata (dict): Workflow metadata - query (str): SQL query - model_name (str): Model name for object naming - object_prefix (str, optional): Prefix inside bucket (default: datasets/retrain) - - Returns: - dict: { success: bool, object_name: str, uri: str } - """ - - if self.minio_repository is None: - raise ValueError('Minio repository not initialized') - - metadata = input_data.get('metadata', {}) - model_name = input_data.get('model_name') or metadata.get('model_name') or 'unknown' - object_prefix = input_data.get('object_prefix', 'datasets/retrain') - - timestamp = now().strftime(DATETIME_FORMAT_FILENAME) - # Keep a stable model-level layout for minimal_retrain: - # training_datasets// - # Sanitize object_prefix to avoid extra subdirectories in the relative key. - safe_prefix = str(object_prefix).strip().strip('/').replace('/', '_') - filename = f'{safe_prefix}_{timestamp}.parquet' - relative_key = f'training_datasets/{model_name}/{filename}' - bucket = getattr(self.minio_repository, 'bucket', 'streamlit-connectors') - uri = f's3://{bucket}/{relative_key}' - - try: - data = await self.load_custom_query(input_data) - if not data: - self.error('query_to_minio failed: No data returned from query', metadata) - return {'success': False, 'message': 'No data returned from query'} - - # Ensure we have a DataFrame - data = pd.DataFrame(data) - - # Convert DataFrame -> parquet bytes, then upload using the new MinIO interface. - parquet_buffer = BytesIO() - data.to_parquet(parquet_buffer, engine='pyarrow', index=True) - file_bytes = parquet_buffer.getvalue() - - upload_result = await self.minio_repository.upload_file( - file_bytes=file_bytes, - relative_key=relative_key, - metadata=metadata, - ) - - object_key_full = upload_result.get('minio_object_name', relative_key) - uri = f's3://{bucket}/{object_key_full}' - return {'success': True, 'object_key': object_key_full, 'uri': uri} - except Exception as e: - trace = traceback.format_exc() - await self.send_notification_async( - metadata=metadata, - notification_id='ERROR_STORING_QUERY_TO_MINIO', - message=f'Error storing query to MinIO: {e}', - block='query_to_minio', - level=NotificationLevel.ERROR, - attachment_content=trace, - ) - - self.error(trace, metadata) - - return {'success': False, 'message': str(e)} def close(self) -> None: """Close Storage resources (MinIO client and Postgres engine).""" diff --git a/laborious/worker/worker.py b/laborious/worker/worker.py index 8c63b34..5202f2d 100644 --- a/laborious/worker/worker.py +++ b/laborious/worker/worker.py @@ -152,9 +152,7 @@ async def main(): main_workflow=MinimalRetrain, other_workflows=[], activities=[ - activities.load_custom_query, activities.load_query_with_minio_offload, - activities.query_to_minio, activities.retrain_model, activities.update_production_model, activities.format_retrain_report, @@ -200,7 +198,6 @@ async def main(): activities.format_transformed_data, activities.format_prediction, activities.format_default_prediction, - activities.get_last_timestamp, # OPC activities.write_opc_data, # Postgres / MinIO offload