Skip to content

Commit 1e1ad53

Browse files
committed
Expose eager activity concurrency limit
1 parent 4804646 commit 1e1ad53

6 files changed

Lines changed: 133 additions & 11 deletions

File tree

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

Lines changed: 62 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -218,8 +218,9 @@ public WorkerLifecycleState getLifecycleState() {
218218
return poller.getLifecycleState();
219219
}
220220

221-
public EagerActivityDispatcher getEagerActivityDispatcher() {
222-
return new EagerActivityDispatcherImpl();
221+
public EagerActivityDispatcher getEagerActivityDispatcher(
222+
int maxConcurrentEagerActivityExecutionSize) {
223+
return new EagerActivityDispatcherImpl(maxConcurrentEagerActivityExecutionSize);
223224
}
224225

225226
private PollerOptions getPollerOptions(SingleWorkerOptions options) {
@@ -485,6 +486,13 @@ private void logExceptionDuringResultReporting(
485486
}
486487

487488
private final class EagerActivityDispatcherImpl implements EagerActivityDispatcher {
489+
private final EagerActivitySlotLimiter eagerActivitySlotLimiter;
490+
491+
private EagerActivityDispatcherImpl(int maxConcurrentEagerActivityExecutionSize) {
492+
this.eagerActivitySlotLimiter =
493+
new EagerActivitySlotLimiter(maxConcurrentEagerActivityExecutionSize);
494+
}
495+
488496
@Override
489497
public Optional<SlotPermit> tryReserveActivitySlot(
490498
ScheduleActivityTaskCommandAttributesOrBuilder commandAttributes) {
@@ -493,15 +501,31 @@ public Optional<SlotPermit> tryReserveActivitySlot(
493501
commandAttributes.getTaskQueue().getName(), ActivityWorker.this.taskQueue)) {
494502
return Optional.empty();
495503
}
496-
return ActivityWorker.this.slotSupplier.tryReserveSlot(
497-
new SlotReservationData(
498-
ActivityWorker.this.taskQueue, options.getIdentity(), options.getBuildId()));
504+
if (!eagerActivitySlotLimiter.tryReserve()) {
505+
return Optional.empty();
506+
}
507+
Optional<SlotPermit> permit = Optional.empty();
508+
try {
509+
permit =
510+
ActivityWorker.this.slotSupplier.tryReserveSlot(
511+
new SlotReservationData(
512+
ActivityWorker.this.taskQueue, options.getIdentity(), options.getBuildId()));
513+
return permit;
514+
} finally {
515+
if (!permit.isPresent()) {
516+
eagerActivitySlotLimiter.release();
517+
}
518+
}
499519
}
500520

501521
@Override
502522
public void releaseActivitySlotReservations(Iterable<SlotPermit> permits) {
503523
for (SlotPermit permit : permits) {
504-
ActivityWorker.this.slotSupplier.releaseSlot(SlotReleaseReason.neverUsed(), permit);
524+
try {
525+
ActivityWorker.this.slotSupplier.releaseSlot(SlotReleaseReason.neverUsed(), permit);
526+
} finally {
527+
eagerActivitySlotLimiter.release();
528+
}
505529
}
506530
}
507531

@@ -511,9 +535,39 @@ public void dispatchActivity(PollActivityTaskQueueResponse activity, SlotPermit
511535
new ActivityTask(
512536
activity,
513537
permit,
514-
() ->
538+
() -> {
539+
try {
515540
ActivityWorker.this.slotSupplier.releaseSlot(
516-
SlotReleaseReason.taskComplete(), permit)));
541+
SlotReleaseReason.taskComplete(), permit);
542+
} finally {
543+
eagerActivitySlotLimiter.release();
544+
}
545+
}));
546+
}
547+
}
548+
549+
static final class EagerActivitySlotLimiter {
550+
private final int maxConcurrent;
551+
private int heldSlotCount;
552+
553+
EagerActivitySlotLimiter(int maxConcurrent) {
554+
this.maxConcurrent = maxConcurrent;
555+
}
556+
557+
synchronized boolean tryReserve() {
558+
if (maxConcurrent > 0 && heldSlotCount >= maxConcurrent) {
559+
return false;
560+
}
561+
heldSlotCount++;
562+
return true;
563+
}
564+
565+
synchronized void release() {
566+
if (heldSlotCount <= 0) {
567+
throw new IllegalStateException(
568+
"Trying to release an unreserved eager activity slot. This is an SDK bug.");
569+
}
570+
heldSlotCount--;
517571
}
518572
}
519573
}

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -163,8 +163,9 @@ public WorkerLifecycleState getLifecycleState() {
163163
}
164164
}
165165

