SIENTIAPDE-1478
Update environment variables in values.yaml and refactor worker.py for improved worker preparation - Removed KAFKA_BOOTSTRAP_SERVERS from environment variables in values.yaml. - Added PYPI_SERVER environment variable for library distribution. - Refactored worker.py to replace resource tuner and poller behavior with a new prepare_worker function, streamlining worker initialization and enhancing code clarity.
This commit is contained in:
72
laborious/worker/prepare_worker.py
Normal file
72
laborious/worker/prepare_worker.py
Normal file
@@ -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'],
|
||||
),
|
||||
)
|
||||
@@ -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 = []
|
||||
|
||||
Reference in New Issue
Block a user