SIENTIAPDE-1199

Refactor logging in Gates and MLFlow activities to use info level for key operations

- Updated logging statements in the Gates class to replace debug logs with info logs for input and output gate operations, enhancing visibility.
- Modified MLFlow class to use info logs for data transformation and prediction processes, improving clarity in the logging output.
- Adjusted OPC class to return the count of successfully written tags, providing better insight into data writing operations.
This commit is contained in:
vitor-aignosi
2025-08-20 10:38:08 -03:00
parent 89384d7783
commit 8b9bb8d720
4 changed files with 67 additions and 30 deletions

View File

@@ -68,7 +68,7 @@ class Gates(BaseActivity):
""" """
metadata = input_data['metadata'] metadata = input_data['metadata']
self.debug("Performing input gate...", metadata) self.info("Performing input gate...", metadata)
self.debug(f"Input data: {input_data}", metadata) self.debug(f"Input data: {input_data}", metadata)
@@ -103,11 +103,11 @@ class Gates(BaseActivity):
for path_flag in path_priority: for path_flag in path_priority:
if path_flag in filter_output: 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], \ return path_flag, input_filter_functions['path_confidence'][path_flag], \
"Input data with bad quality" "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, "" return None, 0, ""
@activity.defn(name="mlflow_response_gate") @activity.defn(name="mlflow_response_gate")
@@ -128,7 +128,7 @@ class Gates(BaseActivity):
""" """
metadata = input_data['metadata'] metadata = input_data['metadata']
self.debug("Performing mlflow response gate...", metadata) self.info("Performing mlflow response gate...", metadata)
filters = input_data['filters'] filters = input_data['filters']
data = input_data['data'] data = input_data['data']
@@ -169,12 +169,12 @@ class Gates(BaseActivity):
for path_flag in path_priority: for path_flag in path_priority:
if path_flag in filter_output: if path_flag in filter_output:
self.debug( self.info(
f"Mlflow response gate result: {path_flag}", metadata) f"Mlflow response gate result: {path_flag}", metadata)
return path_flag, mlflow_response_filter_functions['path_confidence'][path_flag], \ return path_flag, mlflow_response_filter_functions['path_confidence'][path_flag], \
", ".join(comments) ", ".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, "" return None, 0, ""
@activity.defn(name="mlflow_content_gate") @activity.defn(name="mlflow_content_gate")
@@ -195,7 +195,7 @@ class Gates(BaseActivity):
""" """
metadata = input_data['metadata'] metadata = input_data['metadata']
self.debug("Performing mlflow content gate...", metadata) self.info("Performing mlflow content gate...", metadata)
filters = input_data['filters'] filters = input_data['filters']
data = DataFrame(input_data['data']) data = DataFrame(input_data['data'])
@@ -234,12 +234,12 @@ class Gates(BaseActivity):
for path_flag in path_priority: for path_flag in path_priority:
if path_flag in filter_output: if path_flag in filter_output:
self.debug( self.info(
f"Mlflow content gate result: {path_flag}", metadata) f"Mlflow content gate result: {path_flag}", metadata)
return path_flag, mlflow_content_filter_functions['path_confidence'][path_flag], \ return path_flag, mlflow_content_filter_functions['path_confidence'][path_flag], \
"Transformed data not passed the content filter" "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, "" return None, 0, ""
@activity.defn(name="format_prediction") @activity.defn(name="format_prediction")
@@ -256,7 +256,7 @@ class Gates(BaseActivity):
dict: The formatted data. dict: The formatted data.
""" """
metadata = input_data['metadata'] metadata = input_data['metadata']
self.debug("Formatting prediction...", metadata) self.info("Formatting prediction...", metadata)
data = DataFrame(input_data['data']) data = DataFrame(input_data['data'])
data['timestamp'] = input_data['timestamp'] data['timestamp'] = input_data['timestamp']
@@ -266,6 +266,8 @@ class Gates(BaseActivity):
data['comments'] = "" data['comments'] = ""
data = data.sort_values(by='timestamp') data = data.sort_values(by='timestamp')
self.info(f"Prediction formatted: {data.size} rows", metadata)
return data.to_dict() return data.to_dict()
@activity.defn(name="format_default_prediction") @activity.defn(name="format_default_prediction")
@@ -287,7 +289,7 @@ class Gates(BaseActivity):
metadata = input_data['metadata'] metadata = input_data['metadata']
self.debug("Formatting default prediction...", metadata) self.debug("Formatting default prediction...", metadata)
return DataFrame({ data = DataFrame({
'prediction': [0], 'prediction': [0],
'response_time': [0], 'response_time': [0],
'timestamp': [input_data['timestamp']], 'timestamp': [input_data['timestamp']],
@@ -295,7 +297,10 @@ class Gates(BaseActivity):
'prediction_confidence': [input_data['prediction_confidence']], 'prediction_confidence': [input_data['prediction_confidence']],
'prediction_status': ['Bad'], 'prediction_status': ['Bad'],
'comments': [input_data['comment']] '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") @activity.defn(name="get_last_timestamp")
async def get_last_timestamp(self, input_data: dict[str, Any]) -> str: async def get_last_timestamp(self, input_data: dict[str, Any]) -> str:
@@ -307,9 +312,18 @@ class Gates(BaseActivity):
Returns: Returns:
str: The last timestamp of the data. str: The last timestamp of the data.
""" """
metadata = input_data['metadata']
self.info("Getting last timestamp...", metadata)
data = DataFrame(input_data['data']) data = DataFrame(input_data['data'])
if data.empty: if data.empty:
return datetime.now().strftime('%Y-%m-%d %H:%M:%S') 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()) return max(data['timestamp'].values.tolist())
@activity.defn(name="write_metrics") @activity.defn(name="write_metrics")
@@ -325,6 +339,9 @@ class Gates(BaseActivity):
prediction_confidence = prediction['prediction_confidence'].values[0] prediction_confidence = prediction['prediction_confidence'].values[0]
response_time = prediction['response_time'].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( metrics.PREDICTIONS_WRITTEN_COUNT.labels(
pod_id=self.pod_id, pod_id=self.pod_id,
model_name=metadata['model_name'], model_name=metadata['model_name'],
@@ -342,3 +359,6 @@ class Gates(BaseActivity):
model_name=metadata['model_name'], model_name=metadata['model_name'],
pipeline_name=metadata['workflow_name'] pipeline_name=metadata['workflow_name']
).observe(response_time) ).observe(response_time)
self.info(
f"Metrics written for model {metadata['model_name']}", metadata)

View File

@@ -41,7 +41,7 @@ class MLFlow(BaseActivity):
dict[str, Any]: The transformed data. dict[str, Any]: The transformed data.
""" """
metadata = input_data['metadata'] metadata = input_data['metadata']
self.debug('Transforming data...', metadata) self.info('Transforming data...', metadata)
data = DataFrame(input_data['data']) data = DataFrame(input_data['data'])
model_name = input_data['model_name'] model_name = input_data['model_name']
model_retention = input_data['model_retention'] model_retention = input_data['model_retention']
@@ -70,6 +70,8 @@ class MLFlow(BaseActivity):
self.debug("Transform response data:", metadata) self.debug("Transform response data:", metadata)
self.debug(json.dumps(response_data, indent=4), metadata) self.debug(json.dumps(response_data, indent=4), metadata)
self.info("Data transformed successfully", metadata)
return response_data return response_data
@activity.defn(name="request_predict") @activity.defn(name="request_predict")
@@ -85,7 +87,7 @@ class MLFlow(BaseActivity):
dict[str, Any]: The predicted data. dict[str, Any]: The predicted data.
""" """
metadata = input_data['metadata'] metadata = input_data['metadata']
self.debug('Predicting data...', metadata) self.info('Predicting data...', metadata)
data = DataFrame(input_data['data']) data = DataFrame(input_data['data'])
model_name = input_data['model_name'] model_name = input_data['model_name']
model_retention = input_data['model_retention'] model_retention = input_data['model_retention']
@@ -100,6 +102,8 @@ class MLFlow(BaseActivity):
self.debug("Prediction response data:", metadata) self.debug("Prediction response data:", metadata)
self.debug(json.dumps(response_data, indent=4), metadata) self.debug(json.dumps(response_data, indent=4), metadata)
self.info("Prediction completed successfully", metadata)
return response_data return response_data
@activity.defn(name="retrain_model") @activity.defn(name="retrain_model")

View File

@@ -115,7 +115,9 @@ class OPC(BaseActivity):
def manage_output_tags( def manage_output_tags(
self, server_id: str, config: dict[str, Any], data: DataFrame, 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: if 'prediction_tags' in config:
for tag, tag_config in config['prediction_tags'].items(): for tag, tag_config in config['prediction_tags'].items():
local_success = self.write_data( local_success = self.write_data(
@@ -129,6 +131,7 @@ class OPC(BaseActivity):
if local_success: if local_success:
self.info( self.info(
f"Prediction data written to OPC server {server_id} for tag {tag}.", metadata) f"Prediction data written to OPC server {server_id} for tag {tag}.", metadata)
count += 1
success = success and local_success success = success and local_success
if 'confidence_tags' in config: if 'confidence_tags' in config:
@@ -144,9 +147,10 @@ class OPC(BaseActivity):
if local_success: if local_success:
self.info( self.info(
f"Confidence data written to OPC server {server_id} for tag {tag}.", metadata) f"Confidence data written to OPC server {server_id} for tag {tag}.", metadata)
count += 1
success = success and local_success success = success and local_success
return success return success, count
@activity.defn(name='write_opc_data') @activity.defn(name='write_opc_data')
async def write_opc_data(self, input_data: dict[str, Any]) -> dict[Any, Any]: 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'] 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']) data = DataFrame(input_data['data'])
opc_output_config = input_data['opc_output_config'] 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 = True
success_count = 0
for server_id, config in opc_output_config.items(): for server_id, config in opc_output_config.items():
if not self.validate_server(server_id, metadata): if not self.validate_server(server_id, metadata):
success = False success = False
continue continue
success = success and self.manage_output_tags( local_success, local_count = self.manage_output_tags(
server_id, config, data, metadata, success) 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) return self.process_confidence(data, success, metadata)

View File

@@ -1,5 +1,5 @@
from temporalio import workflow, client from temporalio import workflow, client
from temporalio.worker import Worker from temporalio.worker import Worker, PollerBehaviorAutoscaling
from temporalio.runtime import Runtime, TelemetryConfig, PrometheusConfig from temporalio.runtime import Runtime, TelemetryConfig, PrometheusConfig
with workflow.unsafe.imports_passed_through(): with workflow.unsafe.imports_passed_through():
@@ -86,11 +86,12 @@ async def main():
activities.update_production_model, activities.update_production_model,
activities.export_data_to_postgres activities.export_data_to_postgres
], ],
max_concurrent_workflow_tasks=100, max_concurrent_workflow_tasks=50,
max_concurrent_activities=100, max_concurrent_activities=50,
max_concurrent_local_activities=100, max_concurrent_local_activities=50,
max_concurrent_workflow_task_polls=100, max_cached_workflows=200,
max_cached_workflows=50, workflow_task_poller_behavior=PollerBehaviorAutoscaling(),
activity_task_poller_behavior=PollerBehaviorAutoscaling()
), ),
Worker( Worker(
temporal_client, temporal_client,
@@ -116,11 +117,12 @@ async def main():
activities.export_data_to_postgres, activities.export_data_to_postgres,
activities.write_metrics activities.write_metrics
], ],
max_concurrent_workflow_tasks=100, max_concurrent_workflow_tasks=50,
max_concurrent_activities=100, max_concurrent_activities=50,
max_concurrent_local_activities=100, max_concurrent_local_activities=50,
max_concurrent_workflow_task_polls=100, max_cached_workflows=200,
max_cached_workflows=50, workflow_task_poller_behavior=PollerBehaviorAutoscaling(),
activity_task_poller_behavior=PollerBehaviorAutoscaling()
) )
] ]