Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions contrib/temporal-workflowstreams/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,11 @@ Items are buffered and flushed automatically every batch interval (default 2s),
when the buffer reaches the max batch size, on `forceFlush`, on an explicit
`flush()`, or on `close()`.

Background flushes run on the client's publish executor: a single daemon thread
owned by each client by default. Applications running many clients can supply a
shared executor via `publishExecutor` (see the options table); it is never shut
down by the client.

## Subscribing

There are two subscriber APIs over one shared poll engine: a non-blocking
Expand Down Expand Up @@ -196,6 +201,7 @@ unrecoverable poll failure is rethrown from `hasNext()`.
| `maxRetryDuration` | 10m | Max time to retry a failed flush before `FlushTimeoutException`. Must be < the workflow's publisher TTL (15m) to preserve exactly-once delivery |
| `payloadConverters` | standard set | Per-item serialization. Payload conversion only — the client's codec chain runs once on the envelope, never per item |
| `pollExecutor` | 2 daemon threads, client-owned | Scheduler shared by the client's subscriptions. It runs the short update-admission and delivery steps and poll cooldowns — never held during the long poll itself. A user-supplied executor is never shut down by the client; supply a bigger pool for many subscriptions against slow workflows |
| `publishExecutor` | 1 daemon thread, client-owned | Scheduler driving the client's background flushes (periodic ticks and full-buffer/`forceFlush` triggers). A flush occupies a thread while signaling the workflow. A user-supplied executor is never shut down by the client; share one across clients instead of paying a thread per client |
| `SubscribeOptions.pollCooldown` | 100ms | Min interval between polls |

## Cross-language protocol
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -101,7 +101,8 @@ private WorkflowStreamClient(
dataConverter,
options.getBatchInterval(),
options.getMaxBatchSize(),
options.getMaxRetryDuration());
options.getMaxRetryDuration(),
options.getPublishExecutor());
}

