SIENTIAPDE-1174
Integrate write_metrics activity into CoreScouter for enhanced metric logging and update worker.py to include write_metrics in workflow tasks.
This commit is contained in:
@@ -84,7 +84,8 @@ async def main():
|
|||||||
activities.aggregate_data,
|
activities.aggregate_data,
|
||||||
activities.group_and_hold_data,
|
activities.group_and_hold_data,
|
||||||
activities.export_data_to_postgres,
|
activities.export_data_to_postgres,
|
||||||
activities.store_data_package
|
activities.write_metrics,
|
||||||
|
activities.store_data_package,
|
||||||
],
|
],
|
||||||
max_concurrent_workflow_tasks=100,
|
max_concurrent_workflow_tasks=100,
|
||||||
max_concurrent_activities=100,
|
max_concurrent_activities=100,
|
||||||
|
|||||||
@@ -88,11 +88,14 @@ class CoreScouter:
|
|||||||
start_to_close_timeout=timedelta(seconds=60)
|
start_to_close_timeout=timedelta(seconds=60)
|
||||||
)
|
)
|
||||||
|
|
||||||
metrics.LABORIOUS_DATA_WRITTEN_COUNT.labels(
|
await workflow.execute_activity_method(
|
||||||
pod_id=metadata['pod_id'],
|
Activities.write_metrics,
|
||||||
model_name=input_data['model_name'],
|
{
|
||||||
pipeline_name=input_data['workflow_name']
|
**metadata,
|
||||||
).inc()
|
},
|
||||||
|
retry_policy=retry_policy,
|
||||||
|
start_to_close_timeout=timedelta(seconds=60)
|
||||||
|
)
|
||||||
|
|
||||||
if input_data.get('debug_data_package', False):
|
if input_data.get('debug_data_package', False):
|
||||||
await workflow.execute_activity_method(
|
await workflow.execute_activity_method(
|
||||||
|
|||||||
Reference in New Issue
Block a user