|
6 | 6 | import io.temporal.spring.boot.autoconfigure.properties.WorkerProperties; |
7 | 7 | import io.temporal.worker.WorkerDeploymentOptions; |
8 | 8 | import io.temporal.worker.WorkerOptions; |
| 9 | +import io.temporal.worker.tuning.PollerBehaviorAutoscaling; |
9 | 10 | import java.util.Optional; |
10 | 11 | import javax.annotation.Nonnull; |
11 | 12 | import javax.annotation.Nullable; |
@@ -50,70 +51,92 @@ WorkerOptions createWorkerOptions() { |
50 | 51 | Optional.ofNullable(threadsConfiguration.getMaxConcurrentNexusTaskPollers()) |
51 | 52 | .ifPresent(options::setMaxConcurrentNexusTaskPollers); |
52 | 53 | if (threadsConfiguration.getWorkflowTaskPollersConfiguration() != null) { |
53 | | - Optional.ofNullable( |
| 54 | + WorkerProperties.PollerConfigurationProperties.PollerBehaviorAutoscalingConfiguration |
| 55 | + pollerBehaviorAutoscaling = |
54 | 56 | threadsConfiguration |
55 | 57 | .getWorkflowTaskPollersConfiguration() |
56 | | - .getPollerBehaviorAutoscaling()) |
57 | | - .ifPresent(options::setWorkflowTaskPollersBehavior); |
| 58 | + .getPollerBehaviorAutoscaling(); |
| 59 | + if (pollerBehaviorAutoscaling != null && pollerBehaviorAutoscaling.isEnabled()) { |
| 60 | + options.setWorkflowTaskPollersBehavior( |
| 61 | + new PollerBehaviorAutoscaling( |
| 62 | + pollerBehaviorAutoscaling.getMinConcurrentTaskPollers(), |
| 63 | + pollerBehaviorAutoscaling.getMaxConcurrentTaskPollers(), |
| 64 | + pollerBehaviorAutoscaling.getInitialConcurrentTaskPollers())); |
| 65 | + } |
58 | 66 | } |
59 | 67 | if (threadsConfiguration.getActivityTaskPollersConfiguration() != null) { |
60 | | - Optional.ofNullable( |
| 68 | + WorkerProperties.PollerConfigurationProperties.PollerBehaviorAutoscalingConfiguration |
| 69 | + pollerBehaviorAutoscaling = |
61 | 70 | threadsConfiguration |
62 | 71 | .getActivityTaskPollersConfiguration() |
63 | | - .getPollerBehaviorAutoscaling()) |
64 | | - .ifPresent(options::setActivityTaskPollersBehavior); |
| 72 | + .getPollerBehaviorAutoscaling(); |
| 73 | + if (pollerBehaviorAutoscaling != null && pollerBehaviorAutoscaling.isEnabled()) { |
| 74 | + options.setActivityTaskPollersBehavior( |
| 75 | + new PollerBehaviorAutoscaling( |
| 76 | + pollerBehaviorAutoscaling.getMinConcurrentTaskPollers(), |
| 77 | + pollerBehaviorAutoscaling.getMaxConcurrentTaskPollers(), |
| 78 | + pollerBehaviorAutoscaling.getInitialConcurrentTaskPollers())); |
| 79 | + } |
65 | 80 | } |
66 | 81 | if (threadsConfiguration.getNexusTaskPollersConfiguration() != null) { |
67 | | - Optional.ofNullable( |
| 82 | + WorkerProperties.PollerConfigurationProperties.PollerBehaviorAutoscalingConfiguration |
| 83 | + pollerBehaviorAutoscaling = |
68 | 84 | threadsConfiguration |
69 | 85 | .getNexusTaskPollersConfiguration() |
70 | | - .getPollerBehaviorAutoscaling()) |
71 | | - .ifPresent(options::setNexusTaskPollersBehavior); |
| 86 | + .getPollerBehaviorAutoscaling(); |
| 87 | + if (pollerBehaviorAutoscaling != null && pollerBehaviorAutoscaling.isEnabled()) { |
| 88 | + options.setNexusTaskPollersBehavior( |
| 89 | + new PollerBehaviorAutoscaling( |
| 90 | + pollerBehaviorAutoscaling.getMinConcurrentTaskPollers(), |
| 91 | + pollerBehaviorAutoscaling.getMaxConcurrentTaskPollers(), |
| 92 | + pollerBehaviorAutoscaling.getInitialConcurrentTaskPollers())); |
| 93 | + } |
72 | 94 | } |
73 | | - } |
74 | 95 |
|
75 | | - WorkerProperties.RateLimitsConfigurationProperties rateLimitConfiguration = |
76 | | - workerProperties.getRateLimits(); |
77 | | - if (rateLimitConfiguration != null) { |
78 | | - Optional.ofNullable(rateLimitConfiguration.getMaxWorkerActivitiesPerSecond()) |
79 | | - .ifPresent(options::setMaxWorkerActivitiesPerSecond); |
80 | | - Optional.ofNullable(rateLimitConfiguration.getMaxTaskQueueActivitiesPerSecond()) |
81 | | - .ifPresent(options::setMaxTaskQueueActivitiesPerSecond); |
82 | | - } |
| 96 | + WorkerProperties.RateLimitsConfigurationProperties rateLimitConfiguration = |
| 97 | + workerProperties.getRateLimits(); |
| 98 | + if (rateLimitConfiguration != null) { |
| 99 | + Optional.ofNullable(rateLimitConfiguration.getMaxWorkerActivitiesPerSecond()) |
| 100 | + .ifPresent(options::setMaxWorkerActivitiesPerSecond); |
| 101 | + Optional.ofNullable(rateLimitConfiguration.getMaxTaskQueueActivitiesPerSecond()) |
| 102 | + .ifPresent(options::setMaxTaskQueueActivitiesPerSecond); |
| 103 | + } |
83 | 104 |
|
84 | | - WorkerProperties.BuildIdConfigurationProperties buildIdConfigurations = |
85 | | - workerProperties.getBuildId(); |
86 | | - if (buildIdConfigurations != null) { |
87 | | - Optional.ofNullable(buildIdConfigurations.getWorkerBuildId()) |
88 | | - .ifPresent(options::setBuildId); |
89 | | - options.setUseBuildIdForVersioning(buildIdConfigurations.getEnabledWorkerVersioning()); |
90 | | - } |
| 105 | + WorkerProperties.BuildIdConfigurationProperties buildIdConfigurations = |
| 106 | + workerProperties.getBuildId(); |
| 107 | + if (buildIdConfigurations != null) { |
| 108 | + Optional.ofNullable(buildIdConfigurations.getWorkerBuildId()) |
| 109 | + .ifPresent(options::setBuildId); |
| 110 | + options.setUseBuildIdForVersioning(buildIdConfigurations.getEnabledWorkerVersioning()); |
| 111 | + } |
91 | 112 |
|
92 | | - WorkerProperties.VirtualThreadConfigurationProperties virtualThreadConfiguration = |
93 | | - workerProperties.getVirtualThreads(); |
94 | | - if (virtualThreadConfiguration != null) { |
95 | | - Optional.ofNullable(virtualThreadConfiguration.isUsingVirtualThreads()) |
96 | | - .ifPresent(options::setUsingVirtualThreads); |
97 | | - Optional.ofNullable(virtualThreadConfiguration.isUsingVirtualThreadsOnWorkflowWorker()) |
98 | | - .ifPresent(options::setUsingVirtualThreadsOnWorkflowWorker); |
99 | | - Optional.ofNullable(virtualThreadConfiguration.isUsingVirtualThreadsOnActivityWorker()) |
100 | | - .ifPresent(options::setUsingVirtualThreadsOnActivityWorker); |
101 | | - Optional.ofNullable(virtualThreadConfiguration.isUsingVirtualThreadsOnLocalActivityWorker()) |
102 | | - .ifPresent(options::setUsingVirtualThreadsOnLocalActivityWorker); |
103 | | - Optional.ofNullable(virtualThreadConfiguration.isUsingVirtualThreadsOnNexusWorker()) |
104 | | - .ifPresent(options::setUsingVirtualThreadsOnNexusWorker); |
105 | | - } |
106 | | - WorkerProperties.WorkerDeploymentConfigurationProperties workerDeploymentConfiguration = |
107 | | - workerProperties.getDeploymentProperties(); |
108 | | - if (workerDeploymentConfiguration != null) { |
109 | | - WorkerDeploymentOptions.Builder opts = WorkerDeploymentOptions.newBuilder(); |
110 | | - Optional.ofNullable(workerDeploymentConfiguration.getUseVersioning()) |
111 | | - .ifPresent(opts::setUseVersioning); |
112 | | - Optional.ofNullable(workerDeploymentConfiguration.getDeploymentVersion()) |
113 | | - .ifPresent((v) -> opts.setVersion(WorkerDeploymentVersion.fromCanonicalString(v))); |
114 | | - Optional.ofNullable(workerDeploymentConfiguration.getDefaultVersioningBehavior()) |
115 | | - .ifPresent(opts::setDefaultVersioningBehavior); |
116 | | - options.setDeploymentOptions(opts.build()); |
| 113 | + WorkerProperties.VirtualThreadConfigurationProperties virtualThreadConfiguration = |
| 114 | + workerProperties.getVirtualThreads(); |
| 115 | + if (virtualThreadConfiguration != null) { |
| 116 | + Optional.ofNullable(virtualThreadConfiguration.isUsingVirtualThreads()) |
| 117 | + .ifPresent(options::setUsingVirtualThreads); |
| 118 | + Optional.ofNullable(virtualThreadConfiguration.isUsingVirtualThreadsOnWorkflowWorker()) |
| 119 | + .ifPresent(options::setUsingVirtualThreadsOnWorkflowWorker); |
| 120 | + Optional.ofNullable(virtualThreadConfiguration.isUsingVirtualThreadsOnActivityWorker()) |
| 121 | + .ifPresent(options::setUsingVirtualThreadsOnActivityWorker); |
| 122 | + Optional.ofNullable( |
| 123 | + virtualThreadConfiguration.isUsingVirtualThreadsOnLocalActivityWorker()) |
| 124 | + .ifPresent(options::setUsingVirtualThreadsOnLocalActivityWorker); |
| 125 | + Optional.ofNullable(virtualThreadConfiguration.isUsingVirtualThreadsOnNexusWorker()) |
| 126 | + .ifPresent(options::setUsingVirtualThreadsOnNexusWorker); |
| 127 | + } |
| 128 | + WorkerProperties.WorkerDeploymentConfigurationProperties workerDeploymentConfiguration = |
| 129 | + workerProperties.getDeploymentProperties(); |
| 130 | + if (workerDeploymentConfiguration != null) { |
| 131 | + WorkerDeploymentOptions.Builder opts = WorkerDeploymentOptions.newBuilder(); |
| 132 | + Optional.ofNullable(workerDeploymentConfiguration.getUseVersioning()) |
| 133 | + .ifPresent(opts::setUseVersioning); |
| 134 | + Optional.ofNullable(workerDeploymentConfiguration.getDeploymentVersion()) |
| 135 | + .ifPresent((v) -> opts.setVersion(WorkerDeploymentVersion.fromCanonicalString(v))); |
| 136 | + Optional.ofNullable(workerDeploymentConfiguration.getDefaultVersioningBehavior()) |
| 137 | + .ifPresent(opts::setDefaultVersioningBehavior); |
| 138 | + options.setDeploymentOptions(opts.build()); |
| 139 | + } |
117 | 140 | } |
118 | 141 | } |
119 | 142 |
|
|
0 commit comments