/**
Expand Down Expand Up @@ -206,7 +207,8 @@ private ScheduledExecutorService pollExecutor() {
*
* <p>Also stops this client's live subscriptions (their done futures complete normally, without
* {@link WorkflowStreamListener#onCompleted}) and, if the client owns the default poll executor,
* shuts it down. A user-supplied poll executor is never shut down.
* shuts it down. A user-supplied poll or publish executor is never shut down — only this client's
* own tasks on it are stopped.
*/
@Override
public void close() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,18 +24,21 @@ public static WorkflowStreamClientOptions getDefaultInstance() {
private final Duration maxRetryDuration;
private final PayloadConverter[] payloadConverters;
@Nullable private final ScheduledExecutorService pollExecutor;
@Nullable private final ScheduledExecutorService publishExecutor;

private WorkflowStreamClientOptions(
Duration batchInterval,
int maxBatchSize,
Duration maxRetryDuration,
PayloadConverter[] payloadConverters,
@Nullable ScheduledExecutorService pollExecutor) {
@Nullable ScheduledExecutorService pollExecutor,
@Nullable ScheduledExecutorService publishExecutor) {
this.batchInterval = batchInterval;
this.maxBatchSize = maxBatchSize;
this.maxRetryDuration = maxRetryDuration;
this.payloadConverters = payloadConverters.clone();
this.pollExecutor = pollExecutor;
this.publishExecutor = publishExecutor;
}

public Duration getBatchInterval() {
Expand All @@ -59,12 +62,18 @@ public ScheduledExecutorService getPollExecutor() {
return pollExecutor;
}

@Nullable
public ScheduledExecutorService getPublishExecutor() {
return publishExecutor;
}

public static final class Builder {
private Duration batchInterval = WorkflowStreamConstants.DEFAULT_BATCH_INTERVAL;
private int maxBatchSize;
private Duration maxRetryDuration = WorkflowStreamConstants.DEFAULT_MAX_RETRY_DURATION;
private PayloadConverter[] payloadConverters = new PayloadConverter[0];
@Nullable private ScheduledExecutorService pollExecutor;
@Nullable private ScheduledExecutorService publishExecutor;

private Builder() {}

Expand Down Expand Up @@ -128,9 +137,29 @@ public Builder setPollExecutor(ScheduledExecutorService pollExecutor) {
return this;
}

/**
* Executor that drives the client's background publish path: the periodic flushes, and the
* flushes triggered by a full buffer or {@code forceFlush}. The caller owns its lifecycle; it
* is shared across all publishes of this client and must have at least one thread. A flush
* blocks while signaling the workflow, so it occupies an executor thread for the duration of
* each send — supply a pool sized for the number of clients that may flush concurrently.
*
* <p>Default: a single-thread daemon executor created lazily and owned by the client's
* publisher (shut down by {@link WorkflowStreamClient#close}).
*/
public Builder setPublishExecutor(ScheduledExecutorService publishExecutor) {
this.publishExecutor = publishExecutor;
return this;
}

public WorkflowStreamClientOptions build() {
return new WorkflowStreamClientOptions(
batchInterval, maxBatchSize, maxRetryDuration, payloadConverters, pollExecutor);
batchInterval,
maxBatchSize,
maxRetryDuration,
payloadConverters,
pollExecutor,
publishExecutor);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,9 @@
import java.util.UUID;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit;
import javax.annotation.Nullable;

/**
* Owns the client-side publish path: it buffers published values, batches them, and sends each
Expand All @@ -36,6 +38,8 @@ public interface SignalFunction {
private final long batchIntervalMs;
private final int maxBatchSize;
private final long maxRetryDurationMs;
// When null, the publisher creates a single-thread executor it owns and shuts down in close().
@Nullable private final ScheduledExecutorService userExecutor;

private final Object stateLock = new Object();
private List<PublishEntry> buffer = new ArrayList<>();
Expand All @@ -46,7 +50,12 @@ public interface SignalFunction {
private boolean started;
private boolean closed;
private FlushTimeoutException deferredError;
// The executor driving the flush loop once started; the owned one when no user executor was
// supplied. Guarded by stateLock.
private ScheduledExecutorService scheduler;
// The periodic flush tick, tracked so it can be cancelled without shutting down a user-supplied
// executor. Guarded by stateLock.
private ScheduledFuture<?> flushTask;

/** Serializes doFlush so concurrent callers send sequentially. */
private final Object flushLock = new Object();
Expand All @@ -57,12 +66,31 @@ public StreamPublisher(
Duration batchInterval,
int maxBatchSize,
Duration maxRetryDuration) {
this(signal, dataConverter, batchInterval, maxBatchSize, maxRetryDuration, null);
}

/**
* @param executor drives the background flush loop (the periodic ticks and the flushes triggered
* by a full buffer or {@code forceFlush}). When non-null the caller owns its lifecycle and it
* is never shut down by this publisher, so many publishers can share one executor; when null
* a single-thread executor is created lazily, owned by this publisher, and shut down by
* {@link #close}. Flushes block while signaling the workflow, so each in-flight flush
* occupies an executor thread for the duration of the send.
*/
public StreamPublisher(
SignalFunction signal,
DataConverter dataConverter,
Duration batchInterval,
int maxBatchSize,
Duration maxRetryDuration,
@Nullable ScheduledExecutorService executor) {
this.signal = signal;
this.dataConverter = dataConverter;
this.publisherId = UUID.randomUUID().toString().replace("-", "").substring(0, 16);
this.batchIntervalMs = batchInterval.toMillis();
this.maxBatchSize = maxBatchSize;
this.maxRetryDurationMs = maxRetryDuration.toMillis();
this.userExecutor = executor;
}

/**
Expand Down Expand Up @@ -97,15 +125,20 @@ private void ensureStartedLocked() {
return;
}
started = true;
scheduler =
Executors.newSingleThreadScheduledExecutor(
r -> {
Thread t = new Thread(r, "temporal-workflow-stream-publisher");
t.setDaemon(true);
return t;
});
scheduler.scheduleWithFixedDelay(
this::backgroundFlush, batchIntervalMs, batchIntervalMs, TimeUnit.MILLISECONDS);
if (userExecutor != null) {
scheduler = userExecutor;
} else {
scheduler =
Executors.newSingleThreadScheduledExecutor(
r -> {
Thread t = new Thread(r, "temporal-workflow-stream-publisher");
t.setDaemon(true);
return t;
});
}
flushTask =
scheduler.scheduleWithFixedDelay(
this::backgroundFlush, batchIntervalMs, batchIntervalMs, TimeUnit.MILLISECONDS);
}

private void backgroundFlush() {
Expand All @@ -114,10 +147,15 @@ private void backgroundFlush() {
} catch (FlushTimeoutException e) {
// The pending batch was dropped and can't be recovered. Stash the error so
// flush/close surface it and stop the loop.
ScheduledFuture<?> toCancel;
ScheduledExecutorService toStop;
synchronized (stateLock) {
deferredError = e;
toStop = scheduler;
toCancel = flushTask;
toStop = ownedSchedulerLocked();
}
if (toCancel != null) {
toCancel.cancel(false);
}
if (toStop != null) {
toStop.shutdown();
Expand Down Expand Up @@ -226,18 +264,24 @@ public void flush() {

/**
* Stops the background flush loop and drains any remaining items, surfacing a deferred {@link
* FlushTimeoutException} from a prior background failure.
* FlushTimeoutException} from a prior background failure. A user-supplied executor is never shut
* down; only the periodic flush task is cancelled, leaving the executor free for its other work.
*/
public void close() {
ScheduledFuture<?> toCancel;
ScheduledExecutorService toStop;
synchronized (stateLock) {
if (closed) {
return;
}
closed = true;
toStop = scheduler;
toCancel = flushTask;
toStop = ownedSchedulerLocked();
}

if (toCancel != null) {
toCancel.cancel(false);
}
if (toStop != null) {
toStop.shutdownNow();
try {
Expand All @@ -259,6 +303,11 @@ public void close() {
throwDeferred();
}

/** Returns the executor to shut down on stop, or null when a user executor must be left alone. */
private ScheduledExecutorService ownedSchedulerLocked() {
return userExecutor == null ? scheduler : null;
}

private void throwDeferred() {
synchronized (stateLock) {
if (deferredError != null) {
Expand Down
Loading