From 7167d2316d9783df70ecb29ea9ba5205974a631d Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Fri, 31 Jul 2026 00:14:53 +0000 Subject: [PATCH 1/4] Part 2: Support systemName in WindmillStateCache and ComputationStateCache views --- .../dataflow/worker/StreamingDataflowWorker.java | 2 +- .../worker/streaming/ActiveWorkState.java | 2 +- .../worker/streaming/ComputationStateCache.java | 12 ++++++++---- .../worker/windmill/state/WindmillStateCache.java | 15 +++++++++++++-- .../streaming/ComputationStateCacheTest.java | 5 ++++- 5 files changed, 27 insertions(+), 9 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java index 2339430464c7..8d5122b0e6c8 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java @@ -924,7 +924,7 @@ static StreamingDataflowWorker forTesting( mapTask, workExecutor, stateNameMap, - stateCache.forComputation(mapTask.getStageName()))); + stateCache.forComputation(mapTask.getStageName(), mapTask.getSystemName()))); MemoryMonitor memoryMonitor = MemoryMonitor.fromOptions(options); FailureTracker failureTracker = options.isEnableStreamingEngine() diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkState.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkState.java index de4082581293..f0150cf73eb3 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkState.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkState.java @@ -185,7 +185,7 @@ synchronized void failWorkForKey(ImmutableList failedWork executableWork.work().setFailed(); LOG.debug( "Failing work {} {}. The work will be retried and is not lost.", - computationStateCache.getComputation(), + computationStateCache.getSystemName(), failedId); } } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCache.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCache.java index 4b4acb73f4a7..e6f902a65bdc 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCache.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCache.java @@ -28,6 +28,7 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ExecutionException; +import java.util.function.BiFunction; import java.util.function.Function; import javax.annotation.concurrent.ThreadSafe; import org.apache.beam.runners.dataflow.worker.apiary.FixMultiOutputInfosOnParDoInstructions; @@ -77,7 +78,8 @@ private ComputationStateCache( public static ComputationStateCache create( ComputationConfig.Fetcher computationConfigFetcher, BoundedQueueExecutor workUnitExecutor, - Function perComputationStateCacheViewFactory, + BiFunction + perComputationStateCacheViewFactory, IdGenerator idGenerator) { Function fixMultiOutputInfosOnParDoInstructions = new FixMultiOutputInfosOnParDoInstructions(idGenerator); @@ -105,7 +107,8 @@ public ComputationState load(String computationId) { fixMultiOutputInfosOnParDoInstructions.apply(computationConfig.mapTask()), workUnitExecutor, transformUserNameToStateFamilyForComputation, - perComputationStateCacheViewFactory.apply(computationId)); + perComputationStateCacheViewFactory.apply( + computationId, computationConfig.mapTask().getSystemName())); } }), fixMultiOutputInfosOnParDoInstructions, @@ -116,7 +119,8 @@ public ComputationState load(String computationId) { public static ComputationStateCache forTesting( ComputationConfig.Fetcher computationConfigFetcher, BoundedQueueExecutor workUnitExecutor, - Function perComputationStateCacheViewFactory, + BiFunction + perComputationStateCacheViewFactory, IdGenerator idGenerator, ConcurrentMap pipelineUserNameToStateFamilyNameMap) { ComputationStateCache cache = @@ -205,7 +209,7 @@ public void closeAndInvalidateAll() { public void appendSummaryHtml(PrintWriter writer) { writer.println("

Specs

"); for (ComputationState computationState : getAllPresentComputations()) { - writer.println("

" + computationState.getComputationId() + "

"); + writer.println("

" + computationState.getSystemName() + "

"); writer.print(""); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java index 7515db000852..818b5e50898f 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java @@ -170,8 +170,12 @@ public CacheStats getCacheStats() { } /** Returns a per-computation view of the state cache. */ + public ForComputation forComputation(String computation, String systemName) { + return new ForComputation(computation, systemName); + } + public ForComputation forComputation(String computation) { - return new ForComputation(computation); + return new ForComputation(computation, computation); } /** Print summary statistics of the cache to the given {@link PrintWriter}. */ @@ -353,9 +357,11 @@ private Optional value() { public class ForComputation { private final String computation; + private final String systemName; - private ForComputation(String computation) { + private ForComputation(String computation, String systemName) { this.computation = computation; + this.systemName = systemName; } /** Returns the computation associated to this class. */ @@ -363,6 +369,11 @@ public String getComputation() { return this.computation; } + /** Returns the system name associated to this class. */ + public String getSystemName() { + return this.systemName; + } + /** Invalidate all cache entries for this computation and {@code processingKey}. */ public void invalidate(ByteString processingKey, long shardingKey) { WindmillComputationKey key = diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java index f57e20d4b5fb..6785ce47d0f6 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java @@ -86,7 +86,10 @@ private static ExecutableWork createWork(ShardedKey shardedKey, long workToken, public void setUp() { computationStateCache = ComputationStateCache.create( - configFetcher, workExecutor, ignored -> stateCache, IdGenerators.decrementingLongs()); + configFetcher, + workExecutor, + (ignored1, ignored2) -> stateCache, + IdGenerators.decrementingLongs()); } @Test From ef3786218f0f67f7ab946d45b1f5ccc1a8aac67d Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Wed, 5 Aug 2026 01:17:25 +0000 Subject: [PATCH 2/4] misc fixes --- .../worker/StreamingDataflowWorker.java | 2 +- .../client/commits/CompleteCommit.java | 5 +- .../StreamingApplianceWorkCommitter.java | 1 + .../commits/StreamingEngineWorkCommitter.java | 3 + .../windmill/state/WindmillStateCache.java | 4 - .../ComputationWorkExecutorFactory.java | 7 +- .../processing/StreamingWorkScheduler.java | 4 +- .../StreamingModeExecutionContextTest.java | 2 +- .../worker/WorkerCustomSourcesTest.java | 5 +- .../StreamingEngineWorkCommitterTest.java | 35 +++- .../state/WindmillStateCacheTest.java | 153 +++++++++++++----- .../state/WindmillStateInternalsTest.java | 8 +- .../work/refresh/ActiveWorkRefresherTest.java | 3 +- 13 files changed, 164 insertions(+), 68 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java index 8d5122b0e6c8..d57d54eb1814 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java @@ -1202,7 +1202,7 @@ private void onCompleteCommit(CompleteCommit completeCommit) { WindmillComputationKey.create( completeCommit.computationId(), completeCommit.shardedKey())); stateCache - .forComputation(completeCommit.computationId()) + .forComputation(completeCommit.computationId(), completeCommit.systemName()) .invalidate(completeCommit.shardedKey()); } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/CompleteCommit.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/CompleteCommit.java index 7e2be8308954..8fd672b8f493 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/CompleteCommit.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/CompleteCommit.java @@ -39,16 +39,19 @@ public abstract class CompleteCommit { public static CompleteCommit create( String computationId, + String systemName, ShardedKey shardedKey, WorkId workId, CommitStatus status, boolean retryableFailure) { return new AutoValue_CompleteCommit( - computationId, shardedKey, workId, status, retryableFailure); + computationId, systemName, shardedKey, workId, status, retryableFailure); } public abstract String computationId(); + public abstract String systemName(); + public abstract ShardedKey shardedKey(); public abstract WorkId workId(); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingApplianceWorkCommitter.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingApplianceWorkCommitter.java index ffb9b64595c5..ea18e91c6c25 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingApplianceWorkCommitter.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingApplianceWorkCommitter.java @@ -154,6 +154,7 @@ private void completeWork( onCommitComplete.accept( CompleteCommit.create( entry.getKey().getComputationId(), + entry.getKey().getSystemName(), ShardedKey.create(workRequest.getKey(), workRequest.getShardingKey()), WorkId.builder() .setCacheToken(workRequest.getCacheToken()) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitter.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitter.java index 8ac9b1593c54..66e918b9ceba 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitter.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitter.java @@ -159,6 +159,7 @@ private void failQueuedCommit(Commit commit) { onCommitComplete.accept( CompleteCommit.create( commit.computationId(), + commit.systemName(), w.getShardedKey(), w.id(), CommitStatus.ABORTED, @@ -241,6 +242,7 @@ private boolean tryAddToCommitBatch(Commit commit, CommitWorkStream.RequestBatch onCommitComplete.accept( CompleteCommit.create( commit.computationId(), + commit.systemName(), w.getShardedKey(), w.id(), commitStatus, @@ -258,6 +260,7 @@ private boolean tryAddToCommitBatch(Commit commit, CommitWorkStream.RequestBatch onCommitComplete.accept( CompleteCommit.create( commit.computationId(), + commit.systemName(), w.getShardedKey(), w.id(), commitStatus, diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java index 818b5e50898f..ff62d12a8fa2 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java @@ -174,10 +174,6 @@ public ForComputation forComputation(String computation, String systemName) { return new ForComputation(computation, systemName); } - public ForComputation forComputation(String computation) { - return new ForComputation(computation, computation); - } - /** Print summary statistics of the cache to the given {@link PrintWriter}. */ @Override public void appendSummaryHtml(PrintWriter response) { diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java index b51512252e37..c51cdeafbc6b 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java @@ -20,6 +20,7 @@ import static org.apache.beam.runners.dataflow.DataflowRunner.hasExperiment; import com.google.api.services.dataflow.model.MapTask; +import java.util.function.BiFunction; import java.util.function.Function; import org.apache.beam.runners.dataflow.internal.CustomSources; import org.apache.beam.runners.dataflow.options.DataflowWorkerHarnessOptions; @@ -82,7 +83,7 @@ final class ComputationWorkExecutorFactory { private final DataflowWorkerHarnessOptions options; private final DataflowMapTaskExecutorFactory mapTaskExecutorFactory; private final ReaderCache readerCache; - private final Function stateCacheFactory; + private final BiFunction stateCacheFactory; private final ReaderRegistry readerRegistry; private final SinkRegistry sinkRegistry; private final DataflowExecutionStateSampler sampler; @@ -112,7 +113,7 @@ final class ComputationWorkExecutorFactory { DataflowWorkerHarnessOptions options, DataflowMapTaskExecutorFactory mapTaskExecutorFactory, ReaderCache readerCache, - Function stateCacheFactory, + BiFunction stateCacheFactory, DataflowExecutionStateSampler sampler, StreamingCounters streamingCounters, FailureTracker failureTracker, @@ -287,7 +288,7 @@ private StreamingModeExecutionContext createExecutionContext( computationId, readerCache, computationState.getTransformUserNameToStateFamily(), - stateCacheFactory.apply(computationId), + stateCacheFactory.apply(computationId, computationState.getSystemName()), stageInfo.metricsContainerRegistry(), executionStateTracker, stageInfo.executionStateRegistry(), diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java index 7c65c3326c99..b311ff8e0812 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java @@ -27,7 +27,7 @@ import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; -import java.util.function.Function; +import java.util.function.BiFunction; import java.util.function.Supplier; import javax.annotation.concurrent.ThreadSafe; import org.apache.beam.repackaged.core.org.apache.commons.lang3.tuple.Pair; @@ -119,7 +119,7 @@ public static StreamingWorkScheduler create( DataflowMapTaskExecutorFactory mapTaskExecutorFactory, BoundedQueueExecutor workExecutor, ScheduledExecutorService commitFinalizerCleanupExecutor, - Function stateCacheFactory, + BiFunction stateCacheFactory, FailureTracker failureTracker, WorkFailureProcessor workFailureProcessor, StreamingCounters streamingCounters, diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java index c5efcea4e47c..6f10f6e3749f 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java @@ -134,7 +134,7 @@ private StreamingModeExecutionContext createExecutionContext( WindmillStateCache.builder() .setSizeMb(options.getWorkerCacheMb()) .build() - .forComputation("comp"), + .forComputation("comp", "systemName"), StreamingStepMetricsContainer.createRegistry(), new DataflowExecutionStateTracker( ExecutionStateSampler.newForTest(), diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java index 679227a11dc0..6be50aba0b6c 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java @@ -152,6 +152,7 @@ public class WorkerCustomSourcesTest { @Rule public ExpectedLogs logged = ExpectedLogs.none(WorkerCustomSources.class); private static final String COMPUTATION_ID = "computationId"; + private static final String SYSTEM_NAME = "systemName"; private DataflowPipelineOptions options; @@ -1003,7 +1004,7 @@ public void testFailedWorkItemsAbort() throws Exception { WindmillStateCache.builder() .setSizeMb(options.getWorkerCacheMb()) .build() - .forComputation(COMPUTATION_ID), + .forComputation(COMPUTATION_ID, SYSTEM_NAME), StreamingStepMetricsContainer.createRegistry(), new DataflowExecutionStateTracker( ExecutionStateSampler.newForTest(), @@ -1019,7 +1020,7 @@ public void testFailedWorkItemsAbort() throws Exception { new HotKeyLogger(), /*hotKeyLoggingEnabled=*/ false, /*stepName=*/ "stepName", - /*systemName=*/ "systemName", + /*systemName=*/ SYSTEM_NAME, StreamingCounters.create(), mock(FailureTracker.class), "sourceBytesProcessCounterName", diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitterTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitterTest.java index 4c2b8a9f44fb..559925f383dd 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitterTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitterTest.java @@ -139,10 +139,15 @@ private static ComputationState createComputationState(String computationId) { } private static CompleteCommit asCompleteCommit( - String computationId, Work work, Windmill.CommitStatus status) { + String computationId, String systemName, Work work, Windmill.CommitStatus status) { Windmill.CommitStatus finalStatus = work.isFailed() ? Windmill.CommitStatus.ABORTED : status; return CompleteCommit.create( - computationId, work.getShardedKey(), work.id(), finalStatus, /* retryableFailure= */ false); + computationId, + systemName, + work.getShardedKey(), + work.id(), + finalStatus, + /* retryableFailure= */ false); } @Before @@ -197,7 +202,10 @@ public void testCommit_sendsCommitsToStreamingEngine() { assertThat(completeCommits) .contains( asCompleteCommit( - commit.computationId(), commit.workBatch().get(0), Windmill.CommitStatus.OK)); + commit.computationId(), + commit.systemName(), + commit.workBatch().get(0), + Windmill.CommitStatus.OK)); } workCommitter.stop(); @@ -238,6 +246,7 @@ public void testCommit_handlesFailedCommits() { .contains( asCompleteCommit( commit.computationId(), + commit.systemName(), commit.workBatch().get(0), Windmill.CommitStatus.ABORTED)); assertThat(committed) @@ -246,7 +255,10 @@ public void testCommit_handlesFailedCommits() { assertThat(completeCommits) .contains( asCompleteCommit( - commit.computationId(), commit.workBatch().get(0), Windmill.CommitStatus.OK)); + commit.computationId(), + commit.systemName(), + commit.workBatch().get(0), + Windmill.CommitStatus.OK)); assertThat(committed) .containsEntry( commit.workBatch().get(0).getWorkItem().getWorkToken(), commit.singleKeyRequest()); @@ -309,6 +321,7 @@ public void testCommit_handlesCompleteCommits_commitStatusNotOK() { .contains( asCompleteCommit( commit.computationId(), + commit.systemName(), commit.workBatch().get(0), expectedCommitStatus.get(commit.workBatch().get(0).id()))); } @@ -451,7 +464,10 @@ public void testMultipleCommitSendersSingleStream() { assertThat(completeCommits) .contains( asCompleteCommit( - commit.computationId(), commit.workBatch().get(0), Windmill.CommitStatus.OK)); + commit.computationId(), + commit.systemName(), + commit.workBatch().get(0), + Windmill.CommitStatus.OK)); } workCommitter.stop(); @@ -572,18 +588,21 @@ public void testCommit_multiKeyCommitSuccess() { .containsExactly( CompleteCommit.create( "computationId", + "system", workA.getShardedKey(), workA.id(), CommitStatus.OK, /* retryableFailure= */ false), CompleteCommit.create( "computationId", + "system", workB.getShardedKey(), workB.id(), CommitStatus.OK, /* retryableFailure= */ false), CompleteCommit.create( "computationId", + "system", workC.getShardedKey(), workC.id(), CommitStatus.OK, @@ -649,18 +668,21 @@ public void testCommit_multiKeyCommitFailedWork() { .containsExactly( CompleteCommit.create( "computationId", + "system", workA.getShardedKey(), workA.id(), CommitStatus.ABORTED, /* retryableFailure= */ true), CompleteCommit.create( "computationId", + "system", workB.getShardedKey(), workB.id(), CommitStatus.ABORTED, /* retryableFailure= */ false), CompleteCommit.create( "computationId", + "system", workC.getShardedKey(), workC.id(), CommitStatus.ABORTED, @@ -732,18 +754,21 @@ public void testCommit_multiKeyCommitStatusNotOK() { .containsExactly( CompleteCommit.create( "computationId", + "system", workA.getShardedKey(), workA.id(), CommitStatus.NOT_FOUND, /* retryableFailure= */ false), CompleteCommit.create( "computationId", + "system", workB.getShardedKey(), workB.id(), CommitStatus.NOT_FOUND, /* retryableFailure= */ false), CompleteCommit.create( "computationId", + "system", workC.getShardedKey(), workC.id(), CommitStatus.NOT_FOUND, diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCacheTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCacheTest.java index 2d3d9b5ccff2..788c407da1c3 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCacheTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCacheTest.java @@ -53,6 +53,7 @@ public class WindmillStateCacheTest { @Rule public transient Timeout globalTimeout = Timeout.seconds(600); private static final String COMPUTATION = "computation"; + private static final String SYSTEM_NAME = "systemName"; private static final long SHARDING_KEY = 123; private static final WindmillComputationKey COMPUTATION_KEY = WindmillComputationKey.create(COMPUTATION, ByteString.copyFromUtf8("key"), SHARDING_KEY); @@ -181,7 +182,10 @@ public void setUp() { @Test public void conflictingUserAndSystemTags() { WindmillStateCache.ForKeyAndFamily keyCache = - cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 1L).forFamily(STATE_FAMILY); + cache + .forComputation(COMPUTATION, SYSTEM_NAME) + .forKey(COMPUTATION_KEY, 0L, 1L) + .forFamily(STATE_FAMILY); StateTag> userTag = StateTags.value("tag1", StringUtf8Coder.of()); StateTag> systemTag = StateTags.makeSystemTagInternal(userTag); assertEquals(Optional.empty(), getFromCache(keyCache, StateNamespaces.global(), userTag)); @@ -220,7 +224,10 @@ public void conflictingUserAndSystemTags() { @Test public void testBasic() throws Exception { WindmillStateCache.ForKeyAndFamily keyCache = - cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 1L).forFamily(STATE_FAMILY); + cache + .forComputation(COMPUTATION, SYSTEM_NAME) + .forKey(COMPUTATION_KEY, 0L, 1L) + .forFamily(STATE_FAMILY); assertEquals( Optional.empty(), getFromCache(keyCache, StateNamespaces.global(), new TestStateTag("tag1"))); @@ -249,7 +256,10 @@ public void testBasic() throws Exception { assertEquals(482, cache.getWeight()); keyCache = - cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 2L).forFamily(STATE_FAMILY); + cache + .forComputation(COMPUTATION, SYSTEM_NAME) + .forKey(COMPUTATION_KEY, 0L, 2L) + .forFamily(STATE_FAMILY); assertEquals( Optional.of(new TestState("g1")), getFromCache(keyCache, StateNamespaces.global(), new TestStateTag("tag1"))); @@ -276,7 +286,10 @@ public void testMaxCachedEntryBytes() throws Exception { 100); // Set limit to 100 bytes, per cache entry overhead is 136. WindmillStateCache.ForKeyAndFamily keyCache = - cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 1L).forFamily(STATE_FAMILY); + cache + .forComputation(COMPUTATION, SYSTEM_NAME) + .forKey(COMPUTATION_KEY, 0L, 1L) + .forFamily(STATE_FAMILY); TestStateTag tag1 = new TestStateTag("tag1"); TestStateTag tag2 = new TestStateTag("tag2"); @@ -286,7 +299,10 @@ public void testMaxCachedEntryBytes() throws Exception { // It should not be in global cache because it's too large. keyCache = - cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 2L).forFamily(STATE_FAMILY); + cache + .forComputation(COMPUTATION, SYSTEM_NAME) + .forKey(COMPUTATION_KEY, 0L, 2L) + .forFamily(STATE_FAMILY); assertEquals(Optional.empty(), getFromCache(keyCache, StateNamespaces.global(), tag1)); // Now set limit larger. @@ -297,7 +313,10 @@ public void testMaxCachedEntryBytes() throws Exception { // It should be in global cache. keyCache = - cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 3L).forFamily(STATE_FAMILY); + cache + .forComputation(COMPUTATION, SYSTEM_NAME) + .forKey(COMPUTATION_KEY, 0L, 3L) + .forFamily(STATE_FAMILY); assertEquals( Optional.of(new TestState("g2")), getFromCache(keyCache, StateNamespaces.global(), tag2)); @@ -307,7 +326,10 @@ public void testMaxCachedEntryBytes() throws Exception { // It should be removed from global cache. keyCache = - cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 4L).forFamily(STATE_FAMILY); + cache + .forComputation(COMPUTATION, SYSTEM_NAME) + .forKey(COMPUTATION_KEY, 0L, 4L) + .forFamily(STATE_FAMILY); assertEquals(Optional.empty(), getFromCache(keyCache, StateNamespaces.global(), tag2)); } @@ -317,7 +339,7 @@ public void testDisableHistogram() throws Exception { WindmillStateCache.builder().setSizeMb(400).setEnableHistogram(false).build(); WindmillStateCache.ForKeyAndFamily keyCache = noHistogramCache - .forComputation(COMPUTATION) + .forComputation(COMPUTATION, SYSTEM_NAME) .forKey(COMPUTATION_KEY, 0L, 1L) .forFamily(STATE_FAMILY); @@ -336,7 +358,10 @@ public void testDisableHistogram() throws Exception { @Test public void testInvalidation() throws Exception { WindmillStateCache.ForKeyAndFamily keyCache = - cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 1L).forFamily(STATE_FAMILY); + cache + .forComputation(COMPUTATION, SYSTEM_NAME) + .forKey(COMPUTATION_KEY, 0L, 1L) + .forFamily(STATE_FAMILY); assertEquals( Optional.empty(), getFromCache(keyCache, StateNamespaces.global(), new TestStateTag("tag1"))); @@ -345,14 +370,20 @@ public void testInvalidation() throws Exception { keyCache.persist(); keyCache = - cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 2L).forFamily(STATE_FAMILY); + cache + .forComputation(COMPUTATION, SYSTEM_NAME) + .forKey(COMPUTATION_KEY, 0L, 2L) + .forFamily(STATE_FAMILY); assertEquals(207, cache.getWeight()); assertEquals( Optional.of(new TestState("g1")), getFromCache(keyCache, StateNamespaces.global(), new TestStateTag("tag1"))); keyCache = - cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 1L, 3L).forFamily(STATE_FAMILY); + cache + .forComputation(COMPUTATION, SYSTEM_NAME) + .forKey(COMPUTATION_KEY, 1L, 3L) + .forFamily(STATE_FAMILY); assertEquals( Optional.empty(), getFromCache(keyCache, StateNamespaces.global(), new TestStateTag("tag1"))); @@ -363,7 +394,10 @@ public void testInvalidation() throws Exception { @Test public void testEviction() throws Exception { WindmillStateCache.ForKeyAndFamily keyCache = - cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 1L).forFamily(STATE_FAMILY); + cache + .forComputation(COMPUTATION, SYSTEM_NAME) + .forKey(COMPUTATION_KEY, 0L, 1L) + .forFamily(STATE_FAMILY); putInCache(keyCache, windowNamespace(0), new TestStateTag("tag2"), new TestState("w2"), 2); putInCache( keyCache, @@ -376,7 +410,10 @@ public void testEviction() throws Exception { // Eviction is atomic across the whole window. keyCache = - cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 2L).forFamily(STATE_FAMILY); + cache + .forComputation(COMPUTATION, SYSTEM_NAME) + .forKey(COMPUTATION_KEY, 0L, 2L) + .forFamily(STATE_FAMILY); assertEquals( Optional.empty(), getFromCache(keyCache, windowNamespace(0), new TestStateTag("tag2"))); assertEquals( @@ -389,7 +426,10 @@ public void testStaleWorkItem() throws Exception { TestStateTag tag = new TestStateTag("tag2"); WindmillStateCache.ForKeyAndFamily keyCache = - cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 2L).forFamily(STATE_FAMILY); + cache + .forComputation(COMPUTATION, SYSTEM_NAME) + .forKey(COMPUTATION_KEY, 0L, 2L) + .forFamily(STATE_FAMILY); putInCache(keyCache, windowNamespace(0), tag, new TestState("w2"), 2); // Same cache. @@ -401,26 +441,41 @@ public void testStaleWorkItem() throws Exception { // Previous work token. keyCache = - cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 1L).forFamily(STATE_FAMILY); + cache + .forComputation(COMPUTATION, SYSTEM_NAME) + .forKey(COMPUTATION_KEY, 0L, 1L) + .forFamily(STATE_FAMILY); assertEquals(Optional.empty(), getFromCache(keyCache, windowNamespace(0), tag)); // Retry of work token that inserted. keyCache = - cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 2L).forFamily(STATE_FAMILY); + cache + .forComputation(COMPUTATION, SYSTEM_NAME) + .forKey(COMPUTATION_KEY, 0L, 2L) + .forFamily(STATE_FAMILY); assertEquals(Optional.empty(), getFromCache(keyCache, windowNamespace(0), tag)); keyCache = - cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 10L).forFamily(STATE_FAMILY); + cache + .forComputation(COMPUTATION, SYSTEM_NAME) + .forKey(COMPUTATION_KEY, 0L, 10L) + .forFamily(STATE_FAMILY); assertEquals(Optional.empty(), getFromCache(keyCache, windowNamespace(0), tag)); putInCache(keyCache, windowNamespace(0), tag, new TestState("w3"), 2); // Ensure that second put updated work token. keyCache = - cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 5L).forFamily(STATE_FAMILY); + cache + .forComputation(COMPUTATION, SYSTEM_NAME) + .forKey(COMPUTATION_KEY, 0L, 5L) + .forFamily(STATE_FAMILY); assertEquals(Optional.empty(), getFromCache(keyCache, windowNamespace(0), tag)); keyCache = - cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 15L).forFamily(STATE_FAMILY); + cache + .forComputation(COMPUTATION, SYSTEM_NAME) + .forKey(COMPUTATION_KEY, 0L, 15L) + .forFamily(STATE_FAMILY); assertEquals(Optional.empty(), getFromCache(keyCache, windowNamespace(0), tag)); } @@ -431,17 +486,17 @@ public void testMultipleKeys() throws Exception { WindmillStateCache.ForKeyAndFamily keyCache1 = cache - .forComputation("comp1") + .forComputation("comp1", "system1") .forKey(computationKey("comp1", "key1", SHARDING_KEY), 0L, 0L) .forFamily(STATE_FAMILY); WindmillStateCache.ForKeyAndFamily keyCache2 = cache - .forComputation("comp1") + .forComputation("comp1", "system1") .forKey(computationKey("comp1", "key2", SHARDING_KEY), 0L, 10L) .forFamily(STATE_FAMILY); WindmillStateCache.ForKeyAndFamily keyCache3 = cache - .forComputation("comp2") + .forComputation("comp2", "system2") .forKey(computationKey("comp2", "key1", SHARDING_KEY), 0L, 0L) .forFamily(STATE_FAMILY); @@ -452,7 +507,7 @@ public void testMultipleKeys() throws Exception { keyCache1 = cache - .forComputation("comp1") + .forComputation("comp1", "system1") .forKey(computationKey("comp1", "key1", SHARDING_KEY), 0L, 1L) .forFamily(STATE_FAMILY); assertEquals(Optional.of(state1), getFromCache(keyCache1, StateNamespaces.global(), tag)); @@ -465,7 +520,7 @@ public void testMultipleKeys() throws Exception { assertEquals(Optional.of(state2), getFromCache(keyCache2, StateNamespaces.global(), tag)); keyCache2 = cache - .forComputation("comp1") + .forComputation("comp1", "system1") .forKey(computationKey("comp1", "key2", SHARDING_KEY), 0L, 20L) .forFamily(STATE_FAMILY); assertEquals(Optional.of(state2), getFromCache(keyCache2, StateNamespaces.global(), tag)); @@ -480,17 +535,17 @@ public void testMultipleShardsOfKey() throws Exception { WindmillStateCache.ForKeyAndFamily key1CacheShard1 = cache - .forComputation(COMPUTATION) + .forComputation(COMPUTATION, SYSTEM_NAME) .forKey(computationKey(COMPUTATION, "key1", 1), 0L, 0L) .forFamily(STATE_FAMILY); WindmillStateCache.ForKeyAndFamily key1CacheShard2 = cache - .forComputation(COMPUTATION) + .forComputation(COMPUTATION, SYSTEM_NAME) .forKey(computationKey(COMPUTATION, "key1", 2), 0L, 0L) .forFamily(STATE_FAMILY); WindmillStateCache.ForKeyAndFamily key2CacheShard1 = cache - .forComputation(COMPUTATION) + .forComputation(COMPUTATION, SYSTEM_NAME) .forKey(computationKey(COMPUTATION, "key2", 1), 0L, 0L) .forFamily(STATE_FAMILY); @@ -500,7 +555,7 @@ public void testMultipleShardsOfKey() throws Exception { assertEquals(Optional.of(state1), getFromCache(key1CacheShard1, StateNamespaces.global(), tag)); key1CacheShard1 = cache - .forComputation(COMPUTATION) + .forComputation(COMPUTATION, SYSTEM_NAME) .forKey(computationKey(COMPUTATION, "key1", 1), 0L, 1L) .forFamily(STATE_FAMILY); assertEquals(Optional.of(state1), getFromCache(key1CacheShard1, StateNamespaces.global(), tag)); @@ -513,7 +568,7 @@ public void testMultipleShardsOfKey() throws Exception { key1CacheShard2.persist(); key1CacheShard2 = cache - .forComputation(COMPUTATION) + .forComputation(COMPUTATION, SYSTEM_NAME) .forKey(computationKey(COMPUTATION, "key1", 2), 0L, 20L) .forFamily(STATE_FAMILY); assertEquals(Optional.of(state2), getFromCache(key1CacheShard2, StateNamespaces.global(), tag)); @@ -527,7 +582,9 @@ public void testMultipleFamilies() throws Exception { TestStateTag tag = new TestStateTag("tag1"); WindmillStateCache.ForKey keyCache = - cache.forComputation("comp1").forKey(computationKey("comp1", "key1", SHARDING_KEY), 0L, 0L); + cache + .forComputation("comp1", "system1") + .forKey(computationKey("comp1", "key1", SHARDING_KEY), 0L, 0L); WindmillStateCache.ForKeyAndFamily family1 = keyCache.forFamily("family1"); WindmillStateCache.ForKeyAndFamily family2 = keyCache.forFamily("family2"); @@ -542,7 +599,9 @@ public void testMultipleFamilies() throws Exception { assertEquals(Optional.of(state2), getFromCache(family2, StateNamespaces.global(), tag)); keyCache = - cache.forComputation("comp1").forKey(computationKey("comp1", "key1", SHARDING_KEY), 0L, 1L); + cache + .forComputation("comp1", "system1") + .forKey(computationKey("comp1", "key1", SHARDING_KEY), 0L, 1L); family1 = keyCache.forFamily("family1"); family2 = keyCache.forFamily("family2"); WindmillStateCache.ForKeyAndFamily family3 = keyCache.forFamily("family3"); @@ -556,22 +615,22 @@ public void testMultipleFamilies() throws Exception { public void testExplicitInvalidation() throws Exception { WindmillStateCache.ForKeyAndFamily keyCache1 = cache - .forComputation("comp1") + .forComputation("comp1", "system1") .forKey(computationKey("comp1", "key1", 1), 0L, 0L) .forFamily(STATE_FAMILY); WindmillStateCache.ForKeyAndFamily keyCache2 = cache - .forComputation("comp1") + .forComputation("comp1", "system1") .forKey(computationKey("comp1", "key2", SHARDING_KEY), 0L, 0L) .forFamily(STATE_FAMILY); WindmillStateCache.ForKeyAndFamily keyCache3 = cache - .forComputation("comp2") + .forComputation("comp2", "system2") .forKey(computationKey("comp2", "key1", SHARDING_KEY), 0L, 0L) .forFamily(STATE_FAMILY); WindmillStateCache.ForKeyAndFamily keyCache4 = cache - .forComputation("comp1") + .forComputation("comp1", "system1") .forKey(computationKey("comp1", "key1", 2), 0L, 0L) .forFamily(STATE_FAMILY); @@ -589,22 +648,22 @@ public void testExplicitInvalidation() throws Exception { keyCache4.persist(); keyCache1 = cache - .forComputation("comp1") + .forComputation("comp1", "system1") .forKey(computationKey("comp1", "key1", 1), 0L, 1L) .forFamily(STATE_FAMILY); keyCache2 = cache - .forComputation("comp1") + .forComputation("comp1", "system1") .forKey(computationKey("comp1", "key2", SHARDING_KEY), 0L, 1L) .forFamily(STATE_FAMILY); keyCache3 = cache - .forComputation("comp2") + .forComputation("comp2", "system2") .forKey(computationKey("comp2", "key1", SHARDING_KEY), 0L, 1L) .forFamily(STATE_FAMILY); keyCache4 = cache - .forComputation("comp1") + .forComputation("comp1", "system1") .forKey(computationKey("comp1", "key1", 2), 0L, 1L) .forFamily(STATE_FAMILY); assertEquals( @@ -621,10 +680,10 @@ public void testExplicitInvalidation() throws Exception { getFromCache(keyCache4, StateNamespaces.global(), new TestStateTag("tag4"))); // Invalidation of key 1 shard 1 does not affect another shard of key 1 or other keys. - cache.forComputation("comp1").invalidate(ByteString.copyFromUtf8("key1"), 1); + cache.forComputation("comp1", "system1").invalidate(ByteString.copyFromUtf8("key1"), 1); keyCache1 = cache - .forComputation("comp1") + .forComputation("comp1", "system1") .forKey(computationKey("comp1", "key1", 1), 0L, 2L) .forFamily(STATE_FAMILY); @@ -642,7 +701,7 @@ public void testExplicitInvalidation() throws Exception { getFromCache(keyCache4, StateNamespaces.global(), new TestStateTag("tag4"))); // Invalidation of an non-existing key affects nothing. - cache.forComputation("comp1").invalidate(ByteString.copyFromUtf8("key1"), 3); + cache.forComputation("comp1", "system1").invalidate(ByteString.copyFromUtf8("key1"), 3); assertEquals( Optional.of(new TestState("g2")), @@ -679,14 +738,20 @@ public int hashCode() { @Test public void testBadCoderEquality() throws Exception { WindmillStateCache.ForKeyAndFamily keyCache1 = - cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 0L).forFamily(STATE_FAMILY); + cache + .forComputation(COMPUTATION, SYSTEM_NAME) + .forKey(COMPUTATION_KEY, 0L, 0L) + .forFamily(STATE_FAMILY); StateTag tag = new TestStateTagWithBadEquality("tag1"); putInCache(keyCache1, StateNamespaces.global(), tag, new TestState("g1"), 1); keyCache1.persist(); keyCache1 = - cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 1L).forFamily(STATE_FAMILY); + cache + .forComputation(COMPUTATION, SYSTEM_NAME) + .forKey(COMPUTATION_KEY, 0L, 1L) + .forFamily(STATE_FAMILY); assertEquals( Optional.of(new TestState("g1")), getFromCache(keyCache1, StateNamespaces.global(), tag)); assertEquals( diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java index 87b746089f11..0b55a5119564 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java @@ -225,7 +225,7 @@ public void resetUnderTest() { mockReader, false, cache - .forComputation("comp") + .forComputation("comp", "systemName") .forKey( WindmillComputationKey.create( "comp", ByteString.copyFrom("dummyKey", StandardCharsets.UTF_8), 123), @@ -241,7 +241,7 @@ public void resetUnderTest() { mockReader, true, cache - .forComputation("comp") + .forComputation("comp", "systemName") .forKey( WindmillComputationKey.create( "comp", ByteString.copyFrom("dummyNewKey", StandardCharsets.UTF_8), 123), @@ -257,7 +257,7 @@ public void resetUnderTest() { mockReader, false, cacheViaMultimap - .forComputation("comp") + .forComputation("comp", "systemName") .forKey( WindmillComputationKey.create( "comp", ByteString.copyFrom("dummyNewKey", StandardCharsets.UTF_8), 123), @@ -2049,7 +2049,7 @@ false, key(NAMESPACE, tag), STATE_FAMILY, VarIntCoder.of())) // clear cache and recreate multimapState cache - .forComputation("comp") + .forComputation("comp", "systemName") .invalidate(ByteString.copyFrom("dummyKey", StandardCharsets.UTF_8), 123); resetUnderTest(); multimapState = underTest.state(NAMESPACE, addr); diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresherTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresherTest.java index caa25bf83090..f4cad86c150a 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresherTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresherTest.java @@ -71,6 +71,7 @@ private static Instant aLongTimeAgo() { } private static final String COMPUTATION_ID_PREFIX = "ComputationId-"; + private static final String SYSTEM_NAME_PREFIX = "SystemName-"; private final HeartbeatSender heartbeatSender = mock(HeartbeatSender.class); private static BoundedQueueExecutor workExecutor() { @@ -257,7 +258,7 @@ public void testInvalidateStuckCommits() throws InterruptedException { ByteString key = ByteString.EMPTY; for (int i = 0; i < 5; i++) { WindmillStateCache.ForComputation perComputationStateCache = - spy(stateCache.forComputation(COMPUTATION_ID_PREFIX + i)); + spy(stateCache.forComputation(COMPUTATION_ID_PREFIX + i, SYSTEM_NAME_PREFIX + i)); ComputationState computationState = spy(createComputationState(i, perComputationStateCache)); ExecutableWork fakeWork = createOldWork(ShardedKey.create(key, i), i, ignored -> {}); fakeWork.work().setState(Work.State.COMMITTING); From 5650829f85da5390ed8eb81d435b24c0611a9b48 Mon Sep 17 00:00:00 2001 From: rwiggles Date: Wed, 5 Aug 2026 20:40:41 +0000 Subject: [PATCH 3/4] Simplify the pr by removing a complex effort that only touched a debug log --- .../dataflow/worker/StreamingDataflowWorker.java | 4 ++-- .../dataflow/worker/streaming/ActiveWorkState.java | 2 +- .../worker/streaming/ComputationStateCache.java | 12 ++++-------- .../processing/ComputationWorkExecutorFactory.java | 9 ++++----- .../work/processing/StreamingWorkScheduler.java | 4 ++-- .../worker/streaming/ComputationStateCacheTest.java | 5 +---- 6 files changed, 14 insertions(+), 22 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java index d57d54eb1814..2339430464c7 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java @@ -924,7 +924,7 @@ static StreamingDataflowWorker forTesting( mapTask, workExecutor, stateNameMap, - stateCache.forComputation(mapTask.getStageName(), mapTask.getSystemName()))); + stateCache.forComputation(mapTask.getStageName()))); MemoryMonitor memoryMonitor = MemoryMonitor.fromOptions(options); FailureTracker failureTracker = options.isEnableStreamingEngine() @@ -1202,7 +1202,7 @@ private void onCompleteCommit(CompleteCommit completeCommit) { WindmillComputationKey.create( completeCommit.computationId(), completeCommit.shardedKey())); stateCache - .forComputation(completeCommit.computationId(), completeCommit.systemName()) + .forComputation(completeCommit.computationId()) .invalidate(completeCommit.shardedKey()); } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkState.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkState.java index f0150cf73eb3..de4082581293 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkState.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkState.java @@ -185,7 +185,7 @@ synchronized void failWorkForKey(ImmutableList failedWork executableWork.work().setFailed(); LOG.debug( "Failing work {} {}. The work will be retried and is not lost.", - computationStateCache.getSystemName(), + computationStateCache.getComputation(), failedId); } } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCache.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCache.java index e6f902a65bdc..4b4acb73f4a7 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCache.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCache.java @@ -28,7 +28,6 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ExecutionException; -import java.util.function.BiFunction; import java.util.function.Function; import javax.annotation.concurrent.ThreadSafe; import org.apache.beam.runners.dataflow.worker.apiary.FixMultiOutputInfosOnParDoInstructions; @@ -78,8 +77,7 @@ private ComputationStateCache( public static ComputationStateCache create( ComputationConfig.Fetcher computationConfigFetcher, BoundedQueueExecutor workUnitExecutor, - BiFunction - perComputationStateCacheViewFactory, + Function perComputationStateCacheViewFactory, IdGenerator idGenerator) { Function fixMultiOutputInfosOnParDoInstructions = new FixMultiOutputInfosOnParDoInstructions(idGenerator); @@ -107,8 +105,7 @@ public ComputationState load(String computationId) { fixMultiOutputInfosOnParDoInstructions.apply(computationConfig.mapTask()), workUnitExecutor, transformUserNameToStateFamilyForComputation, - perComputationStateCacheViewFactory.apply( - computationId, computationConfig.mapTask().getSystemName())); + perComputationStateCacheViewFactory.apply(computationId)); } }), fixMultiOutputInfosOnParDoInstructions, @@ -119,8 +116,7 @@ public ComputationState load(String computationId) { public static ComputationStateCache forTesting( ComputationConfig.Fetcher computationConfigFetcher, BoundedQueueExecutor workUnitExecutor, - BiFunction - perComputationStateCacheViewFactory, + Function perComputationStateCacheViewFactory, IdGenerator idGenerator, ConcurrentMap pipelineUserNameToStateFamilyNameMap) { ComputationStateCache cache = @@ -209,7 +205,7 @@ public void closeAndInvalidateAll() { public void appendSummaryHtml(PrintWriter writer) { writer.println("

Specs

"); for (ComputationState computationState : getAllPresentComputations()) { - writer.println("

" + computationState.getSystemName() + "

"); + writer.println("

" + computationState.getComputationId() + "

"); writer.print(""); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java index c51cdeafbc6b..faa267966dd6 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java @@ -1,4 +1,4 @@ -/* +h/* * Licensed to the Apache Software Foundation (ASF) under one * or more contributor license agreements. See the NOTICE file * distributed with this work for additional information @@ -20,7 +20,6 @@ import static org.apache.beam.runners.dataflow.DataflowRunner.hasExperiment; import com.google.api.services.dataflow.model.MapTask; -import java.util.function.BiFunction; import java.util.function.Function; import org.apache.beam.runners.dataflow.internal.CustomSources; import org.apache.beam.runners.dataflow.options.DataflowWorkerHarnessOptions; @@ -83,7 +82,7 @@ final class ComputationWorkExecutorFactory { private final DataflowWorkerHarnessOptions options; private final DataflowMapTaskExecutorFactory mapTaskExecutorFactory; private final ReaderCache readerCache; - private final BiFunction stateCacheFactory; + private final Function stateCacheFactory; private final ReaderRegistry readerRegistry; private final SinkRegistry sinkRegistry; private final DataflowExecutionStateSampler sampler; @@ -113,7 +112,7 @@ final class ComputationWorkExecutorFactory { DataflowWorkerHarnessOptions options, DataflowMapTaskExecutorFactory mapTaskExecutorFactory, ReaderCache readerCache, - BiFunction stateCacheFactory, + Function stateCacheFactory, DataflowExecutionStateSampler sampler, StreamingCounters streamingCounters, FailureTracker failureTracker, @@ -288,7 +287,7 @@ private StreamingModeExecutionContext createExecutionContext( computationId, readerCache, computationState.getTransformUserNameToStateFamily(), - stateCacheFactory.apply(computationId, computationState.getSystemName()), + stateCacheFactory.apply(computationId), stageInfo.metricsContainerRegistry(), executionStateTracker, stageInfo.executionStateRegistry(), diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java index b311ff8e0812..7c65c3326c99 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java @@ -27,7 +27,7 @@ import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; -import java.util.function.BiFunction; +import java.util.function.Function; import java.util.function.Supplier; import javax.annotation.concurrent.ThreadSafe; import org.apache.beam.repackaged.core.org.apache.commons.lang3.tuple.Pair; @@ -119,7 +119,7 @@ public static StreamingWorkScheduler create( DataflowMapTaskExecutorFactory mapTaskExecutorFactory, BoundedQueueExecutor workExecutor, ScheduledExecutorService commitFinalizerCleanupExecutor, - BiFunction stateCacheFactory, + Function stateCacheFactory, FailureTracker failureTracker, WorkFailureProcessor workFailureProcessor, StreamingCounters streamingCounters, diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java index 6785ce47d0f6..f57e20d4b5fb 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java @@ -86,10 +86,7 @@ private static ExecutableWork createWork(ShardedKey shardedKey, long workToken, public void setUp() { computationStateCache = ComputationStateCache.create( - configFetcher, - workExecutor, - (ignored1, ignored2) -> stateCache, - IdGenerators.decrementingLongs()); + configFetcher, workExecutor, ignored -> stateCache, IdGenerators.decrementingLongs()); } @Test From 178159939a8fbe31583652c50eab860cfdecb88a Mon Sep 17 00:00:00 2001 From: rwiggles Date: Wed, 5 Aug 2026 21:21:07 +0000 Subject: [PATCH 4/4] refactoring --- .../client/commits/CompleteCommit.java | 5 +- .../StreamingApplianceWorkCommitter.java | 1 - .../commits/StreamingEngineWorkCommitter.java | 3 - .../windmill/state/WindmillStateCache.java | 13 +- .../ComputationWorkExecutorFactory.java | 2 +- .../StreamingModeExecutionContextTest.java | 2 +- .../worker/WorkerCustomSourcesTest.java | 2 +- .../StreamingEngineWorkCommitterTest.java | 35 +--- .../state/WindmillStateCacheTest.java | 153 +++++------------- .../state/WindmillStateInternalsTest.java | 8 +- .../work/refresh/ActiveWorkRefresherTest.java | 3 +- 11 files changed, 61 insertions(+), 166 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/CompleteCommit.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/CompleteCommit.java index 8fd672b8f493..7e2be8308954 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/CompleteCommit.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/CompleteCommit.java @@ -39,19 +39,16 @@ public abstract class CompleteCommit { public static CompleteCommit create( String computationId, - String systemName, ShardedKey shardedKey, WorkId workId, CommitStatus status, boolean retryableFailure) { return new AutoValue_CompleteCommit( - computationId, systemName, shardedKey, workId, status, retryableFailure); + computationId, shardedKey, workId, status, retryableFailure); } public abstract String computationId(); - public abstract String systemName(); - public abstract ShardedKey shardedKey(); public abstract WorkId workId(); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingApplianceWorkCommitter.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingApplianceWorkCommitter.java index ea18e91c6c25..ffb9b64595c5 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingApplianceWorkCommitter.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingApplianceWorkCommitter.java @@ -154,7 +154,6 @@ private void completeWork( onCommitComplete.accept( CompleteCommit.create( entry.getKey().getComputationId(), - entry.getKey().getSystemName(), ShardedKey.create(workRequest.getKey(), workRequest.getShardingKey()), WorkId.builder() .setCacheToken(workRequest.getCacheToken()) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitter.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitter.java index 66e918b9ceba..8ac9b1593c54 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitter.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitter.java @@ -159,7 +159,6 @@ private void failQueuedCommit(Commit commit) { onCommitComplete.accept( CompleteCommit.create( commit.computationId(), - commit.systemName(), w.getShardedKey(), w.id(), CommitStatus.ABORTED, @@ -242,7 +241,6 @@ private boolean tryAddToCommitBatch(Commit commit, CommitWorkStream.RequestBatch onCommitComplete.accept( CompleteCommit.create( commit.computationId(), - commit.systemName(), w.getShardedKey(), w.id(), commitStatus, @@ -260,7 +258,6 @@ private boolean tryAddToCommitBatch(Commit commit, CommitWorkStream.RequestBatch onCommitComplete.accept( CompleteCommit.create( commit.computationId(), - commit.systemName(), w.getShardedKey(), w.id(), commitStatus, diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java index ff62d12a8fa2..7515db000852 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java @@ -170,8 +170,8 @@ public CacheStats getCacheStats() { } /** Returns a per-computation view of the state cache. */ - public ForComputation forComputation(String computation, String systemName) { - return new ForComputation(computation, systemName); + public ForComputation forComputation(String computation) { + return new ForComputation(computation); } /** Print summary statistics of the cache to the given {@link PrintWriter}. */ @@ -353,11 +353,9 @@ private Optional value() { public class ForComputation { private final String computation; - private final String systemName; - private ForComputation(String computation, String systemName) { + private ForComputation(String computation) { this.computation = computation; - this.systemName = systemName; } /** Returns the computation associated to this class. */ @@ -365,11 +363,6 @@ public String getComputation() { return this.computation; } - /** Returns the system name associated to this class. */ - public String getSystemName() { - return this.systemName; - } - /** Invalidate all cache entries for this computation and {@code processingKey}. */ public void invalidate(ByteString processingKey, long shardingKey) { WindmillComputationKey key = diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java index faa267966dd6..b51512252e37 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java @@ -1,4 +1,4 @@ -h/* +/* * Licensed to the Apache Software Foundation (ASF) under one * or more contributor license agreements. See the NOTICE file * distributed with this work for additional information diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java index 6f10f6e3749f..c5efcea4e47c 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java @@ -134,7 +134,7 @@ private StreamingModeExecutionContext createExecutionContext( WindmillStateCache.builder() .setSizeMb(options.getWorkerCacheMb()) .build() - .forComputation("comp", "systemName"), + .forComputation("comp"), StreamingStepMetricsContainer.createRegistry(), new DataflowExecutionStateTracker( ExecutionStateSampler.newForTest(), diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java index 6be50aba0b6c..28d351670021 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java @@ -1004,7 +1004,7 @@ public void testFailedWorkItemsAbort() throws Exception { WindmillStateCache.builder() .setSizeMb(options.getWorkerCacheMb()) .build() - .forComputation(COMPUTATION_ID, SYSTEM_NAME), + .forComputation(COMPUTATION_ID), StreamingStepMetricsContainer.createRegistry(), new DataflowExecutionStateTracker( ExecutionStateSampler.newForTest(), diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitterTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitterTest.java index 559925f383dd..4c2b8a9f44fb 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitterTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitterTest.java @@ -139,15 +139,10 @@ private static ComputationState createComputationState(String computationId) { } private static CompleteCommit asCompleteCommit( - String computationId, String systemName, Work work, Windmill.CommitStatus status) { + String computationId, Work work, Windmill.CommitStatus status) { Windmill.CommitStatus finalStatus = work.isFailed() ? Windmill.CommitStatus.ABORTED : status; return CompleteCommit.create( - computationId, - systemName, - work.getShardedKey(), - work.id(), - finalStatus, - /* retryableFailure= */ false); + computationId, work.getShardedKey(), work.id(), finalStatus, /* retryableFailure= */ false); } @Before @@ -202,10 +197,7 @@ public void testCommit_sendsCommitsToStreamingEngine() { assertThat(completeCommits) .contains( asCompleteCommit( - commit.computationId(), - commit.systemName(), - commit.workBatch().get(0), - Windmill.CommitStatus.OK)); + commit.computationId(), commit.workBatch().get(0), Windmill.CommitStatus.OK)); } workCommitter.stop(); @@ -246,7 +238,6 @@ public void testCommit_handlesFailedCommits() { .contains( asCompleteCommit( commit.computationId(), - commit.systemName(), commit.workBatch().get(0), Windmill.CommitStatus.ABORTED)); assertThat(committed) @@ -255,10 +246,7 @@ public void testCommit_handlesFailedCommits() { assertThat(completeCommits) .contains( asCompleteCommit( - commit.computationId(), - commit.systemName(), - commit.workBatch().get(0), - Windmill.CommitStatus.OK)); + commit.computationId(), commit.workBatch().get(0), Windmill.CommitStatus.OK)); assertThat(committed) .containsEntry( commit.workBatch().get(0).getWorkItem().getWorkToken(), commit.singleKeyRequest()); @@ -321,7 +309,6 @@ public void testCommit_handlesCompleteCommits_commitStatusNotOK() { .contains( asCompleteCommit( commit.computationId(), - commit.systemName(), commit.workBatch().get(0), expectedCommitStatus.get(commit.workBatch().get(0).id()))); } @@ -464,10 +451,7 @@ public void testMultipleCommitSendersSingleStream() { assertThat(completeCommits) .contains( asCompleteCommit( - commit.computationId(), - commit.systemName(), - commit.workBatch().get(0), - Windmill.CommitStatus.OK)); + commit.computationId(), commit.workBatch().get(0), Windmill.CommitStatus.OK)); } workCommitter.stop(); @@ -588,21 +572,18 @@ public void testCommit_multiKeyCommitSuccess() { .containsExactly( CompleteCommit.create( "computationId", - "system", workA.getShardedKey(), workA.id(), CommitStatus.OK, /* retryableFailure= */ false), CompleteCommit.create( "computationId", - "system", workB.getShardedKey(), workB.id(), CommitStatus.OK, /* retryableFailure= */ false), CompleteCommit.create( "computationId", - "system", workC.getShardedKey(), workC.id(), CommitStatus.OK, @@ -668,21 +649,18 @@ public void testCommit_multiKeyCommitFailedWork() { .containsExactly( CompleteCommit.create( "computationId", - "system", workA.getShardedKey(), workA.id(), CommitStatus.ABORTED, /* retryableFailure= */ true), CompleteCommit.create( "computationId", - "system", workB.getShardedKey(), workB.id(), CommitStatus.ABORTED, /* retryableFailure= */ false), CompleteCommit.create( "computationId", - "system", workC.getShardedKey(), workC.id(), CommitStatus.ABORTED, @@ -754,21 +732,18 @@ public void testCommit_multiKeyCommitStatusNotOK() { .containsExactly( CompleteCommit.create( "computationId", - "system", workA.getShardedKey(), workA.id(), CommitStatus.NOT_FOUND, /* retryableFailure= */ false), CompleteCommit.create( "computationId", - "system", workB.getShardedKey(), workB.id(), CommitStatus.NOT_FOUND, /* retryableFailure= */ false), CompleteCommit.create( "computationId", - "system", workC.getShardedKey(), workC.id(), CommitStatus.NOT_FOUND, diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCacheTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCacheTest.java index 788c407da1c3..2d3d9b5ccff2 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCacheTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCacheTest.java @@ -53,7 +53,6 @@ public class WindmillStateCacheTest { @Rule public transient Timeout globalTimeout = Timeout.seconds(600); private static final String COMPUTATION = "computation"; - private static final String SYSTEM_NAME = "systemName"; private static final long SHARDING_KEY = 123; private static final WindmillComputationKey COMPUTATION_KEY = WindmillComputationKey.create(COMPUTATION, ByteString.copyFromUtf8("key"), SHARDING_KEY); @@ -182,10 +181,7 @@ public void setUp() { @Test public void conflictingUserAndSystemTags() { WindmillStateCache.ForKeyAndFamily keyCache = - cache - .forComputation(COMPUTATION, SYSTEM_NAME) - .forKey(COMPUTATION_KEY, 0L, 1L) - .forFamily(STATE_FAMILY); + cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 1L).forFamily(STATE_FAMILY); StateTag> userTag = StateTags.value("tag1", StringUtf8Coder.of()); StateTag> systemTag = StateTags.makeSystemTagInternal(userTag); assertEquals(Optional.empty(), getFromCache(keyCache, StateNamespaces.global(), userTag)); @@ -224,10 +220,7 @@ public void conflictingUserAndSystemTags() { @Test public void testBasic() throws Exception { WindmillStateCache.ForKeyAndFamily keyCache = - cache - .forComputation(COMPUTATION, SYSTEM_NAME) - .forKey(COMPUTATION_KEY, 0L, 1L) - .forFamily(STATE_FAMILY); + cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 1L).forFamily(STATE_FAMILY); assertEquals( Optional.empty(), getFromCache(keyCache, StateNamespaces.global(), new TestStateTag("tag1"))); @@ -256,10 +249,7 @@ public void testBasic() throws Exception { assertEquals(482, cache.getWeight()); keyCache = - cache - .forComputation(COMPUTATION, SYSTEM_NAME) - .forKey(COMPUTATION_KEY, 0L, 2L) - .forFamily(STATE_FAMILY); + cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 2L).forFamily(STATE_FAMILY); assertEquals( Optional.of(new TestState("g1")), getFromCache(keyCache, StateNamespaces.global(), new TestStateTag("tag1"))); @@ -286,10 +276,7 @@ public void testMaxCachedEntryBytes() throws Exception { 100); // Set limit to 100 bytes, per cache entry overhead is 136. WindmillStateCache.ForKeyAndFamily keyCache = - cache - .forComputation(COMPUTATION, SYSTEM_NAME) - .forKey(COMPUTATION_KEY, 0L, 1L) - .forFamily(STATE_FAMILY); + cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 1L).forFamily(STATE_FAMILY); TestStateTag tag1 = new TestStateTag("tag1"); TestStateTag tag2 = new TestStateTag("tag2"); @@ -299,10 +286,7 @@ public void testMaxCachedEntryBytes() throws Exception { // It should not be in global cache because it's too large. keyCache = - cache - .forComputation(COMPUTATION, SYSTEM_NAME) - .forKey(COMPUTATION_KEY, 0L, 2L) - .forFamily(STATE_FAMILY); + cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 2L).forFamily(STATE_FAMILY); assertEquals(Optional.empty(), getFromCache(keyCache, StateNamespaces.global(), tag1)); // Now set limit larger. @@ -313,10 +297,7 @@ public void testMaxCachedEntryBytes() throws Exception { // It should be in global cache. keyCache = - cache - .forComputation(COMPUTATION, SYSTEM_NAME) - .forKey(COMPUTATION_KEY, 0L, 3L) - .forFamily(STATE_FAMILY); + cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 3L).forFamily(STATE_FAMILY); assertEquals( Optional.of(new TestState("g2")), getFromCache(keyCache, StateNamespaces.global(), tag2)); @@ -326,10 +307,7 @@ public void testMaxCachedEntryBytes() throws Exception { // It should be removed from global cache. keyCache = - cache - .forComputation(COMPUTATION, SYSTEM_NAME) - .forKey(COMPUTATION_KEY, 0L, 4L) - .forFamily(STATE_FAMILY); + cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 4L).forFamily(STATE_FAMILY); assertEquals(Optional.empty(), getFromCache(keyCache, StateNamespaces.global(), tag2)); } @@ -339,7 +317,7 @@ public void testDisableHistogram() throws Exception { WindmillStateCache.builder().setSizeMb(400).setEnableHistogram(false).build(); WindmillStateCache.ForKeyAndFamily keyCache = noHistogramCache - .forComputation(COMPUTATION, SYSTEM_NAME) + .forComputation(COMPUTATION) .forKey(COMPUTATION_KEY, 0L, 1L) .forFamily(STATE_FAMILY); @@ -358,10 +336,7 @@ public void testDisableHistogram() throws Exception { @Test public void testInvalidation() throws Exception { WindmillStateCache.ForKeyAndFamily keyCache = - cache - .forComputation(COMPUTATION, SYSTEM_NAME) - .forKey(COMPUTATION_KEY, 0L, 1L) - .forFamily(STATE_FAMILY); + cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 1L).forFamily(STATE_FAMILY); assertEquals( Optional.empty(), getFromCache(keyCache, StateNamespaces.global(), new TestStateTag("tag1"))); @@ -370,20 +345,14 @@ public void testInvalidation() throws Exception { keyCache.persist(); keyCache = - cache - .forComputation(COMPUTATION, SYSTEM_NAME) - .forKey(COMPUTATION_KEY, 0L, 2L) - .forFamily(STATE_FAMILY); + cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 2L).forFamily(STATE_FAMILY); assertEquals(207, cache.getWeight()); assertEquals( Optional.of(new TestState("g1")), getFromCache(keyCache, StateNamespaces.global(), new TestStateTag("tag1"))); keyCache = - cache - .forComputation(COMPUTATION, SYSTEM_NAME) - .forKey(COMPUTATION_KEY, 1L, 3L) - .forFamily(STATE_FAMILY); + cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 1L, 3L).forFamily(STATE_FAMILY); assertEquals( Optional.empty(), getFromCache(keyCache, StateNamespaces.global(), new TestStateTag("tag1"))); @@ -394,10 +363,7 @@ public void testInvalidation() throws Exception { @Test public void testEviction() throws Exception { WindmillStateCache.ForKeyAndFamily keyCache = - cache - .forComputation(COMPUTATION, SYSTEM_NAME) - .forKey(COMPUTATION_KEY, 0L, 1L) - .forFamily(STATE_FAMILY); + cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 1L).forFamily(STATE_FAMILY); putInCache(keyCache, windowNamespace(0), new TestStateTag("tag2"), new TestState("w2"), 2); putInCache( keyCache, @@ -410,10 +376,7 @@ public void testEviction() throws Exception { // Eviction is atomic across the whole window. keyCache = - cache - .forComputation(COMPUTATION, SYSTEM_NAME) - .forKey(COMPUTATION_KEY, 0L, 2L) - .forFamily(STATE_FAMILY); + cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 2L).forFamily(STATE_FAMILY); assertEquals( Optional.empty(), getFromCache(keyCache, windowNamespace(0), new TestStateTag("tag2"))); assertEquals( @@ -426,10 +389,7 @@ public void testStaleWorkItem() throws Exception { TestStateTag tag = new TestStateTag("tag2"); WindmillStateCache.ForKeyAndFamily keyCache = - cache - .forComputation(COMPUTATION, SYSTEM_NAME) - .forKey(COMPUTATION_KEY, 0L, 2L) - .forFamily(STATE_FAMILY); + cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 2L).forFamily(STATE_FAMILY); putInCache(keyCache, windowNamespace(0), tag, new TestState("w2"), 2); // Same cache. @@ -441,41 +401,26 @@ public void testStaleWorkItem() throws Exception { // Previous work token. keyCache = - cache - .forComputation(COMPUTATION, SYSTEM_NAME) - .forKey(COMPUTATION_KEY, 0L, 1L) - .forFamily(STATE_FAMILY); + cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 1L).forFamily(STATE_FAMILY); assertEquals(Optional.empty(), getFromCache(keyCache, windowNamespace(0), tag)); // Retry of work token that inserted. keyCache = - cache - .forComputation(COMPUTATION, SYSTEM_NAME) - .forKey(COMPUTATION_KEY, 0L, 2L) - .forFamily(STATE_FAMILY); + cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 2L).forFamily(STATE_FAMILY); assertEquals(Optional.empty(), getFromCache(keyCache, windowNamespace(0), tag)); keyCache = - cache - .forComputation(COMPUTATION, SYSTEM_NAME) - .forKey(COMPUTATION_KEY, 0L, 10L) - .forFamily(STATE_FAMILY); + cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 10L).forFamily(STATE_FAMILY); assertEquals(Optional.empty(), getFromCache(keyCache, windowNamespace(0), tag)); putInCache(keyCache, windowNamespace(0), tag, new TestState("w3"), 2); // Ensure that second put updated work token. keyCache = - cache - .forComputation(COMPUTATION, SYSTEM_NAME) - .forKey(COMPUTATION_KEY, 0L, 5L) - .forFamily(STATE_FAMILY); + cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 5L).forFamily(STATE_FAMILY); assertEquals(Optional.empty(), getFromCache(keyCache, windowNamespace(0), tag)); keyCache = - cache - .forComputation(COMPUTATION, SYSTEM_NAME) - .forKey(COMPUTATION_KEY, 0L, 15L) - .forFamily(STATE_FAMILY); + cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 15L).forFamily(STATE_FAMILY); assertEquals(Optional.empty(), getFromCache(keyCache, windowNamespace(0), tag)); } @@ -486,17 +431,17 @@ public void testMultipleKeys() throws Exception { WindmillStateCache.ForKeyAndFamily keyCache1 = cache - .forComputation("comp1", "system1") + .forComputation("comp1") .forKey(computationKey("comp1", "key1", SHARDING_KEY), 0L, 0L) .forFamily(STATE_FAMILY); WindmillStateCache.ForKeyAndFamily keyCache2 = cache - .forComputation("comp1", "system1") + .forComputation("comp1") .forKey(computationKey("comp1", "key2", SHARDING_KEY), 0L, 10L) .forFamily(STATE_FAMILY); WindmillStateCache.ForKeyAndFamily keyCache3 = cache - .forComputation("comp2", "system2") + .forComputation("comp2") .forKey(computationKey("comp2", "key1", SHARDING_KEY), 0L, 0L) .forFamily(STATE_FAMILY); @@ -507,7 +452,7 @@ public void testMultipleKeys() throws Exception { keyCache1 = cache - .forComputation("comp1", "system1") + .forComputation("comp1") .forKey(computationKey("comp1", "key1", SHARDING_KEY), 0L, 1L) .forFamily(STATE_FAMILY); assertEquals(Optional.of(state1), getFromCache(keyCache1, StateNamespaces.global(), tag)); @@ -520,7 +465,7 @@ public void testMultipleKeys() throws Exception { assertEquals(Optional.of(state2), getFromCache(keyCache2, StateNamespaces.global(), tag)); keyCache2 = cache - .forComputation("comp1", "system1") + .forComputation("comp1") .forKey(computationKey("comp1", "key2", SHARDING_KEY), 0L, 20L) .forFamily(STATE_FAMILY); assertEquals(Optional.of(state2), getFromCache(keyCache2, StateNamespaces.global(), tag)); @@ -535,17 +480,17 @@ public void testMultipleShardsOfKey() throws Exception { WindmillStateCache.ForKeyAndFamily key1CacheShard1 = cache - .forComputation(COMPUTATION, SYSTEM_NAME) + .forComputation(COMPUTATION) .forKey(computationKey(COMPUTATION, "key1", 1), 0L, 0L) .forFamily(STATE_FAMILY); WindmillStateCache.ForKeyAndFamily key1CacheShard2 = cache - .forComputation(COMPUTATION, SYSTEM_NAME) + .forComputation(COMPUTATION) .forKey(computationKey(COMPUTATION, "key1", 2), 0L, 0L) .forFamily(STATE_FAMILY); WindmillStateCache.ForKeyAndFamily key2CacheShard1 = cache - .forComputation(COMPUTATION, SYSTEM_NAME) + .forComputation(COMPUTATION) .forKey(computationKey(COMPUTATION, "key2", 1), 0L, 0L) .forFamily(STATE_FAMILY); @@ -555,7 +500,7 @@ public void testMultipleShardsOfKey() throws Exception { assertEquals(Optional.of(state1), getFromCache(key1CacheShard1, StateNamespaces.global(), tag)); key1CacheShard1 = cache - .forComputation(COMPUTATION, SYSTEM_NAME) + .forComputation(COMPUTATION) .forKey(computationKey(COMPUTATION, "key1", 1), 0L, 1L) .forFamily(STATE_FAMILY); assertEquals(Optional.of(state1), getFromCache(key1CacheShard1, StateNamespaces.global(), tag)); @@ -568,7 +513,7 @@ public void testMultipleShardsOfKey() throws Exception { key1CacheShard2.persist(); key1CacheShard2 = cache - .forComputation(COMPUTATION, SYSTEM_NAME) + .forComputation(COMPUTATION) .forKey(computationKey(COMPUTATION, "key1", 2), 0L, 20L) .forFamily(STATE_FAMILY); assertEquals(Optional.of(state2), getFromCache(key1CacheShard2, StateNamespaces.global(), tag)); @@ -582,9 +527,7 @@ public void testMultipleFamilies() throws Exception { TestStateTag tag = new TestStateTag("tag1"); WindmillStateCache.ForKey keyCache = - cache - .forComputation("comp1", "system1") - .forKey(computationKey("comp1", "key1", SHARDING_KEY), 0L, 0L); + cache.forComputation("comp1").forKey(computationKey("comp1", "key1", SHARDING_KEY), 0L, 0L); WindmillStateCache.ForKeyAndFamily family1 = keyCache.forFamily("family1"); WindmillStateCache.ForKeyAndFamily family2 = keyCache.forFamily("family2"); @@ -599,9 +542,7 @@ public void testMultipleFamilies() throws Exception { assertEquals(Optional.of(state2), getFromCache(family2, StateNamespaces.global(), tag)); keyCache = - cache - .forComputation("comp1", "system1") - .forKey(computationKey("comp1", "key1", SHARDING_KEY), 0L, 1L); + cache.forComputation("comp1").forKey(computationKey("comp1", "key1", SHARDING_KEY), 0L, 1L); family1 = keyCache.forFamily("family1"); family2 = keyCache.forFamily("family2"); WindmillStateCache.ForKeyAndFamily family3 = keyCache.forFamily("family3"); @@ -615,22 +556,22 @@ public void testMultipleFamilies() throws Exception { public void testExplicitInvalidation() throws Exception { WindmillStateCache.ForKeyAndFamily keyCache1 = cache - .forComputation("comp1", "system1") + .forComputation("comp1") .forKey(computationKey("comp1", "key1", 1), 0L, 0L) .forFamily(STATE_FAMILY); WindmillStateCache.ForKeyAndFamily keyCache2 = cache - .forComputation("comp1", "system1") + .forComputation("comp1") .forKey(computationKey("comp1", "key2", SHARDING_KEY), 0L, 0L) .forFamily(STATE_FAMILY); WindmillStateCache.ForKeyAndFamily keyCache3 = cache - .forComputation("comp2", "system2") + .forComputation("comp2") .forKey(computationKey("comp2", "key1", SHARDING_KEY), 0L, 0L) .forFamily(STATE_FAMILY); WindmillStateCache.ForKeyAndFamily keyCache4 = cache - .forComputation("comp1", "system1") + .forComputation("comp1") .forKey(computationKey("comp1", "key1", 2), 0L, 0L) .forFamily(STATE_FAMILY); @@ -648,22 +589,22 @@ public void testExplicitInvalidation() throws Exception { keyCache4.persist(); keyCache1 = cache - .forComputation("comp1", "system1") + .forComputation("comp1") .forKey(computationKey("comp1", "key1", 1), 0L, 1L) .forFamily(STATE_FAMILY); keyCache2 = cache - .forComputation("comp1", "system1") + .forComputation("comp1") .forKey(computationKey("comp1", "key2", SHARDING_KEY), 0L, 1L) .forFamily(STATE_FAMILY); keyCache3 = cache - .forComputation("comp2", "system2") + .forComputation("comp2") .forKey(computationKey("comp2", "key1", SHARDING_KEY), 0L, 1L) .forFamily(STATE_FAMILY); keyCache4 = cache - .forComputation("comp1", "system1") + .forComputation("comp1") .forKey(computationKey("comp1", "key1", 2), 0L, 1L) .forFamily(STATE_FAMILY); assertEquals( @@ -680,10 +621,10 @@ public void testExplicitInvalidation() throws Exception { getFromCache(keyCache4, StateNamespaces.global(), new TestStateTag("tag4"))); // Invalidation of key 1 shard 1 does not affect another shard of key 1 or other keys. - cache.forComputation("comp1", "system1").invalidate(ByteString.copyFromUtf8("key1"), 1); + cache.forComputation("comp1").invalidate(ByteString.copyFromUtf8("key1"), 1); keyCache1 = cache - .forComputation("comp1", "system1") + .forComputation("comp1") .forKey(computationKey("comp1", "key1", 1), 0L, 2L) .forFamily(STATE_FAMILY); @@ -701,7 +642,7 @@ public void testExplicitInvalidation() throws Exception { getFromCache(keyCache4, StateNamespaces.global(), new TestStateTag("tag4"))); // Invalidation of an non-existing key affects nothing. - cache.forComputation("comp1", "system1").invalidate(ByteString.copyFromUtf8("key1"), 3); + cache.forComputation("comp1").invalidate(ByteString.copyFromUtf8("key1"), 3); assertEquals( Optional.of(new TestState("g2")), @@ -738,20 +679,14 @@ public int hashCode() { @Test public void testBadCoderEquality() throws Exception { WindmillStateCache.ForKeyAndFamily keyCache1 = - cache - .forComputation(COMPUTATION, SYSTEM_NAME) - .forKey(COMPUTATION_KEY, 0L, 0L) - .forFamily(STATE_FAMILY); + cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 0L).forFamily(STATE_FAMILY); StateTag tag = new TestStateTagWithBadEquality("tag1"); putInCache(keyCache1, StateNamespaces.global(), tag, new TestState("g1"), 1); keyCache1.persist(); keyCache1 = - cache - .forComputation(COMPUTATION, SYSTEM_NAME) - .forKey(COMPUTATION_KEY, 0L, 1L) - .forFamily(STATE_FAMILY); + cache.forComputation(COMPUTATION).forKey(COMPUTATION_KEY, 0L, 1L).forFamily(STATE_FAMILY); assertEquals( Optional.of(new TestState("g1")), getFromCache(keyCache1, StateNamespaces.global(), tag)); assertEquals( diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java index 0b55a5119564..87b746089f11 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java @@ -225,7 +225,7 @@ public void resetUnderTest() { mockReader, false, cache - .forComputation("comp", "systemName") + .forComputation("comp") .forKey( WindmillComputationKey.create( "comp", ByteString.copyFrom("dummyKey", StandardCharsets.UTF_8), 123), @@ -241,7 +241,7 @@ public void resetUnderTest() { mockReader, true, cache - .forComputation("comp", "systemName") + .forComputation("comp") .forKey( WindmillComputationKey.create( "comp", ByteString.copyFrom("dummyNewKey", StandardCharsets.UTF_8), 123), @@ -257,7 +257,7 @@ public void resetUnderTest() { mockReader, false, cacheViaMultimap - .forComputation("comp", "systemName") + .forComputation("comp") .forKey( WindmillComputationKey.create( "comp", ByteString.copyFrom("dummyNewKey", StandardCharsets.UTF_8), 123), @@ -2049,7 +2049,7 @@ false, key(NAMESPACE, tag), STATE_FAMILY, VarIntCoder.of())) // clear cache and recreate multimapState cache - .forComputation("comp", "systemName") + .forComputation("comp") .invalidate(ByteString.copyFrom("dummyKey", StandardCharsets.UTF_8), 123); resetUnderTest(); multimapState = underTest.state(NAMESPACE, addr); diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresherTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresherTest.java index f4cad86c150a..caa25bf83090 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresherTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresherTest.java @@ -71,7 +71,6 @@ private static Instant aLongTimeAgo() { } private static final String COMPUTATION_ID_PREFIX = "ComputationId-"; - private static final String SYSTEM_NAME_PREFIX = "SystemName-"; private final HeartbeatSender heartbeatSender = mock(HeartbeatSender.class); private static BoundedQueueExecutor workExecutor() { @@ -258,7 +257,7 @@ public void testInvalidateStuckCommits() throws InterruptedException { ByteString key = ByteString.EMPTY; for (int i = 0; i < 5; i++) { WindmillStateCache.ForComputation perComputationStateCache = - spy(stateCache.forComputation(COMPUTATION_ID_PREFIX + i, SYSTEM_NAME_PREFIX + i)); + spy(stateCache.forComputation(COMPUTATION_ID_PREFIX + i)); ComputationState computationState = spy(createComputationState(i, perComputationStateCache)); ExecutableWork fakeWork = createOldWork(ShardedKey.create(key, i), i, ignored -> {}); fakeWork.work().setState(Work.State.COMMITTING);