SIENTIAPDE-1645: Add comprehensive model training observability metrics and Grafana dashboard.

This commit is contained in:
vitor-aignosi
2026-06-17 10:34:56 -03:00
parent 445fe643fe
commit 76f926a8ab
12 changed files with 866 additions and 175 deletions

View File

@@ -22,7 +22,6 @@ with workflow.unsafe.imports_passed_through():
from sientia_do.observability.metrics_controller import MetricsController
from sientia_do.observability.sientia_monitoring import SientiaMonitoring
from model_manager.metrics import ACTIVITY_EXECUTION_TOTAL, WORKFLOW_EXECUTION_TOTAL
from model_manager.runtime_paths import REPORTS_TEMP_DIR
RETENTION_HOURS = int(os.getenv('CLEANUP_RETENTION_HOURS', '24'))
@@ -84,7 +83,6 @@ class Cleanup(SientiaMonitoring):
"""
metadata = input_data.get('metadata', {})
temp_path = input_data.get('temp_path', REPORTS_TEMP_DIR)
metrics_status = 'success'
cutoff_time = datetime.now() - timedelta(hours=self.retention_hours)
@@ -164,7 +162,6 @@ class Cleanup(SientiaMonitoring):
)
except Exception as e:
metrics_status = 'error'
error_msg = f'Error in directory cleanup: {str(e)}'
trace = traceback.format_exc()
@@ -178,44 +175,3 @@ class Cleanup(SientiaMonitoring):
)
raise
finally:
self._emit_metrics(
metadata=metadata,
metrics_status=metrics_status,
activity_name='cleanup_temp_directories',
emit_workflow_metric=True,
)
def _emit_metrics(
self,
metadata: dict[str, Any],
metrics_status: str,
activity_name: str,
emit_workflow_metric: bool,
) -> None:
"""
Emit workflow and activity execution metrics.
Args:
metadata: Activity metadata containing pod_id and workflow_name
metrics_status: Execution status ('success' or 'error')
activity_name: Name of the activity being executed
"""
if emit_workflow_metric:
self.emit_metric_sync(
metric_object=WORKFLOW_EXECUTION_TOTAL,
tags={
'pod_id': metadata.get('pod_id'),
'workflow_name': metadata.get('workflow_name'),
'status': metrics_status,
},
)
self.emit_metric_sync(
metric_object=ACTIVITY_EXECUTION_TOTAL,
tags={
'pod_id': metadata.get('pod_id'),
'activity_name': activity_name,
'status': metrics_status,
},
)

View File

@@ -10,6 +10,8 @@ from sientia_model.wrappers.sientia_model import SientiaModel
from temporalio import activity, workflow
with workflow.unsafe.imports_passed_through():
import os
import time
import traceback
from typing import Any
@@ -23,6 +25,7 @@ with workflow.unsafe.imports_passed_through():
from sientia_model.model_repository.mlflow_repository import SientiaMLflowRepository
from sientia_model.model_repository.plugin_store import PluginStore
from model_manager import metrics as mm_metrics
from model_manager.utils.models.train_model_params import TrainModelParams
from model_manager.utils.models.train_model_result import TrainModelResult
from model_manager.utils.repository.data_manager_repository import DataManagerRepository
@@ -184,12 +187,12 @@ class Training(SientiaMonitoring):
"""
metadata = input_data.get('metadata')
train_params = TrainModelParams.from_dict(input_data['train_params'])
labels = self._get_training_labels(train_params)
self.info('Starting train_model process', metadata)
try:
# Download training file bytes from MinIO
self.info(
f'Downloading training file from MinIO for {train_params.file_name}', metadata
)
@@ -211,11 +214,16 @@ class Training(SientiaMonitoring):
)
self.info(f'Preparing training data for {train_params.file_name}', metadata)
train_result = self.data_manager_repository.prepare_training_data(
train_file_bytes=train_bytes,
validation_file_bytes=val_bytes,
params=train_params,
metadata=metadata,
train_result = self._prepare_data(train_bytes, val_bytes, train_params, metadata)
mm_metrics.SIENTIA_TRAINING_DATASET_TRAIN_ROWS.labels(**labels).set(
len(train_result.train_data)
)
mm_metrics.SIENTIA_TRAINING_DATASET_VAL_ROWS.labels(**labels).set(
len(train_result.val_data)
)
mm_metrics.SIENTIA_TRAINING_FEATURE_COUNT.labels(**labels).set(
len(train_params.variable_columns)
)
self.info(f'Getting model wrapper for {train_params.model_type}', metadata)
@@ -232,15 +240,121 @@ class Training(SientiaMonitoring):
wrapper.logger = self.logger.base_logger
self.info(f'Training model for {train_params.model_type}', metadata)
train_data = train_result.train_data
val_data = train_result.val_data
train_result = self._fit_model(wrapper, train_result, train_params, metadata)
self.debug(
f'train_model prepared data (head 10):\ntrain:\n{train_data.head(10).to_string()}'
f'\nval:\n{val_data.head(10).to_string()}',
metadata,
self.info(f'Computing regression metrics for {train_params.model_type}', metadata)
train_result = self.data_manager_repository.compute_regression_metrics(
train_result,
wrapper,
metadata=metadata,
)
if train_result.mse_val is not None:
mm_metrics.SIENTIA_TRAINING_MODEL_QUALITY_MSE.labels(**labels).set(
train_result.mse_val
)
if train_result.mae_val is not None:
mm_metrics.SIENTIA_TRAINING_MODEL_QUALITY_MAE.labels(**labels).set(
train_result.mae_val
)
if train_result.r2_val is not None:
mm_metrics.SIENTIA_TRAINING_MODEL_QUALITY_R2.labels(**labels).set(
train_result.r2_val
)
self.info(f'Starting MLflow run for {train_params.model_type}', metadata)
with self.mlflow_repository.start_run(
model_name=train_params.model_name,
run_name=train_result.run_name,
experiment_name=train_result.experiment_name,
tags=None,
metadata=metadata,
) as run_info:
train_result.run_id = run_info.run_id
self._persist_training_artifacts(train_result, train_params, wrapper, metadata)
self.emit_metric_sync(
metric_object=mm_metrics.SIENTIA_TRAINING_MODEL_TRAINED_TOTAL,
tags=labels,
)
return {
'run_name': train_result.run_name,
'experiment_name': train_result.experiment_name,
'run_id': train_result.run_id,
'run_dir': train_result.run_dir,
}
except Exception as e: # noqa: BLE001
error_msg = f'Error training model - error: {str(e)}'
trace = traceback.format_exc()
self.send_notification(
metadata=metadata or {},
notification_id='TRAIN_MODEL_ERROR',
message=error_msg,
block='train_model',
level=NotificationLevel.ERROR,
attachment_content=trace,
)
raise e
def _get_training_labels(self, train_params: TrainModelParams) -> dict:
return {
'pod_id': os.getenv('POD_ID'),
'model_name': train_params.model_name,
'model_type': train_params.model_type,
}
def _prepare_data(
self,
train_bytes: bytes,
val_bytes: bytes | None,
train_params: TrainModelParams,
metadata: dict | None,
) -> TrainModelResult:
labels = self._get_training_labels(train_params)
start_time = time.monotonic()
try:
return self.data_manager_repository.prepare_training_data(
train_file_bytes=train_bytes,
validation_file_bytes=val_bytes,
params=train_params,
metadata=metadata,
)
except Exception:
self.emit_metric_sync(
metric_object=mm_metrics.SIENTIA_TRAINING_DATA_PREPARATION_ERROR_COUNT_TOTAL,
tags=labels,
)
raise
finally:
self.observe_lag_sync(
start_time,
mm_metrics.SIENTIA_TRAINING_DATA_PREPARATION_LAG,
labels,
)
def _fit_model(
self,
wrapper: Any,
train_result: TrainModelResult,
train_params: TrainModelParams,
metadata: dict | None,
) -> TrainModelResult:
labels = self._get_training_labels(train_params)
train_data = train_result.train_data
val_data = train_result.val_data
self.debug(
f'train_model prepared data (head 10):\ntrain:\n{train_data.head(10).to_string()}'
f'\nval:\n{val_data.head(10).to_string()}',
metadata,
)
start_time = time.monotonic()
try:
wrapper.train(
train_data=train_data,
val_data=val_data,
@@ -251,7 +365,6 @@ class Training(SientiaMonitoring):
f'Generating predictions using the trained wrapper for {train_params.model_type}',
metadata,
)
# Generate predictions using the trained wrapper
transformed_train, _ = wrapper.transform(train_data)
transformed_val, _ = wrapper.transform(val_data)
@@ -275,47 +388,20 @@ class Training(SientiaMonitoring):
train_result.y_train_pred = y_train_pred_df
train_result.y_pred = y_val_pred_df
self.info(f'Computing regression metrics for {train_params.model_type}', metadata)
train_result = self.data_manager_repository.compute_regression_metrics(
train_result,
wrapper,
metadata=metadata,
return train_result
except Exception:
self.emit_metric_sync(
metric_object=mm_metrics.SIENTIA_TRAINING_MODEL_FIT_ERROR_COUNT_TOTAL,
tags=labels,
)
self.info(f'Starting MLflow run for {train_params.model_type}', metadata)
with self.mlflow_repository.start_run(
model_name=train_params.model_name,
run_name=train_result.run_name,
experiment_name=train_result.experiment_name,
tags=None,
metadata=metadata,
) as run_info:
train_result.run_id = run_info.run_id
self._persist_training_artifacts(train_result, train_params, wrapper, metadata)
return {
'run_name': train_result.run_name,
'experiment_name': train_result.experiment_name,
'run_id': train_result.run_id,
'run_dir': train_result.run_dir,
}
except Exception as e: # noqa: BLE001
error_msg = f'Error training model - error: {str(e)}'
trace = traceback.format_exc()
self.send_notification(
metadata=metadata or {},
notification_id='TRAIN_MODEL_ERROR',
message=error_msg,
block='train_model',
level=NotificationLevel.ERROR,
attachment_content=trace,
raise
finally:
self.observe_lag_sync(
start_time,
mm_metrics.SIENTIA_TRAINING_MODEL_FIT_LAG,
labels,
)
raise e
def _persist_training_artifacts(
self,
train_result: TrainModelResult,