Skip to content

Commit f68c9bc

Browse files
authored
Make eager activity reservation limit configurable (#2970)
1 parent fd8b29a commit f68c9bc

9 files changed

Lines changed: 220 additions & 5 deletions

File tree

temporal-sdk/src/main/java/io/temporal/internal/worker/EagerActivitySlotsReservation.java

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,6 @@
77
import io.temporal.api.workflowservice.v1.PollActivityTaskQueueResponse;
88
import io.temporal.api.workflowservice.v1.RespondWorkflowTaskCompletedRequest;
99
import io.temporal.api.workflowservice.v1.RespondWorkflowTaskCompletedResponse;
10-
import io.temporal.internal.Config;
1110
import io.temporal.worker.tuning.SlotPermit;
1211
import java.io.Closeable;
1312
import java.util.ArrayList;
@@ -19,10 +18,13 @@
1918
@NotThreadSafe
2019
class EagerActivitySlotsReservation implements Closeable {
2120
private final EagerActivityDispatcher eagerActivityDispatcher;
21+
private final int maxReservations;
2222
private final List<SlotPermit> reservedSlots = new ArrayList<>();
2323

24-
EagerActivitySlotsReservation(EagerActivityDispatcher eagerActivityDispatcher) {
24+
EagerActivitySlotsReservation(
25+
EagerActivityDispatcher eagerActivityDispatcher, int maxReservations) {
2526
this.eagerActivityDispatcher = eagerActivityDispatcher;
27+
this.maxReservations = maxReservations;
2628
}
2729

2830
public void applyToRequest(RespondWorkflowTaskCompletedRequest.Builder mutableRequest) {
@@ -33,7 +35,7 @@ public void applyToRequest(RespondWorkflowTaskCompletedRequest.Builder mutableRe
3335
ScheduleActivityTaskCommandAttributes commandAttributes =
3436
command.getScheduleActivityTaskCommandAttributes();
3537
if (!commandAttributes.getRequestEagerExecution()) continue;
36-
boolean atLimit = this.reservedSlots.size() >= Config.EAGER_ACTIVITIES_LIMIT;
38+
boolean atLimit = this.reservedSlots.size() >= this.maxReservations;
3739
Optional<SlotPermit> permit = Optional.empty();
3840
if (!atLimit) {
3941
permit = this.eagerActivityDispatcher.tryReserveActivitySlot(commandAttributes);

temporal-sdk/src/main/java/io/temporal/internal/worker/SyncWorkflowWorker.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -67,6 +67,7 @@ public SyncWorkflowWorker(
6767
String stickyTaskQueueName,
6868
@Nonnull WorkflowThreadExecutor workflowThreadExecutor,
6969
@Nonnull EagerActivityDispatcher eagerActivityDispatcher,
70+
int maxEagerActivityReservationsPerWorkflowTask,
7071
@Nonnull SlotSupplier<WorkflowSlotInfo> slotSupplier,
7172
@Nonnull SlotSupplier<LocalActivitySlotInfo> laSlotSupplier,
7273
@Nonnull NamespaceCapabilities namespaceCapabilities) {
@@ -123,6 +124,7 @@ public SyncWorkflowWorker(
123124
cache,
124125
taskHandler,
125126
eagerActivityDispatcher,
127+
maxEagerActivityReservationsPerWorkflowTask,
126128
slotSupplier,
127129
namespaceCapabilities);
128130

temporal-sdk/src/main/java/io/temporal/internal/worker/WorkflowWorker.java

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,7 @@ final class WorkflowWorker implements SuspendableWorker {
5252
private final Scope workerMetricsScope;
5353
private final GrpcRetryer grpcRetryer;
5454
private final EagerActivityDispatcher eagerActivityDispatcher;
55+
private final int maxEagerActivityReservationsPerWorkflowTask;
5556
private final TrackingSlotSupplier<WorkflowSlotInfo> slotSupplier;
5657

5758
private final TaskCounter taskCounter = new TaskCounter();
@@ -77,6 +78,7 @@ public WorkflowWorker(
7778
@Nonnull WorkflowExecutorCache cache,
7879
@Nonnull WorkflowTaskHandler handler,
7980
@Nonnull EagerActivityDispatcher eagerActivityDispatcher,
81+
int maxEagerActivityReservationsPerWorkflowTask,
8082
@Nonnull SlotSupplier<WorkflowSlotInfo> slotSupplier,
8183
@Nonnull NamespaceCapabilities namespaceCapabilities) {
8284
this.service = Objects.requireNonNull(service);
@@ -92,6 +94,7 @@ public WorkflowWorker(
9294
this.handler = Objects.requireNonNull(handler);
9395
this.grpcRetryer = new GrpcRetryer(service.getServerCapabilities());
9496
this.eagerActivityDispatcher = eagerActivityDispatcher;
97+
this.maxEagerActivityReservationsPerWorkflowTask = maxEagerActivityReservationsPerWorkflowTask;
9598
this.slotSupplier = new TrackingSlotSupplier<>(slotSupplier, this.workerMetricsScope);
9699
this.namespaceCapabilities = namespaceCapabilities;
97100
}
@@ -478,7 +481,8 @@ public void handle(WorkflowTask task) throws Exception {
478481
RespondWorkflowTaskCompletedRequest.Builder requestBuilder =
479482
taskCompleted.toBuilder();
480483
try (EagerActivitySlotsReservation activitySlotsReservation =
481-
new EagerActivitySlotsReservation(eagerActivityDispatcher)) {
484+
new EagerActivitySlotsReservation(
485+
eagerActivityDispatcher, maxEagerActivityReservationsPerWorkflowTask)) {
482486
activitySlotsReservation.applyToRequest(requestBuilder);
483487
RespondWorkflowTaskCompletedResponse response =
484488
sendTaskCompleted(

temporal-sdk/src/main/java/io/temporal/worker/Worker.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -244,6 +244,7 @@ private static final class TaskSnapshot {
244244
stickyTaskQueueName,
245245
workflowThreadExecutor,
246246
eagerActivityDispatcher,
247+
this.options.getMaxEagerActivityReservationsPerWorkflowTask(),
247248
workflowSlotSupplier,
248249
localActivitySlotSupplier,
249250
namespaceCapabilities);

temporal-sdk/src/main/java/io/temporal/worker/WorkerOptions.java

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@
44

55
import com.google.common.base.Preconditions;
66
import io.temporal.common.Experimental;
7+
import io.temporal.internal.Config;
78
import io.temporal.serviceclient.WorkflowServiceStubsOptions;
89
import io.temporal.worker.tuning.*;
910
import java.time.Duration;
@@ -64,6 +65,7 @@ public static final class Builder {
6465
private Duration defaultHeartbeatThrottleInterval;
6566
private Duration stickyQueueScheduleToStartTimeout;
6667
private boolean disableEagerExecution;
68+
private int maxEagerActivityReservationsPerWorkflowTask = Config.EAGER_ACTIVITIES_LIMIT;
6769
private String buildId;
6870
private boolean useBuildIdForVersioning;
6971
private Duration stickyTaskQueueDrainTimeout;
@@ -109,6 +111,8 @@ private Builder(WorkerOptions o) {
109111
this.defaultHeartbeatThrottleInterval = o.defaultHeartbeatThrottleInterval;
110112
this.stickyQueueScheduleToStartTimeout = o.stickyQueueScheduleToStartTimeout;
111113
this.disableEagerExecution = o.disableEagerExecution;
114+
this.maxEagerActivityReservationsPerWorkflowTask =
115+
o.maxEagerActivityReservationsPerWorkflowTask;
112116
this.useBuildIdForVersioning = o.useBuildIdForVersioning;
113117
this.buildId = o.buildId;
114118
this.stickyTaskQueueDrainTimeout = o.stickyTaskQueueDrainTimeout;
@@ -400,6 +404,20 @@ public Builder setDisableEagerExecution(boolean disableEagerExecution) {
400404
return this;
401405
}
402406

407+
/**
408+
* Sets the maximum number of activity slots that may be reserved for eager execution when
409+
* completing a workflow task.
410+
*
411+
* <p>The default is 3. The value must be positive. To disable eager activity execution, use
412+
* {@link #setDisableEagerExecution(boolean)}.
413+
*/
414+
public Builder setMaxEagerActivityReservationsPerWorkflowTask(
415+
int maxEagerActivityReservationsPerWorkflowTask) {
416+
this.maxEagerActivityReservationsPerWorkflowTask =
417+
maxEagerActivityReservationsPerWorkflowTask;
418+
return this;
419+
}
420+
403421
/**
404422
* Opts the worker in to the Build-ID-based versioning feature. This ensures that the worker
405423
* will only receive tasks which it is compatible with.
@@ -623,6 +641,7 @@ public WorkerOptions build() {
623641
defaultHeartbeatThrottleInterval,
624642
stickyQueueScheduleToStartTimeout,
625643
disableEagerExecution,
644+
maxEagerActivityReservationsPerWorkflowTask,
626645
useBuildIdForVersioning,
627646
buildId,
628647
stickyTaskQueueDrainTimeout,
@@ -647,6 +666,10 @@ public WorkerOptions validateAndBuildWithDefaults() {
647666
maxWorkerActivitiesPerSecond >= 0, "negative maxActivitiesPerSecond");
648667
Preconditions.checkState(
649668
maxConcurrentActivityExecutionSize >= 0, "negative maxConcurrentActivityExecutionSize");
669+
Preconditions.checkState(
670+
maxEagerActivityReservationsPerWorkflowTask > 0,
671+
"maxEagerActivityReservationsPerWorkflowTask must be positive; use "
672+
+ "setDisableEagerExecution(true) to disable eager activity execution");
650673
Preconditions.checkState(
651674
maxConcurrentWorkflowTaskExecutionSize >= 0,
652675
"negative maxConcurrentWorkflowTaskExecutionSize");
@@ -758,6 +781,7 @@ public WorkerOptions validateAndBuildWithDefaults() {
758781
? DEFAULT_STICKY_SCHEDULE_TO_START_TIMEOUT
759782
: stickyQueueScheduleToStartTimeout,
760783
disableEagerExecution,
784+
maxEagerActivityReservationsPerWorkflowTask,
761785
useBuildIdForVersioning,
762786
buildId,
763787
stickyTaskQueueDrainTimeout == null
@@ -796,6 +820,7 @@ public WorkerOptions validateAndBuildWithDefaults() {
796820
private final Duration defaultHeartbeatThrottleInterval;
797821
private final @Nonnull Duration stickyQueueScheduleToStartTimeout;
798822
private final boolean disableEagerExecution;
823+
private final int maxEagerActivityReservationsPerWorkflowTask;
799824
private final boolean useBuildIdForVersioning;
800825
private final String buildId;
801826
private final Duration stickyTaskQueueDrainTimeout;
@@ -831,6 +856,7 @@ private WorkerOptions(
831856
Duration defaultHeartbeatThrottleInterval,
832857
@Nonnull Duration stickyQueueScheduleToStartTimeout,
833858
boolean disableEagerExecution,
859+
int maxEagerActivityReservationsPerWorkflowTask,
834860
boolean useBuildIdForVersioning,
835861
String buildId,
836862
Duration stickyTaskQueueDrainTimeout,
@@ -864,6 +890,7 @@ private WorkerOptions(
864890
this.defaultHeartbeatThrottleInterval = defaultHeartbeatThrottleInterval;
865891
this.stickyQueueScheduleToStartTimeout = stickyQueueScheduleToStartTimeout;
866892
this.disableEagerExecution = maxTaskQueueActivitiesPerSecond > 0 ? true : disableEagerExecution;
893+
this.maxEagerActivityReservationsPerWorkflowTask = maxEagerActivityReservationsPerWorkflowTask;
867894
this.useBuildIdForVersioning = useBuildIdForVersioning;
868895
this.buildId = buildId;
869896
this.stickyTaskQueueDrainTimeout = stickyTaskQueueDrainTimeout;
@@ -989,6 +1016,10 @@ public boolean isEagerExecutionDisabled() {
9891016
return disableEagerExecution;
9901017
}
9911018

1019+
public int getMaxEagerActivityReservationsPerWorkflowTask() {
1020+
return maxEagerActivityReservationsPerWorkflowTask;
1021+
}
1022+
9921023
public boolean isUsingBuildIdForVersioning() {
9931024
return useBuildIdForVersioning;
9941025
}
@@ -1070,6 +1101,8 @@ && compare(maxTaskQueueActivitiesPerSecond, that.maxTaskQueueActivitiesPerSecond
10701101
&& localActivityWorkerOnly == that.localActivityWorkerOnly
10711102
&& defaultDeadlockDetectionTimeout == that.defaultDeadlockDetectionTimeout
10721103
&& disableEagerExecution == that.disableEagerExecution
1104+
&& maxEagerActivityReservationsPerWorkflowTask
1105+
== that.maxEagerActivityReservationsPerWorkflowTask
10731106
&& useBuildIdForVersioning == that.useBuildIdForVersioning
10741107
&& Objects.equals(workerTuner, that.workerTuner)
10751108
&& Objects.equals(maxHeartbeatThrottleInterval, that.maxHeartbeatThrottleInterval)
@@ -1109,6 +1142,7 @@ public int hashCode() {
11091142
defaultHeartbeatThrottleInterval,
11101143
stickyQueueScheduleToStartTimeout,
11111144
disableEagerExecution,
1145+
maxEagerActivityReservationsPerWorkflowTask,
11121146
useBuildIdForVersioning,
11131147
buildId,
11141148
stickyTaskQueueDrainTimeout,
@@ -1160,6 +1194,8 @@ public String toString() {
11601194
+ stickyQueueScheduleToStartTimeout
11611195
+ ", disableEagerExecution="
11621196
+ disableEagerExecution
1197+
+ ", maxEagerActivityReservationsPerWorkflowTask="
1198+
+ maxEagerActivityReservationsPerWorkflowTask
11631199
+ ", useBuildIdForVersioning="
11641200
+ useBuildIdForVersioning
11651201
+ ", buildId='"
Lines changed: 60 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,60 @@
1+
package io.temporal.internal.worker;
2+
3+
import static org.junit.Assert.assertEquals;
4+
import static org.junit.Assert.assertFalse;
5+
import static org.junit.Assert.assertTrue;
6+
import static org.mockito.ArgumentMatchers.any;
7+
import static org.mockito.Mockito.mock;
8+
import static org.mockito.Mockito.times;
9+
import static org.mockito.Mockito.verify;
10+
import static org.mockito.Mockito.when;
11+
12+
import io.temporal.api.command.v1.Command;
13+
import io.temporal.api.command.v1.ScheduleActivityTaskCommandAttributes;
14+
import io.temporal.api.enums.v1.CommandType;
15+
import io.temporal.api.workflowservice.v1.RespondWorkflowTaskCompletedRequest;
16+
import io.temporal.worker.tuning.SlotPermit;
17+
import java.util.Optional;
18+
import org.junit.Test;
19+
20+
public class EagerActivitySlotsReservationTest {
21+
@Test
22+
public void limitsReservationsPerWorkflowTask() {
23+
EagerActivityDispatcher dispatcher = mock(EagerActivityDispatcher.class);
24+
when(dispatcher.tryReserveActivitySlot(any())).thenReturn(Optional.of(mock(SlotPermit.class)));
25+
RespondWorkflowTaskCompletedRequest.Builder request =
26+
RespondWorkflowTaskCompletedRequest.newBuilder();
27+
for (int i = 0; i < 5; i++) {
28+
request.addCommands(
29+
Command.newBuilder()
30+
.setCommandType(CommandType.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK)
31+
.setScheduleActivityTaskCommandAttributes(
32+
ScheduleActivityTaskCommandAttributes.newBuilder()
33+
.setRequestEagerExecution(true)));
34+
}
35+
36+
try (EagerActivitySlotsReservation reservation =
37+
new EagerActivitySlotsReservation(dispatcher, 2)) {
38+
reservation.applyToRequest(request);
39+
assertEquals(5, request.getCommandsCount());
40+
assertTrue(
41+
request
42+
.getCommands(0)
43+
.getScheduleActivityTaskCommandAttributes()
44+
.getRequestEagerExecution());
45+
assertTrue(
46+
request
47+
.getCommands(1)
48+
.getScheduleActivityTaskCommandAttributes()
49+
.getRequestEagerExecution());
50+
for (int i = 2; i < 5; i++) {
51+
assertFalse(
52+
request
53+
.getCommands(i)
54+
.getScheduleActivityTaskCommandAttributes()
55+
.getRequestEagerExecution());
56+
}
57+
}
58+
verify(dispatcher, times(2)).tryReserveActivitySlot(any());
59+
}
60+
}

temporal-sdk/src/test/java/io/temporal/internal/worker/WorkflowWorkerTest.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -84,6 +84,7 @@ public void concurrentPollRequestLockTest() throws Exception {
8484
cache,
8585
taskHandler,
8686
eagerActivityDispatcher,
87+
3,
8788
slotSupplier,
8889
new NamespaceCapabilities());
8990

@@ -255,6 +256,7 @@ public void respondWorkflowTaskFailureMetricTest() throws Exception {
255256
cache,
256257
taskHandler,
257258
eagerActivityDispatcher,
259+
3,
258260
slotSupplier,
259261
new NamespaceCapabilities());
260262

@@ -399,6 +401,7 @@ public boolean isAnyTypeSupported() {
399401
cache,
400402
taskHandler,
401403
eagerActivityDispatcher,
404+
3,
402405
slotSupplier,
403406
new NamespaceCapabilities());
404407

temporal-sdk/src/test/java/io/temporal/worker/WorkerOptionsTest.java

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ public void build() {
2323
private void verifyBuild(WorkerOptions options) {
2424
assertEquals(10, options.getMaxConcurrentActivityExecutionSize());
2525
assertEquals(11, options.getMaxConcurrentLocalActivityExecutionSize());
26+
assertEquals(3, options.getMaxEagerActivityReservationsPerWorkflowTask());
2627
assertNotNull(options.getPreferredVersionProvider());
2728
}
2829

@@ -55,6 +56,7 @@ public void verifyNewBuilderFromExistingWorkerOptions() {
5556
.setDefaultHeartbeatThrottleInterval(Duration.ofSeconds(7))
5657
.setStickyQueueScheduleToStartTimeout(Duration.ofSeconds(60))
5758
.setDisableEagerExecution(false)
59+
.setMaxEagerActivityReservationsPerWorkflowTask(17)
5860
.setUseBuildIdForVersioning(false)
5961
.setBuildId("build-id")
6062
.setStickyTaskQueueDrainTimeout(Duration.ofSeconds(15))
@@ -90,6 +92,9 @@ public void verifyNewBuilderFromExistingWorkerOptions() {
9092
assertEquals(
9193
w1.getStickyQueueScheduleToStartTimeout(), w2.getStickyQueueScheduleToStartTimeout());
9294
assertEquals(w1.isEagerExecutionDisabled(), w2.isEagerExecutionDisabled());
95+
assertEquals(
96+
w1.getMaxEagerActivityReservationsPerWorkflowTask(),
97+
w2.getMaxEagerActivityReservationsPerWorkflowTask());
9398
assertEquals(w1.isUsingBuildIdForVersioning(), w2.isUsingBuildIdForVersioning());
9499
assertEquals(w1.getBuildId(), w2.getBuildId());
95100
assertEquals(w1.getStickyTaskQueueDrainTimeout(), w2.getStickyTaskQueueDrainTimeout());
@@ -230,4 +235,21 @@ public void verifyMaxTaskQueuePerSecondsDisablesEagerExecution() {
230235
WorkerOptions w2 = WorkerOptions.newBuilder().setMaxTaskQueueActivitiesPerSecond(2.0).build();
231236
assertTrue(w2.isEagerExecutionDisabled());
232237
}
238+
239+
@Test
240+
public void rejectsNonPositiveMaxEagerActivityReservationsPerWorkflowTask() {
241+
for (int value : new int[] {0, -1}) {
242+
IllegalStateException exception =
243+
assertThrows(
244+
IllegalStateException.class,
245+
() ->
246+
WorkerOptions.newBuilder()
247+
.setMaxEagerActivityReservationsPerWorkflowTask(value)
248+
.validateAndBuildWithDefaults());
249+
assertEquals(
250+
"maxEagerActivityReservationsPerWorkflowTask must be positive; use "
251+
+ "setDisableEagerExecution(true) to disable eager activity execution",
252+
exception.getMessage());
253+
}
254+
}
233255
}

0 commit comments

Comments
 (0)