Skip to content

Commit 1d151cf

Browse files
committed
Switch logging to fused stage name in more places. This simplifies operations by removing backend details
1 parent 5c0b302 commit 1d151cf

12 files changed

Lines changed: 49 additions & 37 deletions

File tree

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java

Lines changed: 12 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -924,7 +924,7 @@ static StreamingDataflowWorker forTesting(
924924
mapTask,
925925
workExecutor,
926926
stateNameMap,
927-
stateCache.forComputation(mapTask.getStageName())));
927+
stateCache.forComputation(mapTask.getStageName(), mapTask.getSystemName())));
928928
MemoryMonitor memoryMonitor = MemoryMonitor.fromOptions(options);
929929
FailureTracker failureTracker =
930930
options.isEnableStreamingEngine()
@@ -1197,26 +1197,23 @@ void stop() {
11971197
}
11981198

11991199
private void onCompleteCommit(CompleteCommit completeCommit) {
1200+
Optional<ComputationState> computationState =
1201+
computationStateCache.getIfPresent(completeCommit.computationId());
12001202
if (completeCommit.status() != Windmill.CommitStatus.OK) {
12011203
readerCache.invalidateReader(
12021204
WindmillComputationKey.create(
12031205
completeCommit.computationId(), completeCommit.shardedKey()));
1204-
stateCache
1205-
.forComputation(completeCommit.computationId())
1206-
.invalidate(completeCommit.shardedKey());
1206+
computationState.ifPresent(
1207+
state ->
1208+
stateCache
1209+
.forComputation(completeCommit.computationId(), state.getSystemName())
1210+
.invalidate(completeCommit.shardedKey()));
12071211
}
12081212

1209-
computationStateCache
1210-
.getIfPresent(completeCommit.computationId())
1211-
.ifPresent(
1212-
state -> {
1213-
if (completeCommit.retryableFailure()) {
1214-
state.reexecuteActiveWork(completeCommit.shardedKey(), completeCommit.workId());
1215-
} else {
1216-
state.completeWorkAndScheduleNextWorkForKey(
1217-
completeCommit.shardedKey(), completeCommit.workId());
1218-
}
1219-
});
1213+
computationState.ifPresent(
1214+
state ->
1215+
state.completeWorkAndScheduleNextWorkForKey(
1216+
completeCommit.shardedKey(), completeCommit.workId()));
12201217
}
12211218