166-
public EagerActivityDispatcher getEagerActivityDispatcher() {
167-
return this.worker.getEagerActivityDispatcher();
166+
public EagerActivityDispatcher getEagerActivityDispatcher(
167+
int maxConcurrentEagerActivityExecutionSize) {
168+
return this.worker.getEagerActivityDispatcher(maxConcurrentEagerActivityExecutionSize);
168169
}
169170

170171
public boolean isAnyTypeSupported() {

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

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -173,7 +173,8 @@ private static final class TaskSnapshot {
173173

174174
EagerActivityDispatcher eagerActivityDispatcher =
175175
(activityWorker != null && !this.options.isEagerExecutionDisabled())
176-
? activityWorker.getEagerActivityDispatcher()
176+
? activityWorker.getEagerActivityDispatcher(
177+
this.options.getMaxConcurrentEagerActivityExecutionSize())
177178
: new EagerActivityDispatcher.NoopEagerActivityDispatcher();
178179

179180
SingleWorkerOptions nexusOptions =

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

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,7 @@ public static final class Builder {
6464
private Duration defaultHeartbeatThrottleInterval;
6565
private Duration stickyQueueScheduleToStartTimeout;
6666
private boolean disableEagerExecution;
67+
private int maxConcurrentEagerActivityExecutionSize;
6768
private String buildId;
6869
private boolean useBuildIdForVersioning;
6970
private Duration stickyTaskQueueDrainTimeout;
@@ -109,6 +110,7 @@ private Builder(WorkerOptions o) {
109110
this.defaultHeartbeatThrottleInterval = o.defaultHeartbeatThrottleInterval;
110111
this.stickyQueueScheduleToStartTimeout = o.stickyQueueScheduleToStartTimeout;
111112
this.disableEagerExecution = o.disableEagerExecution;
113+
this.maxConcurrentEagerActivityExecutionSize = o.maxConcurrentEagerActivityExecutionSize;
112114
this.useBuildIdForVersioning = o.useBuildIdForVersioning;
113115
this.buildId = o.buildId;
114116
this.stickyTaskQueueDrainTimeout = o.stickyTaskQueueDrainTimeout;
@@ -400,6 +402,19 @@ public Builder setDisableEagerExecution(boolean disableEagerExecution) {
400402
return this;
401403
}
402404

405+
/**
406+
* Sets the maximum number of eager activities that can be running concurrently.
407+
*
408+
* <p>When nonzero, eager activity execution will not be requested if it would cause the number
409+
* of running eager activities to exceed this value. The default of zero means unlimited and
410+
* therefore only bound by the activity slot supplier.
411+
*/
412+
public Builder setMaxConcurrentEagerActivityExecutionSize(
413+
int maxConcurrentEagerActivityExecutionSize) {
414+
this.maxConcurrentEagerActivityExecutionSize = maxConcurrentEagerActivityExecutionSize;
415+
return this;
416+
}
417+
403418
/**
404419
* Opts the worker in to the Build-ID-based versioning feature. This ensures that the worker
405420
* will only receive tasks which it is compatible with.
@@ -623,6 +638,7 @@ public WorkerOptions build() {
623638
defaultHeartbeatThrottleInterval,
624639
stickyQueueScheduleToStartTimeout,
625640
disableEagerExecution,
641+
maxConcurrentEagerActivityExecutionSize,
626642
useBuildIdForVersioning,
627643
buildId,
628644
stickyTaskQueueDrainTimeout,
@@ -647,6 +663,9 @@ public WorkerOptions validateAndBuildWithDefaults() {
647663
maxWorkerActivitiesPerSecond >= 0, "negative maxActivitiesPerSecond");
648664
Preconditions.checkState(
649665
maxConcurrentActivityExecutionSize >= 0, "negative maxConcurrentActivityExecutionSize");
666+
Preconditions.checkState(
667+
maxConcurrentEagerActivityExecutionSize >= 0,
668+
"negative maxConcurrentEagerActivityExecutionSize");
650669
Preconditions.checkState(
651670
maxConcurrentWorkflowTaskExecutionSize >= 0,
652671
"negative maxConcurrentWorkflowTaskExecutionSize");
@@ -758,6 +777,7 @@ public WorkerOptions validateAndBuildWithDefaults() {
758777
? DEFAULT_STICKY_SCHEDULE_TO_START_TIMEOUT
759778
: stickyQueueScheduleToStartTimeout,
760779
disableEagerExecution,
780+
maxConcurrentEagerActivityExecutionSize,
761781
useBuildIdForVersioning,
762782
buildId,
763783
stickyTaskQueueDrainTimeout == null
@@ -796,6 +816,7 @@ public WorkerOptions validateAndBuildWithDefaults() {
796816
private final Duration defaultHeartbeatThrottleInterval;
797817
private final @Nonnull Duration stickyQueueScheduleToStartTimeout;
798818
private final boolean disableEagerExecution;
819+
private final int maxConcurrentEagerActivityExecutionSize;
799820
private final boolean useBuildIdForVersioning;
800821
private final String buildId;
801822
private final Duration stickyTaskQueueDrainTimeout;
@@ -831,6 +852,7 @@ private WorkerOptions(
831852
Duration defaultHeartbeatThrottleInterval,
832853
@Nonnull Duration stickyQueueScheduleToStartTimeout,
833854
boolean disableEagerExecution,
855+
int maxConcurrentEagerActivityExecutionSize,
834856
boolean useBuildIdForVersioning,
835857
String buildId,
836858
Duration stickyTaskQueueDrainTimeout,
@@ -864,6 +886,7 @@ private WorkerOptions(
864886
this.defaultHeartbeatThrottleInterval = defaultHeartbeatThrottleInterval;
865887
this.stickyQueueScheduleToStartTimeout = stickyQueueScheduleToStartTimeout;
866888
this.disableEagerExecution = maxTaskQueueActivitiesPerSecond > 0 ? true : disableEagerExecution;
889+
this.maxConcurrentEagerActivityExecutionSize = maxConcurrentEagerActivityExecutionSize;
867890
this.useBuildIdForVersioning = useBuildIdForVersioning;
868891
this.buildId = buildId;
869892
this.stickyTaskQueueDrainTimeout = stickyTaskQueueDrainTimeout;
@@ -989,6 +1012,10 @@ public boolean isEagerExecutionDisabled() {
9891012
return disableEagerExecution;
9901013
}
9911014

1015+
public int getMaxConcurrentEagerActivityExecutionSize() {
1016+
return maxConcurrentEagerActivityExecutionSize;
1017+
}
1018+
9921019
public boolean isUsingBuildIdForVersioning() {
9931020
return useBuildIdForVersioning;
9941021
}
@@ -1070,6 +1097,7 @@ && compare(maxTaskQueueActivitiesPerSecond, that.maxTaskQueueActivitiesPerSecond
10701097
&& localActivityWorkerOnly == that.localActivityWorkerOnly
10711098
&& defaultDeadlockDetectionTimeout == that.defaultDeadlockDetectionTimeout
10721099
&& disableEagerExecution == that.disableEagerExecution
1100+
&& maxConcurrentEagerActivityExecutionSize == that.maxConcurrentEagerActivityExecutionSize
10731101
&& useBuildIdForVersioning == that.useBuildIdForVersioning
10741102
&& Objects.equals(workerTuner, that.workerTuner)
10751103
&& Objects.equals(maxHeartbeatThrottleInterval, that.maxHeartbeatThrottleInterval)
@@ -1109,6 +1137,7 @@ public int hashCode() {
11091137
defaultHeartbeatThrottleInterval,
11101138
stickyQueueScheduleToStartTimeout,
11111139
disableEagerExecution,
1140+
maxConcurrentEagerActivityExecutionSize,
11121141
useBuildIdForVersioning,
11131142
buildId,
11141143
stickyTaskQueueDrainTimeout,
@@ -1160,6 +1189,8 @@ public String toString() {
11601189
+ stickyQueueScheduleToStartTimeout
11611190
+ ", disableEagerExecution="
11621191
+ disableEagerExecution
1192+
+ ", maxConcurrentEagerActivityExecutionSize="
1193+
+ maxConcurrentEagerActivityExecutionSize
11631194
+ ", useBuildIdForVersioning="
11641195
+ useBuildIdForVersioning
11651196
+ ", buildId='"
Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,31 @@
1+
package io.temporal.internal.worker;
2+
3+
import static org.junit.Assert.assertFalse;
4+
import static org.junit.Assert.assertTrue;
5+
6+
import org.junit.Test;
7+
8+
public class EagerActivitySlotLimiterTest {
9+
@Test
10+
public void enforcesLimitUntilReservationIsReleased() {
11+
ActivityWorker.EagerActivitySlotLimiter limiter =
12+
new ActivityWorker.EagerActivitySlotLimiter(2);
13+
14+
assertTrue(limiter.tryReserve());
15+
assertTrue(limiter.tryReserve());
16+
assertFalse(limiter.tryReserve());
17+
18+
limiter.release();
19+
assertTrue(limiter.tryReserve());
20+
}
21+
22+
@Test
23+
public void zeroAllowsUnlimitedReservations() {
24+
ActivityWorker.EagerActivitySlotLimiter limiter =
25+
new ActivityWorker.EagerActivitySlotLimiter(0);
26+
27+
for (int i = 0; i < 1000; i++) {
28+
assertTrue(limiter.tryReserve());
29+
}
30+
}
31+
}

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

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,7 @@ public void verifyNewBuilderFromExistingWorkerOptions() {
5555
.setDefaultHeartbeatThrottleInterval(Duration.ofSeconds(7))
5656
.setStickyQueueScheduleToStartTimeout(Duration.ofSeconds(60))
5757
.setDisableEagerExecution(false)
58+
.setMaxConcurrentEagerActivityExecutionSize(17)
5859
.setUseBuildIdForVersioning(false)
5960
.setBuildId("build-id")
6061
.setStickyTaskQueueDrainTimeout(Duration.ofSeconds(15))
@@ -90,6 +91,9 @@ public void verifyNewBuilderFromExistingWorkerOptions() {
9091
assertEquals(
9192
w1.getStickyQueueScheduleToStartTimeout(), w2.getStickyQueueScheduleToStartTimeout());
9293
assertEquals(w1.isEagerExecutionDisabled(), w2.isEagerExecutionDisabled());
94+
assertEquals(
95+
w1.getMaxConcurrentEagerActivityExecutionSize(),
96+
w2.getMaxConcurrentEagerActivityExecutionSize());
9397
assertEquals(w1.isUsingBuildIdForVersioning(), w2.isUsingBuildIdForVersioning());
9498
assertEquals(w1.getBuildId(), w2.getBuildId());
9599
assertEquals(w1.getStickyTaskQueueDrainTimeout(), w2.getStickyTaskQueueDrainTimeout());

0 commit comments

Comments
 (0)