Skip to content

Commit dae5e0b

Browse files
authored
Add option to let activities heartbeat during worker shutdown (#2903)
1 parent a1b6fff commit dae5e0b

7 files changed

Lines changed: 238 additions & 18 deletions

File tree

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

Lines changed: 18 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,7 @@ public static final class Builder {
4141
private boolean usingVirtualThreads;
4242
private WorkerDeploymentOptions deploymentOptions;
4343
private String workerInstanceKey;
44+
private boolean allowActivityHeartbeatDuringShutdown;
4445

4546
private Builder() {}
4647

@@ -66,6 +67,7 @@ private Builder(SingleWorkerOptions options) {
6667
this.usingVirtualThreads = options.isUsingVirtualThreads();
6768
this.deploymentOptions = options.getDeploymentOptions();
6869
this.workerInstanceKey = options.getWorkerInstanceKey();
70+
this.allowActivityHeartbeatDuringShutdown = options.getAllowActivityHeartbeatDuringShutdown();
6971
}
7072

7173
public Builder setIdentity(String identity) {
@@ -162,6 +164,12 @@ public Builder setWorkerInstanceKey(String workerInstanceKey) {
162164
return this;
163165
}
164166

167+
public Builder setAllowActivityHeartbeatDuringShutdown(
168+
boolean allowActivityHeartbeatDuringShutdown) {
169+
this.allowActivityHeartbeatDuringShutdown = allowActivityHeartbeatDuringShutdown;
170+
return this;
171+
}
172+
165173
public SingleWorkerOptions build() {
166174
PollerOptions pollerOptions = this.pollerOptions;
167175
if (pollerOptions == null) {
@@ -201,7 +209,8 @@ public SingleWorkerOptions build() {
201209
drainStickyTaskQueueTimeout,
202210
usingVirtualThreads,
203211
this.deploymentOptions,
204-
this.workerInstanceKey);
212+
this.workerInstanceKey,
213+
this.allowActivityHeartbeatDuringShutdown);
205214
}
206215
}
207216

@@ -223,6 +232,7 @@ public SingleWorkerOptions build() {
223232
private final boolean usingVirtualThreads;
224233
private final WorkerDeploymentOptions deploymentOptions;
225234
private final String workerInstanceKey;
235+
private final boolean allowActivityHeartbeatDuringShutdown;
226236

227237
private SingleWorkerOptions(
228238
String identity,
@@ -242,7 +252,8 @@ private SingleWorkerOptions(
242252
Duration drainStickyTaskQueueTimeout,
243253
boolean usingVirtualThreads,
244254
WorkerDeploymentOptions deploymentOptions,
245-
String workerInstanceKey) {
255+
String workerInstanceKey,
256+
boolean allowActivityHeartbeatDuringShutdown) {
246257
this.identity = identity;
247258
this.binaryChecksum = binaryChecksum;
248259
this.buildId = buildId;
@@ -261,6 +272,7 @@ private SingleWorkerOptions(
261272
this.usingVirtualThreads = usingVirtualThreads;
262273
this.deploymentOptions = deploymentOptions;
263274
this.workerInstanceKey = workerInstanceKey;
275+
this.allowActivityHeartbeatDuringShutdown = allowActivityHeartbeatDuringShutdown;
264276
}
265277

266278
public String getIdentity() {
@@ -291,6 +303,10 @@ public Duration getDrainStickyTaskQueueTimeout() {
291303
return drainStickyTaskQueueTimeout;
292304
}
293305

306+
public boolean getAllowActivityHeartbeatDuringShutdown() {
307+
return allowActivityHeartbeatDuringShutdown;
308+
}
309+
294310
public DataConverter getDataConverter() {
295311
return dataConverter;
296312
}

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

Lines changed: 27 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ public class SyncActivityWorker implements SuspendableWorker {
2626
private final ScheduledExecutorService heartbeatExecutor;
2727
private final ActivityTaskHandlerImpl taskHandler;
2828
private final ActivityWorker worker;
29+
private final boolean allowActivityHeartbeatDuringShutdown;
2930

3031
public SyncActivityWorker(
3132
WorkflowClient client,
@@ -38,6 +39,7 @@ public SyncActivityWorker(
3839
this.identity = options.getIdentity();
3940
this.namespace = namespace;
4041
this.taskQueue = taskQueue;
42+
this.allowActivityHeartbeatDuringShutdown = options.getAllowActivityHeartbeatDuringShutdown();
4143

4244
this.heartbeatExecutor =
4345
Executors.newScheduledThreadPool(
@@ -89,16 +91,31 @@ public boolean start() {
8991

9092
@Override
9193
public CompletableFuture<Void> shutdown(ShutdownManager shutdownManager, boolean interruptTasks) {
92-
return shutdownManager
93-
// we want to shut down heartbeatExecutor before activity worker, so in-flight activities
94-
// could get an ActivityWorkerShutdownException from their heartbeat
95-
.shutdownExecutor(heartbeatExecutor, this + "#heartbeatExecutor", Duration.ofSeconds(5))
96-
.thenCompose(r -> worker.shutdown(shutdownManager, interruptTasks))
97-
.exceptionally(
98-
e -> {
99-
log.error("[BUG] Unexpected exception during shutdown", e);
100-
return null;
101-
});
94+
CompletableFuture<Void> shutdownFuture;
95+
if (allowActivityHeartbeatDuringShutdown && !interruptTasks) {
96+
// we want to shut down heartbeatExecutor only after all outstanding activity tasks have
97+
// finished executing, so in-flight activities can keep heartbeating during the shutdown
98+
shutdownFuture =
99+
worker
100+
.shutdown(shutdownManager, interruptTasks)
101+
.thenCompose(r -> shutdownHeartbeatExecutor(shutdownManager));
102+
} else {
103+
// we want to shut down heartbeatExecutor before activity worker, so in-flight activities
104+
// could get an ActivityWorkerShutdownException from their heartbeat
105+
shutdownFuture =
106+
shutdownHeartbeatExecutor(shutdownManager)
107+
.thenCompose(r -> worker.shutdown(shutdownManager, interruptTasks));
108+
}
109+
return shutdownFuture.exceptionally(
110+
e -> {
111+
log.error("[BUG] Unexpected exception during shutdown", e);
112+
return null;
113+
});
114+
}
115+
116+
private CompletableFuture<Void> shutdownHeartbeatExecutor(ShutdownManager shutdownManager) {
117+
return shutdownManager.shutdownExecutor(
118+
heartbeatExecutor, this + "#heartbeatExecutor", Duration.ofSeconds(5));
102119
}
103120

104121
@Override

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -893,6 +893,7 @@ private static SingleWorkerOptions toActivityOptions(
893893
return toSingleWorkerOptions(
894894
factoryOptions, options, clientOptions, contextPropagators, workerInstanceKey)
895895
.setUsingVirtualThreads(options.isUsingVirtualThreadsOnActivityWorker())
896+
.setAllowActivityHeartbeatDuringShutdown(options.getAllowActivityHeartbeatDuringShutdown())
896897
.setPollerOptions(
897898
PollerOptions.newBuilder()
898899
.setMaximumPollRatePerSecond(options.getMaxWorkerActivitiesPerSecond())

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

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -358,7 +358,9 @@ public WorkflowClient getWorkflowClient() {
358358
* activity tasks are executed. <br>
359359
* After the shutdown, calls to {@link
360360
* io.temporal.activity.ActivityExecutionContext#heartbeat(Object)} start throwing {@link
361-
* io.temporal.client.ActivityWorkerShutdownException}.<br>
361+
* io.temporal.client.ActivityWorkerShutdownException}, unless {@link
362+
* WorkerOptions.Builder#setAllowActivityHeartbeatDuringShutdown(boolean)} is enabled, in which
363+
* case heartbeats keep working until the activity tasks finish executing.<br>
362364
* This method does not wait for the shutdown to complete. Use {@link #awaitTermination(long,
363365
* TimeUnit)} to do that.<br>
364366
* Invocation has no additional effect if already shut down.

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

Lines changed: 43 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -77,6 +77,7 @@ public static final class Builder {
7777
private PollerBehavior workflowTaskPollersBehavior;
7878
private PollerBehavior activityTaskPollersBehavior;
7979
private PollerBehavior nexusTaskPollersBehavior;
80+
private boolean allowActivityHeartbeatDuringShutdown;
8081

8182
private Builder() {}
8283

@@ -112,6 +113,7 @@ private Builder(WorkerOptions o) {
112113
this.workflowTaskPollersBehavior = o.workflowTaskPollersBehavior;
113114
this.activityTaskPollersBehavior = o.activityTaskPollersBehavior;
114115
this.nexusTaskPollersBehavior = o.nexusTaskPollersBehavior;
116+
this.allowActivityHeartbeatDuringShutdown = o.allowActivityHeartbeatDuringShutdown;
115117
}
116118

117119
/**
@@ -524,6 +526,28 @@ public Builder setNexusTaskPollersBehavior(PollerBehavior pollerBehavior) {
524526
return this;
525527
}
526528

529+
/**
530+
* If true, activities can keep heartbeating during graceful worker shutdown (see {@link
531+
* io.temporal.worker.WorkerFactory#shutdown WorkerFactory.shutdown}). Defaults to false, which
532+
* means that after graceful shutdown is requested, calling {@link
533+
* io.temporal.activity.ActivityExecutionContext#heartbeat ActivityExecutionContext.heartbeat}
534+
* does not send a heartbeat and instead throws {@link
535+
* io.temporal.client.ActivityWorkerShutdownException ActivityWorkerShutdownException}. This
536+
* option is ignored by non-graceful shutdown (see {@link
537+
* io.temporal.worker.WorkerFactory#shutdownNow WorkerFactory.shutdownNow}).
538+
*
539+
* <p>Note that with this option enabled, activities are no longer notified of the worker
540+
* shutdown by the {@link io.temporal.client.ActivityWorkerShutdownException
541+
* ActivityWorkerShutdownException} exception, so they are expected to complete within the
542+
* termination grace period on their own.
543+
*/
544+
@Experimental
545+
public Builder setAllowActivityHeartbeatDuringShutdown(
546+
boolean allowActivityHeartbeatDuringShutdown) {
547+
this.allowActivityHeartbeatDuringShutdown = allowActivityHeartbeatDuringShutdown;
548+
return this;
549+
}
550+
527551
public WorkerOptions build() {
528552
return new WorkerOptions(
529553
maxWorkerActivitiesPerSecond,
@@ -553,7 +577,8 @@ public WorkerOptions build() {
553577
deploymentOptions,
554578
workflowTaskPollersBehavior,
555579
activityTaskPollersBehavior,
556-
nexusTaskPollersBehavior);
580+
nexusTaskPollersBehavior,
581+
allowActivityHeartbeatDuringShutdown);
557582
}
558583

559584
public WorkerOptions validateAndBuildWithDefaults() {
@@ -685,7 +710,8 @@ public WorkerOptions validateAndBuildWithDefaults() {
685710
deploymentOptions,
686711
workflowTaskPollersBehavior,
687712
activityTaskPollersBehavior,
688-
nexusTaskPollersBehavior);
713+
nexusTaskPollersBehavior,
714+
allowActivityHeartbeatDuringShutdown);
689715
}
690716
}
691717

@@ -717,6 +743,7 @@ public WorkerOptions validateAndBuildWithDefaults() {
717743
private final PollerBehavior workflowTaskPollersBehavior;
718744
private final PollerBehavior activityTaskPollersBehavior;
719745
private final PollerBehavior nexusTaskPollersBehavior;
746+
private final boolean allowActivityHeartbeatDuringShutdown;
720747

721748
private WorkerOptions(
722749
double maxWorkerActivitiesPerSecond,
@@ -746,7 +773,8 @@ private WorkerOptions(
746773
WorkerDeploymentOptions deploymentOptions,
747774
PollerBehavior workflowTaskPollersBehavior,
748775
PollerBehavior activityTaskPollersBehavior,
749-
PollerBehavior nexusTaskPollersBehavior) {
776+
PollerBehavior nexusTaskPollersBehavior,
777+
boolean allowActivityHeartbeatDuringShutdown) {
750778
this.maxWorkerActivitiesPerSecond = maxWorkerActivitiesPerSecond;
751779
this.maxConcurrentActivityExecutionSize = maxConcurrentActivityExecutionSize;
752780
this.maxConcurrentWorkflowTaskExecutionSize = maxConcurrentWorkflowTaskExecutionSize;
@@ -775,6 +803,7 @@ private WorkerOptions(
775803
this.workflowTaskPollersBehavior = workflowTaskPollersBehavior;
776804
this.activityTaskPollersBehavior = activityTaskPollersBehavior;
777805
this.nexusTaskPollersBehavior = nexusTaskPollersBehavior;
806+
this.allowActivityHeartbeatDuringShutdown = allowActivityHeartbeatDuringShutdown;
778807
}
779808

780809
public double getMaxWorkerActivitiesPerSecond() {
@@ -912,6 +941,11 @@ public PollerBehavior getNexusTaskPollersBehavior() {
912941
return nexusTaskPollersBehavior;
913942
}
914943

944+
@Experimental
945+
public boolean getAllowActivityHeartbeatDuringShutdown() {
946+
return allowActivityHeartbeatDuringShutdown;
947+
}
948+
915949
@Override
916950
public boolean equals(Object o) {
917951
if (this == o) return true;
@@ -944,7 +978,8 @@ && compare(maxTaskQueueActivitiesPerSecond, that.maxTaskQueueActivitiesPerSecond
944978
&& Objects.equals(deploymentOptions, that.deploymentOptions)
945979
&& Objects.equals(workflowTaskPollersBehavior, that.workflowTaskPollersBehavior)
946980
&& Objects.equals(activityTaskPollersBehavior, that.activityTaskPollersBehavior)
947-
&& Objects.equals(nexusTaskPollersBehavior, that.nexusTaskPollersBehavior);
981+
&& Objects.equals(nexusTaskPollersBehavior, that.nexusTaskPollersBehavior)
982+
&& allowActivityHeartbeatDuringShutdown == that.allowActivityHeartbeatDuringShutdown;
948983
}
949984

950985
@Override
@@ -977,7 +1012,8 @@ public int hashCode() {
9771012
deploymentOptions,
9781013
workflowTaskPollersBehavior,
9791014
activityTaskPollersBehavior,
980-
nexusTaskPollersBehavior);
1015+
nexusTaskPollersBehavior,
1016+
allowActivityHeartbeatDuringShutdown);
9811017
}
9821018

9831019
@Override
@@ -1040,6 +1076,8 @@ public String toString() {
10401076
+ activityTaskPollersBehavior
10411077
+ ", nexusTaskPollersBehavior="
10421078
+ nexusTaskPollersBehavior
1079+
+ ", allowActivityHeartbeatDuringShutdown="
1080+
+ allowActivityHeartbeatDuringShutdown
10431081
+ '}';
10441082
}
10451083
}

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,7 @@ public void verifyNewBuilderFromExistingWorkerOptions() {
5656
.setBuildId("build-id")
5757
.setStickyTaskQueueDrainTimeout(Duration.ofSeconds(15))
5858
.setIdentity("worker-identity")
59+
.setAllowActivityHeartbeatDuringShutdown(true)
5960
.build();
6061

6162
WorkerOptions w2 = WorkerOptions.newBuilder(w1).build();
@@ -89,6 +90,8 @@ public void verifyNewBuilderFromExistingWorkerOptions() {
8990
assertEquals(w1.getBuildId(), w2.getBuildId());
9091
assertEquals(w1.getStickyTaskQueueDrainTimeout(), w2.getStickyTaskQueueDrainTimeout());
9192
assertEquals(w1.getIdentity(), w2.getIdentity());
93+
assertEquals(
94+
w1.getAllowActivityHeartbeatDuringShutdown(), w2.getAllowActivityHeartbeatDuringShutdown());
9295
}
9396

9497
@Test

0 commit comments

Comments
 (0)