Skip to content

Commit 02b9a07

Browse files
committed
Switch logging to fused stage name in more places. This simplifies operations by removing backend details
1 parent 9cf599c commit 02b9a07

12 files changed

Lines changed: 49 additions & 32 deletions

File tree

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

Lines changed: 12 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -922,7 +922,7 @@ static StreamingDataflowWorker forTesting(
922922
mapTask,
923923
workExecutor,
924924
stateNameMap,
925-
stateCache.forComputation(mapTask.getStageName())));
925+
stateCache.forComputation(mapTask.getStageName(), mapTask.getSystemName())));
926926
MemoryMonitor memoryMonitor = MemoryMonitor.fromOptions(options);
927927
FailureTracker failureTracker =
928928
options.isEnableStreamingEngine()
@@ -1194,21 +1194,23 @@ void stop() {
11941194
}
11951195

11961196
private void onCompleteCommit(CompleteCommit completeCommit) {
1197+
Optional<ComputationState> computationState =
1198+
computationStateCache.getIfPresent(completeCommit.computationId());
11971199
if (completeCommit.status() != Windmill.CommitStatus.OK) {
11981200
readerCache.invalidateReader(
11991201
WindmillComputationKey.create(
12001202
completeCommit.computationId(), completeCommit.shardedKey()));
1201-
stateCache
1202-
.forComputation(completeCommit.computationId())
1203-
.invalidate(completeCommit.shardedKey());
1203+
computationState.ifPresent(
1204+
state ->
1205+
stateCache
1206+
.forComputation(completeCommit.computationId(), state.getSystemName())
1207+
.invalidate(completeCommit.shardedKey()));
12041208
}
12051209

1206-
computationStateCache
1207-
.getIfPresent(completeCommit.computationId())
1208-
.ifPresent(
1209-
state ->
1210-
state.completeWorkAndScheduleNextWorkForKey(
1211-
completeCommit.shardedKey(), completeCommit.workId()));
1210+
computationState.ifPresent(
1211+
state ->
1212+
state.completeWorkAndScheduleNextWorkForKey(
1213+
completeCommit.shardedKey(), completeCommit.workId()));
12121214
}
12131215

12141216
@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
@@ -180,7 +180,7 @@ synchronized void failWorkForKey(ImmutableList<WorkIdWithShardingKey> failedWork
180180
executableWork.work().setFailed();
181181
LOG.debug(
182182
"Failing work {} {}. The work will be retried and is not lost.",
183-
computationStateCache.getComputation(),
183+
computationStateCache.getSystemName(),
184184
failedId);
185185
}
186186
}

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;
@@ -81,7 +82,7 @@ final class ComputationWorkExecutorFactory {
8182
private final DataflowWorkerHarnessOptions options;
8283
private final DataflowMapTaskExecutorFactory mapTaskExecutorFactory;
8384
private final ReaderCache readerCache;
84-
private final Function<String, WindmillStateCache.ForComputation> stateCacheFactory;
85+
private final BiFunction<String, String, WindmillStateCache.ForComputation> stateCacheFactory;
8586
private final ReaderRegistry readerRegistry;
8687
private final SinkRegistry sinkRegistry;
8788
private final DataflowExecutionStateSampler sampler;
@@ -110,7 +111,7 @@ final class ComputationWorkExecutorFactory {
110111
DataflowWorkerHarnessOptions options,
111112
DataflowMapTaskExecutorFactory mapTaskExecutorFactory,
112113
ReaderCache readerCache,
113-
Function<String, WindmillStateCache.ForComputation> stateCacheFactory,
114+
BiFunction<String, String, WindmillStateCache.ForComputation> stateCacheFactory,
114115
DataflowExecutionStateSampler sampler,
115116
StreamingCounters streamingCounters,
116117
FailureTracker failureTracker,
@@ -283,7 +284,7 @@ private StreamingModeExecutionContext createExecutionContext(
283284
computationId,
284285
readerCache,
285286
computationState.getTransformUserNameToStateFamily(),
286-
stateCacheFactory.apply(computationId),
287+
stateCacheFactory.apply(computationId, stageInfo.systemName()),
287288
stageInfo.metricsContainerRegistry(),
288289
executionStateTracker,
289290
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;
@@ -113,7 +113,7 @@ public static StreamingWorkScheduler create(
113113
DataflowMapTaskExecutorFactory mapTaskExecutorFactory,
114114
BoundedQueueExecutor workExecutor,
115115
ScheduledExecutorService commitFinalizerCleanupExecutor,
116-
Function<String, WindmillStateCache.ForComputation> stateCacheFactory,
116+
BiFunction<String, String, WindmillStateCache.ForComputation> stateCacheFactory,
117117
FailureTracker failureTracker,
118118
WorkFailureProcessor workFailureProcessor,
119119
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
@@ -128,7 +128,7 @@ private StreamingModeExecutionContext createExecutionContext(
128128
WindmillStateCache.builder()
129129
.setSizeMb(options.getWorkerCacheMb())
130130
.build()
131-
.forComputation("comp"),
131+
.forComputation("comp", "systemName"),
132132
StreamingStepMetricsContainer.createRegistry(),
133133
new DataflowExecutionStateTracker(
134134
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
@@ -1003,7 +1003,7 @@ public void testFailedWorkItemsAbort() throws Exception {
10031003
WindmillStateCache.builder()
10041004
.setSizeMb(options.getWorkerCacheMb())
10051005
.build()
1006-
.forComputation(COMPUTATION_ID),
1006+
.forComputation(COMPUTATION_ID, "systemName"),
10071007
StreamingStepMetricsContainer.createRegistry(),
10081008
new DataflowExecutionStateTracker(
10091009
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
@@ -84,7 +84,10 @@ private static ExecutableWork createWork(ShardedKey shardedKey, long workToken,
8484
public void setUp() {
8585
computationStateCache =
8686
ComputationStateCache.create(
87-
configFetcher, workExecutor, ignored -> stateCache, IdGenerators.decrementingLongs());
87+
configFetcher,
88+
workExecutor,
89+
(ignored1, ignored2) -> stateCache,
90+
IdGenerators.decrementingLongs());
8891
}
8992

9093
@Test

0 commit comments

Comments
 (0)