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
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,9 @@ to docs, or any other relevant information.
Child workflow overrides and one-time routing require Temporal Server 1.32.0 or later.
- `WorkerFactoryOptions.Builder.setLoggerTagPrefix` that can be used to customized structured logging tags (MDC keys)
set by Temporal SDK in worker context.
- Workflow tasks taking longer than 5 seconds now log a `[TMPRL1104]` warning reporting the task duration and the
external storage downloads and uploads that contributed to it. The threshold is configurable with the
`TEMPORAL_WORKFLOW_TASK_DURATION_WARN_SECONDS` environment variable.

### Changed
- Release notes for all future releases are now in a single CHANGELOG.md file. `releases` directory with old release
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,14 @@ CompletableFuture<List<Payload>> store(
List<Payload> payloads,
@Nullable StorageDriverTargetInfo target,
CancellationToken<CancellationException> cancellationToken) {
return store(payloads, target, cancellationToken, null);
}

CompletableFuture<List<Payload>> store(
List<Payload> payloads,
@Nullable StorageDriverTargetInfo target,
CancellationToken<CancellationException> cancellationToken,
@Nullable StorageOperationMetrics metrics) {
StorageDriverSelectContext selectContext =
new StorageDriverSelectContextImpl(target, cancellationToken);
Map<String, Batch<Payload>> batches;
Expand All @@ -64,7 +72,7 @@ CompletableFuture<List<Payload>> store(
if (batches.isEmpty()) {
return CompletableFuture.completedFuture(payloads);
}
return runStoreDrivers(batches, target, cancellationToken)
return runStoreDrivers(batches, target, cancellationToken, metrics)
.thenApply(referencePayloads -> applyPayloadReplacements(payloads, referencePayloads));
}

Expand Down Expand Up @@ -94,16 +102,24 @@ private Map<String, Batch<Payload>> buildStoreBatches(
private CompletableFuture<List<IndexedValue<Payload>>> runStoreDrivers(
Map<String, Batch<Payload>> batches,
@Nullable StorageDriverTargetInfo target,
CancellationToken<CancellationException> cancellationToken) {
CancellationToken<CancellationException> cancellationToken,
@Nullable StorageOperationMetrics metrics) {
return withDriverScope(
cancellationToken,
scope -> {
StorageDriverStoreContext context =
new StorageDriverStoreContextImpl(target, scope.token());
for (Batch<Payload> batch : batches.values()) {
long startNanos = System.nanoTime();
scope
.attach(batch.driver.store(context, batch.values()))
.map(claims -> createReferencePayloads(batch, claims));
.map(
claims -> {
List<IndexedValue<Payload>> references =
createReferencePayloads(batch, claims);
recordBatch(metrics, batch, serializedSizeOf(batch.values()), startNanos);
return references;
});
}
return scope.awaitAll(ListUtils::flatten);
});
Expand Down Expand Up @@ -164,6 +180,13 @@ private static List<IndexedValue<Payload>> createReferencePayloads(

CompletableFuture<List<Payload>> retrieve(
List<Payload> payloads, CancellationToken<CancellationException> cancellationToken) {
return retrieve(payloads, cancellationToken, null);
}

CompletableFuture<List<Payload>> retrieve(
List<Payload> payloads,
CancellationToken<CancellationException> cancellationToken,
@Nullable StorageOperationMetrics metrics) {
Map<String, Batch<StorageDriverClaim>> batches;
try {
batches = buildRetrieveBatches(payloads);
Expand All @@ -173,7 +196,7 @@ CompletableFuture<List<Payload>> retrieve(
if (batches.isEmpty()) {
return CompletableFuture.completedFuture(payloads);
}
return runRetrieveDrivers(batches, cancellationToken)
return runRetrieveDrivers(batches, cancellationToken, metrics)
.thenApply(retrievedPayloads -> applyPayloadReplacements(payloads, retrievedPayloads));
}

Expand All @@ -200,16 +223,26 @@ private Map<String, Batch<StorageDriverClaim>> buildRetrieveBatches(List<Payload

private CompletableFuture<List<IndexedValue<Payload>>> runRetrieveDrivers(
Map<String, Batch<StorageDriverClaim>> batches,
CancellationToken<CancellationException> cancellationToken) {
CancellationToken<CancellationException> cancellationToken,
@Nullable StorageOperationMetrics metrics) {
return withDriverScope(
cancellationToken,
scope -> {
StorageDriverRetrieveContext context =
new StorageDriverRetrieveContextImpl(scope.token());
for (Batch<StorageDriverClaim> batch : batches.values()) {
long startNanos = System.nanoTime();
scope
.attach(batch.driver.retrieve(context, batch.values()))
.map(payloads -> mapPayloadsToOriginalPositions(batch, payloads));
.map(
payloads -> {
List<IndexedValue<Payload>> mapped =
mapPayloadsToOriginalPositions(batch, payloads);
// The reference's recorded size is not part of the exchange contract, so
// measure what the driver actually returned.
recordBatch(metrics, batch, serializedSizeOf(payloads), startNanos);
return mapped;
});
}
return scope.awaitAll(ListUtils::flatten);
});
Expand Down Expand Up @@ -237,6 +270,22 @@ private static List<IndexedValue<Payload>> mapPayloadsToOriginalPositions(
return replacements;
}

private static long serializedSizeOf(List<Payload> payloads) {
long total = 0;
for (Payload payload : payloads) {
total += payload.getSerializedSize();
}
return total;
}

private static void recordBatch(
@Nullable StorageOperationMetrics metrics, Batch<?> batch, long sizeBytes, long startNanos) {
if (metrics != null) {
metrics.recordBatch(
batch.size(), sizeBytes, startNanos, System.nanoTime(), batch.driver.getName());
}
}

private static <T> CompletableFuture<T> failedFuture(Throwable t) {
CompletableFuture<T> future = new CompletableFuture<>();
future.completeExceptionally(t);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -44,26 +44,55 @@ public void store(
@Nullable StorageDriverTargetInfo target,
@Nullable MessageVisitor<StorageDriverTargetInfo> targetVisitor,
CancellationToken<CancellationException> cancellationToken) {
store(builder, target, targetVisitor, cancellationToken, null);
}

public void store(
Message.Builder builder,
@Nullable StorageDriverTargetInfo target,
@Nullable MessageVisitor<StorageDriverTargetInfo> targetVisitor,
CancellationToken<CancellationException> cancellationToken,
@Nullable StorageOperationMetrics metrics) {
getOrThrowIfCancelled(
storeAsync(builder, target, targetVisitor, cancellationToken), cancellationToken);
storeAsync(builder, target, targetVisitor, cancellationToken, metrics), cancellationToken);
}

public CompletableFuture<Void> storeAsync(
Message.Builder builder,
@Nullable StorageDriverTargetInfo target,
@Nullable MessageVisitor<StorageDriverTargetInfo> targetVisitor,
CancellationToken<CancellationException> cancellationToken) {
return PayloadVisitors.visit(builder, storeOptions(target, targetVisitor, cancellationToken));
return storeAsync(builder, target, targetVisitor, cancellationToken, null);
}

public CompletableFuture<Void> storeAsync(
Message.Builder builder,
@Nullable StorageDriverTargetInfo target,
@Nullable MessageVisitor<StorageDriverTargetInfo> targetVisitor,
CancellationToken<CancellationException> cancellationToken,
@Nullable StorageOperationMetrics metrics) {
return PayloadVisitors.visit(
builder, storeOptions(target, targetVisitor, cancellationToken, metrics));
}

public <T extends Message> T retrieve(
T message, CancellationToken<CancellationException> cancellationToken) {
return getOrThrowIfCancelled(retrieveAsync(message, cancellationToken), cancellationToken);
return retrieve(message, cancellationToken, null);
}

public <T extends Message> T retrieve(
T message,
CancellationToken<CancellationException> cancellationToken,
@Nullable StorageOperationMetrics metrics) {
return getOrThrowIfCancelled(
retrieveAsync(message, cancellationToken, metrics), cancellationToken);
}

public <T extends Message> CompletableFuture<T> retrieveAsync(
T message, CancellationToken<CancellationException> cancellationToken) {
return PayloadVisitors.visit(message, retrieveOptions(cancellationToken));
T message,
CancellationToken<CancellationException> cancellationToken,
@Nullable StorageOperationMetrics metrics) {
return PayloadVisitors.visit(message, retrieveOptions(cancellationToken, metrics));
}

/**
Expand Down Expand Up @@ -116,10 +145,11 @@ private static <T> T getOrThrowIfCancelled(
private PayloadVisitorOptions<StorageDriverTargetInfo> storeOptions(
@Nullable StorageDriverTargetInfo target,
@Nullable MessageVisitor<StorageDriverTargetInfo> targetVisitor,
CancellationToken<CancellationException> cancellationToken) {
CancellationToken<CancellationException> cancellationToken,
@Nullable StorageOperationMetrics metrics) {
return PayloadVisitorOptions.<StorageDriverTargetInfo>newBuilder(
(visitedTarget, payloads) ->
payloadTransformer.store(payloads, visitedTarget, cancellationToken))
payloadTransformer.store(payloads, visitedTarget, cancellationToken, metrics))
.setInitialContext(target)
.setMessageVisitor(targetVisitor)
.setConcurrency(payloadVisitConcurrency)
Expand All @@ -128,9 +158,11 @@ private PayloadVisitorOptions<StorageDriverTargetInfo> storeOptions(
}

private PayloadVisitorOptions<Void> retrieveOptions(
CancellationToken<CancellationException> cancellationToken) {
CancellationToken<CancellationException> cancellationToken,
@Nullable StorageOperationMetrics metrics) {
return PayloadVisitorOptions.<Void>newBuilder(
(context, payloads) -> payloadTransformer.retrieve(payloads, cancellationToken))
(context, payloads) ->
payloadTransformer.retrieve(payloads, cancellationToken, metrics))
.setConcurrency(payloadVisitConcurrency)
.setSkipSearchAttributes(true)
.build();
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,69 @@
package io.temporal.internal.payload.storage;

import java.time.Duration;
import java.util.ArrayList;
import java.util.Comparator;
import java.util.List;
import java.util.Set;
import java.util.TreeSet;

/**
* Accumulates external storage activity for a single workflow task. Batches may complete
* concurrently on driver threads, so every method is synchronized.
*/
public final class StorageOperationMetrics {
private int payloadCount;
private long totalSizeBytes;
private final Set<String> driverNames = new TreeSet<>();
private final List<long[]> spans = new ArrayList<>();

/**
* Records a completed batch. A batch reports the start and end of its span rather than a
* duration, leaving {@link #getTotalDuration()} to decide how overlapping spans combine.
*/
public synchronized void recordBatch(
int count, long sizeBytes, long startNanos, long endNanos, String driverName) {
payloadCount += count;
totalSizeBytes += sizeBytes;
driverNames.add(driverName);
spans.add(new long[] {startNanos, endNanos});
}

public synchronized int getPayloadCount() {
return payloadCount;
}

public synchronized long getTotalSizeBytes() {
return totalSizeBytes;
}

public synchronized List<String> getDriverNames() {
return new ArrayList<>(driverNames);
}

/**
* Wall-clock time storage was in flight. Batches may run concurrently, so overlapping spans are
* counted once rather than summed.
*/
public synchronized Duration getTotalDuration() {
if (spans.isEmpty()) {
return Duration.ZERO;
}
List<long[]> sorted = new ArrayList<>(spans);
sorted.sort(Comparator.comparingLong(span -> span[0]));
long total = 0;
long currentStart = sorted.get(0)[0];
long currentEnd = sorted.get(0)[1];
for (int i = 1; i < sorted.size(); i++) {
long[] span = sorted.get(i);
if (span[0] > currentEnd) {
total += currentEnd - currentStart;
currentStart = span[0];
currentEnd = span[1];
} else if (span[1] > currentEnd) {
currentEnd = span[1];
}
}
return Duration.ofNanos(total + currentEnd - currentStart);
}
}
Loading
Loading