diff --git a/scouter/worker/worker.py b/scouter/worker/worker.py index 75e6715..228562c 100644 --- a/scouter/worker/worker.py +++ b/scouter/worker/worker.py @@ -26,6 +26,26 @@ POD_ID = os.getenv('HOSTNAME', 'localhost') SDK_METRICS_PORT = int(os.getenv('HTTP_SDK_METRICS_PORT', '9091')) +# For optmized latency, Temporal docs recommends fixed slots, ensuring +# high concurency levels. + +MAX_CONCURRENT_WORKFLOW_TASKS = int(os.getenv('MAX_CONCURRENT_WORKFLOW_TASKS', '200')) +MAX_CONCURRENT_ACTIVITIES = int(os.getenv('MAX_CONCURRENT_ACTIVITIES', '200')) +MAX_CONCURRENT_LOCAL_ACTIVITIES = int(os.getenv('MAX_CONCURRENT_LOCAL_ACTIVITIES', '200')) +MAX_CACHED_WORKFLOWS = int(os.getenv('MAX_CACHED_WORKFLOWS', '200')) + + +# Temporal docs also recommends an autoscaling policy, with agrresive limits to prioritize latency over throughput. + +WORKFLOW_POLLER_BEHAVIUR_MINIMUM = int(os.getenv('WORKFLOW_POLLER_BEHAVIUR_MINIMUM', '10')) +WORKFLOW_POLLER_BEHAVIUR_INITIAL = int(os.getenv('WORKFLOW_POLLER_BEHAVIUR_INITIAL', '100')) +WORKFLOW_POLLER_BEHAVIUR_MAXIMUM = int(os.getenv('WORKFLOW_POLLER_BEHAVIUR_MAXIMUM', '200')) + +ACTIVITY_POLLER_BEHAVIUR_MINIMUM = int(os.getenv('ACTIVITY_POLLER_BEHAVIUR_MINIMUM', '10')) +ACTIVITY_POLLER_BEHAVIUR_INITIAL = int(os.getenv('ACTIVITY_POLLER_BEHAVIUR_INITIAL', '100')) +ACTIVITY_POLLER_BEHAVIUR_MAXIMUM = int(os.getenv('ACTIVITY_POLLER_BEHAVIUR_MAXIMUM', '200')) + + async def main(): """ Main entry point for the Scouter Temporal worker. @@ -116,12 +136,20 @@ async def main(): activities.write_metrics, activities.store_data_package, ], - max_concurrent_workflow_tasks=50, - max_concurrent_activities=50, - max_concurrent_local_activities=50, - max_cached_workflows=200, - workflow_task_poller_behavior=PollerBehaviorAutoscaling(), - activity_task_poller_behavior=PollerBehaviorAutoscaling(), + 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 9941b96..b88a2a2 100644 --- a/values.yaml +++ b/values.yaml @@ -209,6 +209,31 @@ env: - name: PYPI_SERVER value: "http://library-distribution-server.library.svc.cluster.local:5000" + # Temporal worker tuning + - name: MAX_CONCURRENT_WORKFLOW_TASKS + value: 200 + - name: MAX_CONCURRENT_ACTIVITIES + value: 200 + - name: MAX_CONCURRENT_LOCAL_ACTIVITIES + value: 200 + - name: MAX_CACHED_WORKFLOWS + value: 200 + + - name: WORKFLOW_POLLER_BEHAVIUR_MINIMUM + value: 10 + - name: WORKFLOW_POLLER_BEHAVIUR_INITIAL + value: 100 + - name: WORKFLOW_POLLER_BEHAVIUR_MAXIMUM + value: 200 + + - name: ACTIVITY_POLLER_BEHAVIUR_MINIMUM + value: 10 + - name: ACTIVITY_POLLER_BEHAVIUR_INITIAL + value: 100 + - name: ACTIVITY_POLLER_BEHAVIUR_MAXIMUM + value: 200 + + ssh: enabled: true secretName: git-ssh-key-sientia-scouter-worker