From 1695df70aea186d6f8a5570af408fdd35f3b21e4 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Wed, 5 Nov 2025 16:30:51 -0300 Subject: [PATCH] SIENTIAPDE-1325 Enhance metrics handling in Gates and MLFlow classes - Added checks for `None` response times before emitting OPC writing metrics in the Gates class to prevent unnecessary metric emissions. - Updated the MLFlow class to conditionally sort and drop duplicates based on the presence of the 'created_at' column, ensuring robustness in data processing. - Adjusted corresponding tests to validate the new behavior in both classes. --- laborious/activities/gates.py | 45 ++++++++++--------- laborious/activities/mlflow.py | 9 ++-- tests/laborious/activities/test_mlflow.py | 10 +++-- .../utils/repository/test_minio_repository.py | 15 +++++++ .../utils/repository/test_model_repository.py | 10 +++++ 5 files changed, 60 insertions(+), 29 deletions(-) diff --git a/laborious/activities/gates.py b/laborious/activities/gates.py index bc5ea8e..9dc9eb1 100644 --- a/laborious/activities/gates.py +++ b/laborious/activities/gates.py @@ -656,28 +656,29 @@ class Gates(SientiaMonitoring): for server_id, tags in opc_metrics.items(): for tag, response_time in tags.items(): - await self.emit_metric( - metric_object=metrics.PREDICTION_OPC_WRITING_RESPONSE_TIME_MONITOR, - method='observe', - tags={ - 'pod_id': self.pod_id, - 'model_name': metadata['model_name'], - 'workflow_name': metadata['workflow_name'], - 'opc_server_id': server_id, - 'tag': tag, - }, - value=response_time, - ) + if response_time is not None: + await self.emit_metric( + metric_object=metrics.PREDICTION_OPC_WRITING_RESPONSE_TIME_MONITOR, + method='observe', + tags={ + 'pod_id': self.pod_id, + 'model_name': metadata['model_name'], + 'workflow_name': metadata['workflow_name'], + 'opc_server_id': server_id, + 'tag': tag, + }, + value=response_time, + ) - await self.emit_metric( - metric_object=metrics.PREDICTION_OPC_WRITING_COUNT, - tags={ - 'pod_id': self.pod_id, - 'model_name': metadata['model_name'], - 'workflow_name': metadata['workflow_name'], - 'opc_server_id': server_id, - 'tag': tag, - }, - ) + await self.emit_metric( + metric_object=metrics.PREDICTION_OPC_WRITING_COUNT, + tags={ + 'pod_id': self.pod_id, + 'model_name': metadata['model_name'], + 'workflow_name': metadata['workflow_name'], + 'opc_server_id': server_id, + 'tag': tag, + }, + ) self.info(f'Metrics written for model {metadata["model_name"]}', metadata) diff --git a/laborious/activities/mlflow.py b/laborious/activities/mlflow.py index cdeeaea..98be4f2 100644 --- a/laborious/activities/mlflow.py +++ b/laborious/activities/mlflow.py @@ -312,9 +312,12 @@ class MLFlow(SientiaMonitoring): self.debug(f'Timestamp: {timestamp}', metadata) # Sort by created_at in descending order and keep first occurrence of each variable/timestamp pair - data = data.sort_values('created_at', ascending=False).drop_duplicates( - subset=['variable', 'timestamp'], keep='first' - ) + if 'created_at' in data.columns: + data = data.sort_values('created_at', ascending=False).drop_duplicates( + subset=['variable', 'timestamp'], keep='first' + ) + else: + data = data.drop_duplicates(subset=['variable', 'timestamp'], keep='first') data.drop(columns=['model_id'], inplace=True, errors='ignore') data.drop(columns=['created_at'], inplace=True, errors='ignore') diff --git a/tests/laborious/activities/test_mlflow.py b/tests/laborious/activities/test_mlflow.py index d3b3976..9343da9 100644 --- a/tests/laborious/activities/test_mlflow.py +++ b/tests/laborious/activities/test_mlflow.py @@ -254,11 +254,11 @@ async def test_retrain_model_success_data_success_retrain(mock_to_datetime, mlfl timestamp = raw_data.__getitem__.return_value.max.return_value - raw_data.sort_values.assert_called_once_with('created_at', ascending=False) - raw_data.sort_values.return_value.drop_duplicates.assert_called_once_with( + raw_data.sort_values.assert_not_called() + raw_data.drop_duplicates.assert_called_once_with( subset=['variable', 'timestamp'], keep='first' ) - raw_data = raw_data.sort_values.return_value.drop_duplicates.return_value + raw_data = raw_data.drop_duplicates.return_value raw_data.drop.assert_has_calls( [ @@ -313,7 +313,9 @@ async def test_retrain_model_success_data_fail_retrain(mock_to_datetime, mlflow) 'message': 'Model retrained failed.', } - mlflow.minio_repository.get_parquet_as_dataframe.return_value = MagicMock() + mlflow.minio_repository.get_parquet_as_dataframe.return_value = MagicMock( + columns=['variable', 'timestamp', 'value', 'created_at'] + ) response = await mlflow.retrain_model( { diff --git a/tests/laborious/utils/repository/test_minio_repository.py b/tests/laborious/utils/repository/test_minio_repository.py index 5e2e3da..ccbf56f 100644 --- a/tests/laborious/utils/repository/test_minio_repository.py +++ b/tests/laborious/utils/repository/test_minio_repository.py @@ -88,6 +88,21 @@ async def test_create_bucket_success(minio_repository): ) +@mark.asyncio +async def test_create_bucket_error(minio_repository): + minio_repository.s3_client.create_bucket.side_effect = ValueError('test') + + with raises(ValueError): + await minio_repository.create_bucket({}) + + + minio_repository.s3_client.create_bucket.assert_called_once_with(Bucket='test') + minio_repository.emit_metric.assert_called_once_with( + metric_object=metrics.MINIO_WRITE_ERROR_COUNT, tags=ANY + ) + minio_repository.observe_lag.assert_not_called() + + @mark.asyncio async def test_ensure_bucket_exists_bucket_exists(minio_repository): assert await minio_repository.ensure_bucket_exists({}) is None diff --git a/tests/laborious/utils/repository/test_model_repository.py b/tests/laborious/utils/repository/test_model_repository.py index e40af0c..9380853 100644 --- a/tests/laborious/utils/repository/test_model_repository.py +++ b/tests/laborious/utils/repository/test_model_repository.py @@ -286,6 +286,16 @@ def test_get_experiment_error(mlflow, mlflow_repository): raise AssertionError('Expected ValueError') +def test_get_experiment_create_error(mlflow, mlflow_repository): + mlflow.get_experiment_by_name.return_value = None + mlflow.create_experiment.return_value = None + mlflow.get_experiment.return_value = None + with pytest.raises(ValueError) as e: + mlflow_repository.get_experiment('test', create_if_not_exists=True) + + assert str(e) == 'Experiment test not found after creation, unknown reason' + + @pytest.mark.asyncio async def test_load_predict_model_sklearn(mlflow, mlflow_repository): result = await mlflow_repository.load_predict_model('test_model', {}, 'sklearn')