diff --git a/laborious/worker/prepare_worker.py b/laborious/worker/prepare_worker.py new file mode 100644 index 0000000..c8b68f2 --- /dev/null +++ b/laborious/worker/prepare_worker.py @@ -0,0 +1,72 @@ +import os +import re +from collections.abc import Sequence +from typing import Any + +from sientia_do.observability.logger import Logger +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 camel_to_snake(text: str) -> str: + """Convert camelCase or PascalCase to snake_case.""" + text = re.sub('(.)([A-Z][a-z]+)', r'\1_\2', text) + text = re.sub('([a-z0-9])([A-Z])', r'\1_\2', text) + return text.lower() + + +def prepare_worker( + main_workflow: type, + other_workflows: Sequence[type], + activities: Sequence[Any], + temporal_client: Client, + logger: Logger, +) -> Worker: + main_workflow_name = main_workflow.__name__.upper() + + queue_name = f'{camel_to_snake(main_workflow.__name__)}-queue' + + local_workflow_parameters = {} + + for parameter in parameters: + local_workflow_parameters[parameter[0]] = int( + os.getenv(main_workflow_name + '_' + parameter[0], parameter[1]) + ) + + logger.info(f'Preparing worker for {main_workflow_name} with queue {queue_name}') + + 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'], + ), + ) diff --git a/laborious/worker/worker.py b/laborious/worker/worker.py index c6248df..50f1cb5 100644 --- a/laborious/worker/worker.py +++ b/laborious/worker/worker.py @@ -27,21 +27,6 @@ Environment Variables: - HTTP_METRICS_PORT: Prometheus metrics server port (default: 9090) - HTTP_SDK_METRICS_PORT: Temporal SDK metrics port (default: 9091) - PROJECT_NAME: Project name for notifications (default: laborious) - -Tuner Configuration (Resource-based scaling): -- TUNER_TARGET_MEMORY_USAGE: Target memory usage (0.0-1.0, default: 0.75) -- TUNER_TARGET_CPU_USAGE: Target CPU usage (0.0-1.0, default: 0.80) -- TUNER_WORKFLOW_MIN_SLOTS: Minimum workflow slots (default: 5) -- TUNER_WORKFLOW_MAX_SLOTS: Maximum workflow slots (default: 50) -- TUNER_ACTIVITY_MIN_SLOTS: Minimum activity slots (default: 5) -- TUNER_ACTIVITY_MAX_SLOTS: Maximum activity slots (default: 50) -- TUNER_WORKFLOW_RAMP_THROTTLE_MS: Workflow ramp throttle in ms (default: 100) -- TUNER_ACTIVITY_RAMP_THROTTLE_MS: Activity ramp throttle in ms (default: 50) - -Poller Configuration: -- POLLER_MINIMUM: Minimum number of pollers (default: 1) -- POLLER_MAXIMUM: Maximum number of pollers (default: 10) -- POLLER_INITIAL: Initial number of pollers (default: 2) """ from temporalio import client, workflow @@ -83,54 +68,11 @@ with workflow.unsafe.imports_passed_through(): FormatAndExportPrediction, ) from laborious.workflows.sub_workflows.prediction_process import PredictionProcess + from laborious.worker.prepare_worker import prepare_worker POD_ID = os.getenv('POD_ID') SDK_METRICS_PORT = int(os.getenv('HTTP_SDK_METRICS_PORT', '9091')) - -def create_resource_tuner() -> WorkerTuner: - """Create a resource-based tuner from environment variables.""" - target_memory = float(os.getenv('TUNER_TARGET_MEMORY_USAGE', '0.75')) - target_cpu = float(os.getenv('TUNER_TARGET_CPU_USAGE', '0.50')) - workflow_min = int(os.getenv('TUNER_WORKFLOW_MIN_SLOTS', '5')) - workflow_max = int(os.getenv('TUNER_WORKFLOW_MAX_SLOTS', '50')) - activity_min = int(os.getenv('TUNER_ACTIVITY_MIN_SLOTS', '5')) - activity_max = int(os.getenv('TUNER_ACTIVITY_MAX_SLOTS', '50')) - local_activity_min = int(os.getenv('TUNER_LOCAL_ACTIVITY_MIN_SLOTS', '1')) - local_activity_max = int(os.getenv('TUNER_LOCAL_ACTIVITY_MAX_SLOTS', '30')) - workflow_ramp = int(os.getenv('TUNER_WORKFLOW_RAMP_THROTTLE_MS', '100')) - activity_ramp = int(os.getenv('TUNER_ACTIVITY_RAMP_THROTTLE_MS', '50')) - local_activity_ramp = int(os.getenv('TUNER_LOCAL_ACTIVITY_RAMP_THROTTLE_MS', '50')) - - return WorkerTuner.create_resource_based( - target_memory_usage=target_memory, - target_cpu_usage=target_cpu, - workflow_config=ResourceBasedSlotConfig( - minimum_slots=workflow_min, - maximum_slots=workflow_max, - ramp_throttle=timedelta(milliseconds=workflow_ramp), - ), - activity_config=ResourceBasedSlotConfig( - minimum_slots=activity_min, - maximum_slots=activity_max, - ramp_throttle=timedelta(milliseconds=activity_ramp), - ), - local_activity_config=ResourceBasedSlotConfig( - minimum_slots=local_activity_min, - maximum_slots=local_activity_max, - ramp_throttle=timedelta(milliseconds=local_activity_ramp), - ), - ) - - -def create_poller_behavior() -> PollerBehaviorAutoscaling: - """Create poller behavior from environment variables.""" - minimum = int(os.getenv('POLLER_MINIMUM', '1')) - maximum = int(os.getenv('POLLER_MAXIMUM', '10')) - initial = int(os.getenv('POLLER_INITIAL', '2')) - return PollerBehaviorAutoscaling(minimum=minimum, maximum=maximum, initial=initial) - - async def main(): """ Main entry point for the Laborious worker application. @@ -210,14 +152,11 @@ async def main(): logger.custom_info('Starting Workers...', metadata) - tuner = create_resource_tuner() - poller = create_poller_behavior() - workers = [ - Worker( - temporal_client, - task_queue='minimal_retrain-queue', - workflows=[MinimalRetrain], + prepare_worker( + temporal_client=temporal_client, + main_workflow=MinimalRetrain, + other_workflows=[], activities=[ activities.load_custom_query, activities.query_to_minio, @@ -226,44 +165,35 @@ async def main(): activities.format_retrain_report, activities.export_data_to_postgres, ], - tuner=tuner, - max_cached_workflows=2, - workflow_task_poller_behavior=poller, - activity_task_poller_behavior=poller, + logger=logger, ), - Worker( - temporal_client, - task_queue='drift-queue', - workflows=[Drift], + prepare_worker( + temporal_client=temporal_client, + main_workflow=SimpleMetrics, + other_workflows=[], + activities=[ + activities.load_custom_query, + activities.calculate_simple_metrics, + activities.export_data_to_postgres, + ], + logger=logger, + ), + prepare_worker( + temporal_client=temporal_client, + main_workflow=Drift, + other_workflows=[], activities=[ activities.load_custom_query, activities.get_reference_data, activities.calculate_drift, activities.export_data_to_postgres, ], - tuner=tuner, - max_cached_workflows=2, - workflow_task_poller_behavior=poller, - activity_task_poller_behavior=poller, + logger=logger, ), - Worker( - temporal_client, - task_queue='simple_metrics-queue', - workflows=[SimpleMetrics], - activities=[ - activities.load_custom_query, - activities.calculate_simple_metrics, - activities.export_data_to_postgres, - ], - tuner=tuner, - max_cached_workflows=2, - workflow_task_poller_behavior=poller, - activity_task_poller_behavior=poller, - ), - Worker( - temporal_client, - task_queue='predictions_batch-queue', - workflows=[PredictionsBatch, PredictionProcess, FormatAndExportPrediction], + prepare_worker( + temporal_client=temporal_client, + main_workflow=PredictionsBatch, + other_workflows=[PredictionProcess, FormatAndExportPrediction], activities=[ # MLFlow activities.request_predict, @@ -284,11 +214,8 @@ async def main(): activities.export_data_to_postgres, activities.write_metrics, ], - tuner=tuner, - max_cached_workflows=200, - workflow_task_poller_behavior=poller, - activity_task_poller_behavior=poller, - ), + logger=logger, + ) ] handlers = [] diff --git a/values.yaml b/values.yaml index a41347c..2d3a7a4 100644 --- a/values.yaml +++ b/values.yaml @@ -187,8 +187,6 @@ env: - name: OPC_URL value: "opc.tcp://sientia-opc-simulator-opc.sientia.svc.cluster.local:4840" - - name: KAFKA_BOOTSTRAP_SERVERS - value: "kafka.kafka.svc.cluster.local:9092" - name: LOG_LEVEL value: "DEBUG" @@ -226,6 +224,9 @@ env: - name: MINIO_DEFAULT_BUCKET value: "sientia" + - name: PYPI_SERVER + value: "http://library-distribution-server.library.svc.cluster.local:5000" + ssh: enabled: true secretName: git-ssh-key-sientia-laborious-worker