12221219
@AutoValue

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkState.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -185,7 +185,7 @@ synchronized void failWorkForKey(ImmutableList<WorkIdWithShardingKey> failedWork
185185
executableWork.work().setFailed();
186186
LOG.debug(
187187
"Failing work {} {}. The work will be retried and is not lost.",
188-
computationStateCache.getComputation(),
188+
computationStateCache.getSystemName(),
189189
failedId);
190190
}
191191
}

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCache.java

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@
2828
import java.util.concurrent.ConcurrentHashMap;
2929
import java.util.concurrent.ConcurrentMap;
3030
import java.util.concurrent.ExecutionException;
31+
import java.util.function.BiFunction;
3132
import java.util.function.Function;
3233
import javax.annotation.concurrent.ThreadSafe;
3334
import org.apache.beam.runners.dataflow.worker.apiary.FixMultiOutputInfosOnParDoInstructions;
@@ -77,7 +78,8 @@ private ComputationStateCache(
7778
public static ComputationStateCache create(
7879
ComputationConfig.Fetcher computationConfigFetcher,
7980
BoundedQueueExecutor workUnitExecutor,
80-
Function<String, WindmillStateCache.ForComputation> perComputationStateCacheViewFactory,
81+
BiFunction<String, String, WindmillStateCache.ForComputation>
82+
perComputationStateCacheViewFactory,
8183
IdGenerator idGenerator) {
8284
Function<MapTask, MapTask> fixMultiOutputInfosOnParDoInstructions =
8385
new FixMultiOutputInfosOnParDoInstructions(idGenerator);
@@ -105,7 +107,8 @@ public ComputationState load(String computationId) {
105107
fixMultiOutputInfosOnParDoInstructions.apply(computationConfig.mapTask()),
106108
workUnitExecutor,
107109
transformUserNameToStateFamilyForComputation,
108-
perComputationStateCacheViewFactory.apply(computationId));
110+
perComputationStateCacheViewFactory.apply(
111+
computationId, computationConfig.mapTask().getSystemName()));
109112
}
110113
}),
111114
fixMultiOutputInfosOnParDoInstructions,
@@ -116,7 +119,8 @@ public ComputationState load(String computationId) {
116119
public static ComputationStateCache forTesting(
117120
ComputationConfig.Fetcher computationConfigFetcher,
118121
BoundedQueueExecutor workUnitExecutor,
119-
Function<String, WindmillStateCache.ForComputation> perComputationStateCacheViewFactory,
122+
BiFunction<String, String, WindmillStateCache.ForComputation>
123+
perComputationStateCacheViewFactory,
120124
IdGenerator idGenerator,
121125
ConcurrentMap<String, String> pipelineUserNameToStateFamilyNameMap) {
122126
ComputationStateCache cache =
@@ -205,7 +209,7 @@ public void closeAndInvalidateAll() {
205209
public void appendSummaryHtml(PrintWriter writer) {
206210
writer.println("<h1>Specs</h1>");
207211
for (ComputationState computationState : getAllPresentComputations()) {
208-
writer.println("<h3>" + computationState.getComputationId() + "</h3>");
212+
writer.println("<h3>" + computationState.getSystemName() + "</h3>");
209213
writer.print("<script>document.write(JSON.stringify(");
210214
writer.print(computationState.getMapTask().toString());
211215
writer.println(", null, \"&nbsp&nbsp\").replace(/\\n/g, \"<br>\"))</script>");

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/MetricsDataProvider.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -59,7 +59,7 @@ public void appendSummaryHtml(PrintWriter writer) {
5959

6060
writer.println("Active Keys: <br>");
6161
for (ComputationState computationState : allComputationStates.get()) {
62-
writer.print(computationState.getComputationId());
62+
writer.print(computationState.getSystemName());
6363
writer.print(":<br>");
6464
computationState.printActiveWork(writer);
6565
writer.println("<br>");

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java

Lines changed: 10 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -170,8 +170,8 @@ public CacheStats getCacheStats() {
170170
}
171171

172172
/** Returns a per-computation view of the state cache. */
173-
public ForComputation forComputation(String computation) {
174-
return new ForComputation(computation);
173+
public ForComputation forComputation(String computation, String systemName) {
174+
return new ForComputation(computation, systemName);
175175
}
176176

177177
/** Print summary statistics of the cache to the given {@link PrintWriter}. */
@@ -353,16 +353,23 @@ private Optional<T> value() {
353353
public class ForComputation {
354354

355355
private final String computation;
356+
private final String systemName;
356357

357-
private ForComputation(String computation) {
358+
private ForComputation(String computation, String systemName) {
358359
this.computation = computation;
360+
this.systemName = systemName;
359361
}
360362

361363
/** Returns the computation associated to this class. */
362364
public String getComputation() {
363365
return this.computation;
364366
}
365367

368+
/** Returns the system name associated to this class. */
369+
public String getSystemName() {
370+
return this.systemName;
371+
}
372+
366373
/** Invalidate all cache entries for this computation and {@code processingKey}. */
367374
public void invalidate(ByteString processingKey, long shardingKey) {
368375
WindmillComputationKey key =

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
import static org.apache.beam.runners.dataflow.DataflowRunner.hasExperiment;
2121

2222
import com.google.api.services.dataflow.model.MapTask;
23+
import java.util.function.BiFunction;
2324
import java.util.function.Function;
2425
import org.apache.beam.runners.dataflow.internal.CustomSources;
2526
import org.apache.beam.runners.dataflow.options.DataflowWorkerHarnessOptions;
@@ -82,7 +83,7 @@ final class ComputationWorkExecutorFactory {
8283
private final DataflowWorkerHarnessOptions options;
8384
private final DataflowMapTaskExecutorFactory mapTaskExecutorFactory;
8485
private final ReaderCache readerCache;
85-
private final Function<String, WindmillStateCache.ForComputation> stateCacheFactory;
86+
private final BiFunction<String, String, WindmillStateCache.ForComputation> stateCacheFactory;
8687
private final ReaderRegistry readerRegistry;
8788
private final SinkRegistry sinkRegistry;
8889
private final DataflowExecutionStateSampler sampler;
@@ -111,7 +112,7 @@ final class ComputationWorkExecutorFactory {
111112
DataflowWorkerHarnessOptions options,
112113
DataflowMapTaskExecutorFactory mapTaskExecutorFactory,
113114
ReaderCache readerCache,
114-
Function<String, WindmillStateCache.ForComputation> stateCacheFactory,
115+
BiFunction<String, String, WindmillStateCache.ForComputation> stateCacheFactory,
115116
DataflowExecutionStateSampler sampler,
116117
StreamingCounters streamingCounters,
117118
FailureTracker failureTracker,
@@ -286,7 +287,7 @@ private StreamingModeExecutionContext createExecutionContext(
286287
computationId,
287288
readerCache,
288289
computationState.getTransformUserNameToStateFamily(),
289-
stateCacheFactory.apply(computationId),
290+
stateCacheFactory.apply(computationId, stageInfo.systemName()),
290291
stageInfo.metricsContainerRegistry(),
291292
executionStateTracker,
292293
stageInfo.executionStateRegistry(),

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -26,7 +26,7 @@
2626
import java.util.concurrent.ConcurrentMap;
2727
import java.util.concurrent.ScheduledExecutorService;
2828
import java.util.concurrent.TimeUnit;
29-
import java.util.function.Function;
29+
import java.util.function.BiFunction;
3030
import java.util.function.Supplier;
3131
import javax.annotation.concurrent.ThreadSafe;
3232
import org.apache.beam.repackaged.core.org.apache.commons.lang3.tuple.Pair;
@@ -115,7 +115,7 @@ public static StreamingWorkScheduler create(
115115
DataflowMapTaskExecutorFactory mapTaskExecutorFactory,
116116
BoundedQueueExecutor workExecutor,
117117
ScheduledExecutorService commitFinalizerCleanupExecutor,
118-
Function<String, WindmillStateCache.ForComputation> stateCacheFactory,
118+
BiFunction<String, String, WindmillStateCache.ForComputation> stateCacheFactory,
119119
FailureTracker failureTracker,
120120
WorkFailureProcessor workFailureProcessor,
121121
StreamingCounters streamingCounters,

runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -134,7 +134,7 @@ private StreamingModeExecutionContext createExecutionContext(
134134
WindmillStateCache.builder()
135135
.setSizeMb(options.getWorkerCacheMb())
136136
.build()
137-
.forComputation("comp"),
137+
.forComputation("comp", "systemName"),
138138
StreamingStepMetricsContainer.createRegistry(),
139139
new DataflowExecutionStateTracker(
140140
ExecutionStateSampler.newForTest(),

runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1004,7 +1004,7 @@ public void testFailedWorkItemsAbort() throws Exception {
10041004
WindmillStateCache.builder()
10051005
.setSizeMb(options.getWorkerCacheMb())
10061006
.build()
1007-
.forComputation(COMPUTATION_ID),
1007+
.forComputation(COMPUTATION_ID, "systemName"),
10081008
StreamingStepMetricsContainer.createRegistry(),
10091009
new DataflowExecutionStateTracker(
10101010
ExecutionStateSampler.newForTest(),

runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -86,7 +86,10 @@ private static ExecutableWork createWork(ShardedKey shardedKey, long workToken,
8686
public void setUp() {
8787
computationStateCache =
8888
ComputationStateCache.create(
89-
configFetcher, workExecutor, ignored -> stateCache, IdGenerators.decrementingLongs());
89+
configFetcher,
90+
workExecutor,
91+
(ignored1, ignored2) -> stateCache,
92+
IdGenerators.decrementingLongs());
9093
}
9194

9295
@Test

0 commit comments

Comments
 (0)