diff --git a/scouter/worker/prepare_worker.py b/scouter/worker/prepare_worker.py new file mode 100644 index 0000000..6f1391f --- /dev/null +++ b/scouter/worker/prepare_worker.py @@ -0,0 +1,60 @@ +from typing import Sequence, Type, Any + +from temporalio.worker import Worker, PollerBehaviorAutoscaling +from temporalio.client import Client + +import os + + +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, +): + + main_workflow_name = main_workflow.__name__.upper() + + 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='scouter-queue', + 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'], + ), + ) + ] \ No newline at end of file diff --git a/scouter/worker/worker.py b/scouter/worker/worker.py index 228562c..d5f6be8 100644 --- a/scouter/worker/worker.py +++ b/scouter/worker/worker.py @@ -2,6 +2,8 @@ from temporalio import client, workflow from temporalio.runtime import PrometheusConfig, Runtime, TelemetryConfig from temporalio.worker import PollerBehaviorAutoscaling, Worker +from scouter.worker.prepare_worker import prepare_worker + with workflow.unsafe.imports_passed_through(): import asyncio import os @@ -121,10 +123,10 @@ async def main(): logger.custom_info('Starting Workers...', metadata) workers = [ - Worker( - temporal_client, - task_queue='scouter-queue', - workflows=[Scouter, CoreScouter], + prepare_worker( + temporal_client=temporal_client, + main_workflow=Scouter, + other_workflows=[CoreScouter], activities=[ activities.load_latest_data, activities.get_last_data_timestamp, @@ -136,20 +138,6 @@ async def main(): activities.write_metrics, activities.store_data_package, ], - max_concurrent_workflow_tasks=MAX_CONCURRENT_WORKFLOW_TASKS, - max_concurrent_activities=MAX_CONCURRENT_ACTIVITIES, - max_concurrent_local_activities=MAX_CONCURRENT_LOCAL_ACTIVITIES, - max_cached_workflows=MAX_CACHED_WORKFLOWS, - workflow_task_poller_behavior=PollerBehaviorAutoscaling( - minimum=WORKFLOW_POLLER_BEHAVIUR_MINIMUM, - initial=WORKFLOW_POLLER_BEHAVIUR_INITIAL, - maximum=WORKFLOW_POLLER_BEHAVIUR_MAXIMUM, - ), - activity_task_poller_behavior=PollerBehaviorAutoscaling( - minimum=ACTIVITY_POLLER_BEHAVIUR_MINIMUM, - initial=ACTIVITY_POLLER_BEHAVIUR_INITIAL, - maximum=ACTIVITY_POLLER_BEHAVIUR_MAXIMUM, - ), ) ] diff --git a/values.yaml b/values.yaml index ce700c0..56d9b71 100644 --- a/values.yaml +++ b/values.yaml @@ -69,7 +69,7 @@ resources: memory: 2048Mi requests: cpu: 300m - memory: 512Mi + memory: 256Mi # This is to setup the liveness and readiness probes more information can be found here: https://kubernetes.io/docs/tasks/configure-pod-container/configure-liveness-readiness-startup-probes/ livenessProbe: @@ -78,12 +78,7 @@ livenessProbe: - sh - -c - | - APP_UP=$(curl -s http://localhost:9090/metrics | grep '^app_up{' | grep -o '[0-9]' | head -1) - if [ "$APP_UP" = "1" ]; then - exit 0 - else - exit 1 - fi + curl -sf http://localhost:9090/metrics | grep -q '^app_up{.*} 1' initialDelaySeconds: 30 periodSeconds: 15 timeoutSeconds: 5 @@ -95,12 +90,7 @@ readinessProbe: - sh - -c - | - APP_UP=$(curl -s http://localhost:9090/metrics | grep '^app_up{' | grep -o '[0-9]' | head -1) - if [ "$APP_UP" = "1" ]; then - exit 0 - else - exit 1 - fi + curl -sf http://localhost:9090/metrics | grep -q '^app_up{.*} 1' initialDelaySeconds: 20 periodSeconds: 10 timeoutSeconds: 3 @@ -233,27 +223,27 @@ env: value: "http://library-distribution-server.library.svc.cluster.local:5000" # Temporal worker tuning - - name: MAX_CONCURRENT_WORKFLOW_TASKS + - name: SCOUTER_MAX_CONCURRENT_WORKFLOW_TASKS value: "200" - - name: MAX_CONCURRENT_ACTIVITIES + - name: SCOUTER_MAX_CONCURRENT_ACTIVITIES value: "200" - - name: MAX_CONCURRENT_LOCAL_ACTIVITIES + - name: SCOUTER_MAX_CONCURRENT_LOCAL_ACTIVITIES value: "200" - - name: MAX_CACHED_WORKFLOWS + - name: SCOUTER_MAX_CACHED_WORKFLOWS value: "200" - - name: WORKFLOW_POLLER_BEHAVIUR_MINIMUM + - name: SCOUTER_WORKFLOW_POLLER_BEHAVIUR_MINIMUM value: "10" - - name: WORKFLOW_POLLER_BEHAVIUR_INITIAL + - name: SCOUTER_WORKFLOW_POLLER_BEHAVIUR_INITIAL value: "100" - - name: WORKFLOW_POLLER_BEHAVIUR_MAXIMUM + - name: SCOUTER_WORKFLOW_POLLER_BEHAVIUR_MAXIMUM value: "200" - - name: ACTIVITY_POLLER_BEHAVIUR_MINIMUM + - name: SCOUTER_ACTIVITY_POLLER_BEHAVIUR_MINIMUM value: "10" - - name: ACTIVITY_POLLER_BEHAVIUR_INITIAL + - name: SCOUTER_ACTIVITY_POLLER_BEHAVIUR_INITIAL value: "100" - - name: ACTIVITY_POLLER_BEHAVIUR_MAXIMUM + - name: SCOUTER_ACTIVITY_POLLER_BEHAVIUR_MAXIMUM value: "200"