Files
sientia-dataops-scouter_tem…/scouter/worker/prepare_worker.py
vitor-aignosi 034e0b2520 SIENTIAPDE-1445
Refactor workflow parameters in tests.ipynb and prepare_worker.py to enhance clarity and flexibility. Updated workflow_id and task_queue to accept dynamic values, and standardized timeout parameters for improved configuration.
2025-12-17 16:46:15 -03:00

61 lines
2.3 KiB
Python

import os
from collections.abc import Sequence
from typing import Any
from temporalio.client import Client
from temporalio.worker import PollerBehaviorAutoscaling, Worker
parameters = [
('MAX_CONCURRENT_WORKFLOW_TASKS', '200'),
('MAX_CONCURRENT_ACTIVITIES', '200'),
('MAX_CONCURRENT_LOCAL_ACTIVITIES', '200'),
('MAX_CACHED_WORKFLOWS', '200'),
('WORKFLOW_POLLER_BEHAVIUR_MINIMUM', '10'),
('WORKFLOW_POLLER_BEHAVIUR_INITIAL', '100'),
('WORKFLOW_POLLER_BEHAVIUR_MAXIMUM', '200'),
('ACTIVITY_POLLER_BEHAVIUR_MINIMUM', '10'),
('ACTIVITY_POLLER_BEHAVIUR_INITIAL', '100'),
('ACTIVITY_POLLER_BEHAVIUR_MAXIMUM', '200'),
]
def prepare_worker(
main_workflow: type,
other_workflows: Sequence[type],
activities: Sequence[Any],
temporal_client: Client,
) -> Worker:
main_workflow_name = main_workflow.__name__.upper()
queue_name = f'{main_workflow_name.lower().replace("_", "-")}-queue'
local_workflow_parameters = {}
for parameter in parameters:
local_workflow_parameters[parameter[0]] = int(
os.getenv(main_workflow_name + '_' + parameter[0], parameter[1])
)
return Worker(
temporal_client,
task_queue=queue_name,
workflows=[main_workflow, *other_workflows],
activities=[*activities],
max_concurrent_workflow_tasks=local_workflow_parameters['MAX_CONCURRENT_WORKFLOW_TASKS'],
max_concurrent_activities=local_workflow_parameters['MAX_CONCURRENT_ACTIVITIES'],
max_concurrent_local_activities=local_workflow_parameters[
'MAX_CONCURRENT_LOCAL_ACTIVITIES'
],
max_cached_workflows=local_workflow_parameters['MAX_CACHED_WORKFLOWS'],
workflow_task_poller_behavior=PollerBehaviorAutoscaling(
minimum=local_workflow_parameters['WORKFLOW_POLLER_BEHAVIUR_MINIMUM'],
initial=local_workflow_parameters['WORKFLOW_POLLER_BEHAVIUR_INITIAL'],
maximum=local_workflow_parameters['WORKFLOW_POLLER_BEHAVIUR_MAXIMUM'],
),
activity_task_poller_behavior=PollerBehaviorAutoscaling(
minimum=local_workflow_parameters['ACTIVITY_POLLER_BEHAVIUR_MINIMUM'],
initial=local_workflow_parameters['ACTIVITY_POLLER_BEHAVIUR_INITIAL'],
maximum=local_workflow_parameters['ACTIVITY_POLLER_BEHAVIUR_MAXIMUM'],
),
)