diff --git a/laborious/activities/gates.py b/laborious/activities/gates.py index b4a85eb..32d1006 100644 --- a/laborious/activities/gates.py +++ b/laborious/activities/gates.py @@ -68,7 +68,7 @@ class Gates(BaseActivity): """ metadata = input_data['metadata'] - self.debug("Performing input gate...", metadata) + self.info("Performing input gate...", metadata) self.debug(f"Input data: {input_data}", metadata) @@ -103,11 +103,11 @@ class Gates(BaseActivity): for path_flag in path_priority: if path_flag in filter_output: - self.debug(f"Input gate result: {path_flag}", metadata) + self.info(f"Input gate result: {path_flag}", metadata) return path_flag, input_filter_functions['path_confidence'][path_flag], \ "Input data with bad quality" - self.debug("Nothing was filtered by the input gate", metadata) + self.info("Nothing was filtered by the input gate", metadata) return None, 0, "" @activity.defn(name="mlflow_response_gate") @@ -128,7 +128,7 @@ class Gates(BaseActivity): """ metadata = input_data['metadata'] - self.debug("Performing mlflow response gate...", metadata) + self.info("Performing mlflow response gate...", metadata) filters = input_data['filters'] data = input_data['data'] @@ -169,12 +169,12 @@ class Gates(BaseActivity): for path_flag in path_priority: if path_flag in filter_output: - self.debug( + self.info( f"Mlflow response gate result: {path_flag}", metadata) return path_flag, mlflow_response_filter_functions['path_confidence'][path_flag], \ ", ".join(comments) - self.debug("Nothing was filtered by the mlflow response gate", metadata) + self.info("Nothing was filtered by the mlflow response gate", metadata) return None, 0, "" @activity.defn(name="mlflow_content_gate") @@ -195,7 +195,7 @@ class Gates(BaseActivity): """ metadata = input_data['metadata'] - self.debug("Performing mlflow content gate...", metadata) + self.info("Performing mlflow content gate...", metadata) filters = input_data['filters'] data = DataFrame(input_data['data']) @@ -234,12 +234,12 @@ class Gates(BaseActivity): for path_flag in path_priority: if path_flag in filter_output: - self.debug( + self.info( f"Mlflow content gate result: {path_flag}", metadata) return path_flag, mlflow_content_filter_functions['path_confidence'][path_flag], \ "Transformed data not passed the content filter" - self.debug("Nothing was filtered by the mlflow content gate", metadata) + self.info("Nothing was filtered by the mlflow content gate", metadata) return None, 0, "" @activity.defn(name="format_prediction") @@ -256,7 +256,7 @@ class Gates(BaseActivity): dict: The formatted data. """ metadata = input_data['metadata'] - self.debug("Formatting prediction...", metadata) + self.info("Formatting prediction...", metadata) data = DataFrame(input_data['data']) data['timestamp'] = input_data['timestamp'] @@ -266,6 +266,8 @@ class Gates(BaseActivity): data['comments'] = "" data = data.sort_values(by='timestamp') + self.info(f"Prediction formatted: {data.size} rows", metadata) + return data.to_dict() @activity.defn(name="format_default_prediction") @@ -287,7 +289,7 @@ class Gates(BaseActivity): metadata = input_data['metadata'] self.debug("Formatting default prediction...", metadata) - return DataFrame({ + data = DataFrame({ 'prediction': [0], 'response_time': [0], 'timestamp': [input_data['timestamp']], @@ -295,7 +297,10 @@ class Gates(BaseActivity): 'prediction_confidence': [input_data['prediction_confidence']], 'prediction_status': ['Bad'], 'comments': [input_data['comment']] - }).to_dict() + }) + + self.info(f"Default prediction formatted: {data.size} rows", metadata) + return data.to_dict() @activity.defn(name="get_last_timestamp") async def get_last_timestamp(self, input_data: dict[str, Any]) -> str: @@ -307,9 +312,18 @@ class Gates(BaseActivity): Returns: str: The last timestamp of the data. """ + metadata = input_data['metadata'] + + self.info("Getting last timestamp...", metadata) + data = DataFrame(input_data['data']) + if data.empty: return datetime.now().strftime('%Y-%m-%d %H:%M:%S') + + self.info( + f"Last timestamp: {max(data['timestamp'].values.tolist())}", metadata) + return max(data['timestamp'].values.tolist()) @activity.defn(name="write_metrics") @@ -325,6 +339,9 @@ class Gates(BaseActivity): prediction_confidence = prediction['prediction_confidence'].values[0] response_time = prediction['response_time'].values[0] + self.info( + f"Writing metrics for model {metadata['model_name']}", metadata) + metrics.PREDICTIONS_WRITTEN_COUNT.labels( pod_id=self.pod_id, model_name=metadata['model_name'], @@ -342,3 +359,6 @@ class Gates(BaseActivity): model_name=metadata['model_name'], pipeline_name=metadata['workflow_name'] ).observe(response_time) + + self.info( + f"Metrics written for model {metadata['model_name']}", metadata) diff --git a/laborious/activities/mlflow.py b/laborious/activities/mlflow.py index d3e918d..3ad9aae 100644 --- a/laborious/activities/mlflow.py +++ b/laborious/activities/mlflow.py @@ -41,7 +41,7 @@ class MLFlow(BaseActivity): dict[str, Any]: The transformed data. """ metadata = input_data['metadata'] - self.debug('Transforming data...', metadata) + self.info('Transforming data...', metadata) data = DataFrame(input_data['data']) model_name = input_data['model_name'] model_retention = input_data['model_retention'] @@ -70,6 +70,8 @@ class MLFlow(BaseActivity): self.debug("Transform response data:", metadata) self.debug(json.dumps(response_data, indent=4), metadata) + self.info("Data transformed successfully", metadata) + return response_data @activity.defn(name="request_predict") @@ -85,7 +87,7 @@ class MLFlow(BaseActivity): dict[str, Any]: The predicted data. """ metadata = input_data['metadata'] - self.debug('Predicting data...', metadata) + self.info('Predicting data...', metadata) data = DataFrame(input_data['data']) model_name = input_data['model_name'] model_retention = input_data['model_retention'] @@ -100,6 +102,8 @@ class MLFlow(BaseActivity): self.debug("Prediction response data:", metadata) self.debug(json.dumps(response_data, indent=4), metadata) + self.info("Prediction completed successfully", metadata) + return response_data @activity.defn(name="retrain_model") diff --git a/laborious/activities/opc.py b/laborious/activities/opc.py index c01af45..debb011 100644 --- a/laborious/activities/opc.py +++ b/laborious/activities/opc.py @@ -115,7 +115,9 @@ class OPC(BaseActivity): def manage_output_tags( self, server_id: str, config: dict[str, Any], data: DataFrame, - metadata: dict[str, Any], success: bool) -> bool: + metadata: dict[str, Any], success: bool) -> tuple[bool, int]: + + count = 0 if 'prediction_tags' in config: for tag, tag_config in config['prediction_tags'].items(): local_success = self.write_data( @@ -129,6 +131,7 @@ class OPC(BaseActivity): if local_success: self.info( f"Prediction data written to OPC server {server_id} for tag {tag}.", metadata) + count += 1 success = success and local_success if 'confidence_tags' in config: @@ -144,9 +147,10 @@ class OPC(BaseActivity): if local_success: self.info( f"Confidence data written to OPC server {server_id} for tag {tag}.", metadata) + count += 1 success = success and local_success - return success + return success, count @activity.defn(name='write_opc_data') async def write_opc_data(self, input_data: dict[str, Any]) -> dict[Any, Any]: @@ -168,21 +172,28 @@ class OPC(BaseActivity): """ metadata = input_data['metadata'] - self.debug("Writing data to OPC servers...", metadata) + self.info("Writing data to OPC servers...", metadata) data = DataFrame(input_data['data']) opc_output_config = input_data['opc_output_config'] - self.debug(data, metadata) + self.info(f"Data to write: {data.size} rows", metadata) success = True + success_count = 0 + for server_id, config in opc_output_config.items(): if not self.validate_server(server_id, metadata): success = False continue - success = success and self.manage_output_tags( + local_success, local_count = self.manage_output_tags( server_id, config, data, metadata, success) + success = success and local_success + success_count += local_count + + self.info( + f"Data written to OPC server {server_id}: {local_count} of {len(config['prediction_tags'])} prediction tags and {len(config['confidence_tags'])} confidence tags", metadata) return self.process_confidence(data, success, metadata) diff --git a/laborious/worker/worker.py b/laborious/worker/worker.py index 0621aa4..85326df 100644 --- a/laborious/worker/worker.py +++ b/laborious/worker/worker.py @@ -1,5 +1,5 @@ from temporalio import workflow, client -from temporalio.worker import Worker +from temporalio.worker import Worker, PollerBehaviorAutoscaling from temporalio.runtime import Runtime, TelemetryConfig, PrometheusConfig with workflow.unsafe.imports_passed_through(): @@ -86,11 +86,12 @@ async def main(): activities.update_production_model, activities.export_data_to_postgres ], - max_concurrent_workflow_tasks=100, - max_concurrent_activities=100, - max_concurrent_local_activities=100, - max_concurrent_workflow_task_polls=100, - max_cached_workflows=50, + max_concurrent_workflow_tasks=50, + max_concurrent_activities=50, + max_concurrent_local_activities=50, + max_cached_workflows=200, + workflow_task_poller_behavior=PollerBehaviorAutoscaling(), + activity_task_poller_behavior=PollerBehaviorAutoscaling() ), Worker( temporal_client, @@ -116,11 +117,12 @@ async def main(): activities.export_data_to_postgres, activities.write_metrics ], - max_concurrent_workflow_tasks=100, - max_concurrent_activities=100, - max_concurrent_local_activities=100, - max_concurrent_workflow_task_polls=100, - max_cached_workflows=50, + max_concurrent_workflow_tasks=50, + max_concurrent_activities=50, + max_concurrent_local_activities=50, + max_cached_workflows=200, + workflow_task_poller_behavior=PollerBehaviorAutoscaling(), + activity_task_poller_behavior=PollerBehaviorAutoscaling() ) ]