This is an automated email from the ASF dual-hosted git repository.
scwhittle pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new 8ded79b7278 Part 1: Log systemName in DataflowWorkUnitClient, Commit,
and core worker states (#39561)
8ded79b7278 is described below
commit 8ded79b72787e17662f6bf517ff88b63f9cdbc6d
Author: Ryan Wigglesworth <[email protected]>
AuthorDate: Tue Aug 4 07:57:27 2026 +0000
Part 1: Log systemName in DataflowWorkUnitClient, Commit, and core worker
states (#39561)
- Rename DataflowWorkerLoggingMDC stageName methods to systemStageName per
reviewer feedback.
- Log both computationId and systemName in MetricsDataProvider debug output.
- Input the system name for logging instead of the ComputationId
---
.../runners/dataflow/worker/DataflowWorkUnitClient.java | 12 ++++++------
.../dataflow/worker/StreamingModeExecutionContext.java | 6 +++++-
.../apache/beam/runners/dataflow/worker/WindmillSink.java | 14 ++++++++++----
.../worker/logging/DataflowWorkerLoggingHandler.java | 5 +++--
.../dataflow/worker/logging/DataflowWorkerLoggingMDC.java | 15 ++++++++-------
.../dataflow/worker/streaming/ComputationState.java | 4 ++++
.../worker/streaming/harness/MetricsDataProvider.java | 4 +++-
.../streaming/harness/StreamingWorkerStatusReporter.java | 2 +-
.../dataflow/worker/windmill/client/commits/Commit.java | 8 ++++++--
.../windmill/work/processing/StreamingWorkScheduler.java | 6 +++---
.../dataflow/worker/DataflowWorkUnitClientTest.java | 8 ++++----
.../worker/logging/DataflowWorkerLoggingHandlerTest.java | 4 ++--
.../worker/testing/RestoreDataflowLoggingMDC.java | 8 ++++----
.../worker/testing/RestoreDataflowLoggingMDCTest.java | 10 +++++-----
14 files changed, 64 insertions(+), 42 deletions(-)
diff --git
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowWorkUnitClient.java
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowWorkUnitClient.java
index af8e7dd50c9..39d35d5ac94 100644
---
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowWorkUnitClient.java
+++
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowWorkUnitClient.java
@@ -135,18 +135,18 @@ class DataflowWorkUnitClient implements WorkUnitClient {
final String stage;
if (work.getMapTask() != null) {
- stage = work.getMapTask().getStageName();
+ stage = work.getMapTask().getSystemName();
logger.info("Starting MapTask stage {}", stage);
} else if (work.getSeqMapTask() != null) {
- stage = work.getSeqMapTask().getStageName();
+ stage = work.getSeqMapTask().getSystemName();
logger.info("Starting SeqMapTask stage {}", stage);
} else if (work.getSourceOperationTask() != null) {
- stage = work.getSourceOperationTask().getStageName();
+ stage = work.getSourceOperationTask().getSystemName();
logger.info("Starting SourceOperationTask stage {}", stage);
} else {
stage = null;
}
- DataflowWorkerLoggingMDC.setStageName(stage);
+ DataflowWorkerLoggingMDC.setSystemStageName(stage);
stageStartTime.set(DateTime.now());
DataflowWorkerLoggingMDC.setWorkId(Long.toString(work.getId()));
@@ -227,7 +227,7 @@ class DataflowWorkUnitClient implements WorkUnitClient {
// Log the stage execution time of finished stages that have a stage name.
This will not be set
// in the event this status is associated with a dummy work item.
if (firstNonNull(workItemStatus.getCompleted(), Boolean.FALSE)
- && DataflowWorkerLoggingMDC.getStageName() != null) {
+ && DataflowWorkerLoggingMDC.getSystemStageName() != null) {
DateTime startTime = stageStartTime.get();
if (startTime != null) {
// elapsed time can be negative by time correction
@@ -236,7 +236,7 @@ class DataflowWorkUnitClient implements WorkUnitClient {
// This thread should have been tagged with the stage start time
during getWorkItem(),
logger.info(
"Finished processing stage {} with {} errors in {} seconds ",
- DataflowWorkerLoggingMDC.getStageName(),
+ DataflowWorkerLoggingMDC.getSystemStageName(),
numErrors,
(double) elapsed / 1000);
}
diff --git
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java
index 6894ac20ef9..d577b861407 100644
---
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java
+++
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java
@@ -268,6 +268,10 @@ public class StreamingModeExecutionContext
return backlogBytes;
}
+ public String getSystemName() {
+ return systemName;
+ }
+
public long getMaxOutputKeyBytes() {
return operationalLimits.getMaxOutputKeyBytes();
}
@@ -585,7 +589,7 @@ public class StreamingModeExecutionContext
} catch (IOException e) {
Windmill.WorkItem workItem = getWorkItem();
long shardingKey = workItem != null ? workItem.getShardingKey() : -1L;
- LOG.warn("Failed to close reader for {}-{}", computationId,
shardingKey, e);
+ LOG.warn("Failed to close reader for {}-{}", systemName, shardingKey,
e);
}
}
activeReader = null;
diff --git
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java
index abe5f96bb7f..9d8a0f0da30 100644
---
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java
+++
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java
@@ -265,20 +265,26 @@ class WindmillSink<T> extends Sink<WindowedValue<T>> {
}
if (key.size() > context.getMaxOutputKeyBytes()) {
if (context.throwExceptionsForLargeOutput()) {
- throw new OutputTooLargeException("Key too large: " + key.size());
+ throw new OutputTooLargeException(
+ String.format(
+ "Key for fused stage %s too large: %s",
context.getSystemName(), key.size()));
} else {
LOG.error(
- "Trying to output too large key with size {}. Limit is {}. See
https://cloud.google.com/dataflow/docs/guides/common-errors#key-commit-too-large-exception.
Running with --experiments=throw_exceptions_on_large_output will instead throw
an OutputTooLargeException which may be caught in user code.",
+ "Trying to output too large key for fused stage {} with size {}.
Limit is {}. See
https://cloud.google.com/dataflow/docs/guides/common-errors#key-commit-too-large-exception.
Running with --experiments=throw_exceptions_on_large_output will instead throw
an OutputTooLargeException which may be caught in user code.",
+ context.getSystemName(),
key.size(),
context.getMaxOutputKeyBytes());
}
}
if (value.size() > context.getMaxOutputValueBytes()) {
if (context.throwExceptionsForLargeOutput()) {
- throw new OutputTooLargeException("Value too large: " +
value.size());
+ throw new OutputTooLargeException(
+ String.format(
+ "Value for fused stage %s too large: %s",
context.getSystemName(), value.size()));
} else {
LOG.error(
- "Trying to output too large value with size {}. Limit is {}. See
https://cloud.google.com/dataflow/docs/guides/common-errors#key-commit-too-large-exception.
Running with --experiments=throw_exceptions_on_large_output will instead throw
an OutputTooLargeException which may be caught in user code.",
+ "Trying to output too large value for fused stage {} with size
{}. Limit is {}. See
https://cloud.google.com/dataflow/docs/guides/common-errors#key-commit-too-large-exception.
Running with --experiments=throw_exceptions_on_large_output will instead throw
an OutputTooLargeException which may be caught in user code.",
+ context.getSystemName(),
value.size(),
context.getMaxOutputValueBytes());
}
diff --git
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingHandler.java
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingHandler.java
index 62057c22b8d..e8d674af8c9 100644
---
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingHandler.java
+++
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingHandler.java
@@ -350,7 +350,8 @@ public class DataflowWorkerLoggingHandler extends Handler {
addLogField(
payloadBuilder, "exception", formatException(record.getThrown()),
MESSAGE_MAX_LENGTH);
addLogField(payloadBuilder, "thread",
String.valueOf(record.getThreadID()), FIELD_MAX_LENGTH);
- addLogField(payloadBuilder, "stage",
DataflowWorkerLoggingMDC.getStageName(), FIELD_MAX_LENGTH);
+ addLogField(
+ payloadBuilder, "stage",
DataflowWorkerLoggingMDC.getSystemStageName(), FIELD_MAX_LENGTH);
addLogField(payloadBuilder, "worker",
DataflowWorkerLoggingMDC.getWorkerId(), FIELD_MAX_LENGTH);
addLogField(payloadBuilder, "work", DataflowWorkerLoggingMDC.getWorkId(),
FIELD_MAX_LENGTH);
addLogField(payloadBuilder, "job", DataflowWorkerLoggingMDC.getJobId(),
FIELD_MAX_LENGTH);
@@ -593,7 +594,7 @@ public class DataflowWorkerLoggingHandler extends Handler {
writeIfNotEmpty(generator, "message",
getFormatter().formatMessage(record));
writeIfNotEmpty(generator, "thread",
String.valueOf(record.getThreadID()));
writeIfNotEmpty(generator, "job", DataflowWorkerLoggingMDC.getJobId());
- writeIfNotEmpty(generator, "stage",
DataflowWorkerLoggingMDC.getStageName());
+ writeIfNotEmpty(generator, "stage",
DataflowWorkerLoggingMDC.getSystemStageName());
if (currentExecutionState != null) {
NameContext nameContext = currentExecutionState.getStepName();
diff --git
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingMDC.java
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingMDC.java
index 508ef6f4169..38518a70ca6 100644
---
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingMDC.java
+++
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingMDC.java
@@ -25,7 +25,8 @@ import javax.annotation.Nullable;
})
public class DataflowWorkerLoggingMDC {
private static final InheritableThreadLocal<String> jobId = new
InheritableThreadLocal<>();
- private static final InheritableThreadLocal<String> stageName = new
InheritableThreadLocal<>();
+ private static final InheritableThreadLocal<String> systemStageName =
+ new InheritableThreadLocal<>();
private static final InheritableThreadLocal<String> workerId = new
InheritableThreadLocal<>();
private static final InheritableThreadLocal<String> workId = new
InheritableThreadLocal<>();
private static final InheritableThreadLocal<String> sdkHarnessId = new
InheritableThreadLocal<>();
@@ -35,9 +36,9 @@ public class DataflowWorkerLoggingMDC {
jobId.set(newJobId);
}
- /** Sets the Stage Name of the current thread, which will be inherited by
child threads. */
- public static void setStageName(@Nullable String newStageName) {
- stageName.set(newStageName);
+ /** Sets the System Stage Name of the current thread, which will be
inherited by child threads. */
+ public static void setSystemStageName(@Nullable String newSystemStageName) {
+ systemStageName.set(newSystemStageName);
}
/** Sets the Worker ID of the current thread, which will be inherited by
child threads. */
@@ -60,9 +61,9 @@ public class DataflowWorkerLoggingMDC {
return jobId.get();
}
- /** Gets the Stage Name of the current thread. */
- public static String getStageName() {
- return stageName.get();
+ /** Gets the System Stage Name of the current thread. */
+ public static String getSystemStageName() {
+ return systemStageName.get();
}
/** Gets the Worker ID of the current thread. */
diff --git
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationState.java
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationState.java
index 8020eda1b25..5e850d4312e 100644
---
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationState.java
+++
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationState.java
@@ -69,6 +69,10 @@ public class ComputationState {
return computationId;
}
+ public String getSystemName() {
+ return mapTask.getSystemName();
+ }
+
public MapTask getMapTask() {
return mapTask;
}
diff --git
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/MetricsDataProvider.java
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/MetricsDataProvider.java
index 901e2d235f8..f2144fd906a 100644
---
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/MetricsDataProvider.java
+++
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/MetricsDataProvider.java
@@ -60,7 +60,9 @@ class MetricsDataProvider implements StatusDataProvider {
writer.println("Active Keys: <br>");
for (ComputationState computationState : allComputationStates.get()) {
writer.print(computationState.getComputationId());
- writer.print(":<br>");
+ writer.print(" (");
+ writer.print(computationState.getSystemName());
+ writer.print("):<br>");
computationState.printActiveWork(writer);
writer.println("<br>");
}
diff --git
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/StreamingWorkerStatusReporter.java
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/StreamingWorkerStatusReporter.java
index 374dd97a1b1..4b65b263fb9 100644
---
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/StreamingWorkerStatusReporter.java
+++
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/StreamingWorkerStatusReporter.java
@@ -232,7 +232,7 @@ public final class StreamingWorkerStatusReporter {
}
private void reportHarnessStartup() {
- DataflowWorkerLoggingMDC.setStageName("startup");
+ DataflowWorkerLoggingMDC.setSystemStageName("startup");
CounterSet restartCounter = new CounterSet();
restartCounter
.longSum(
diff --git
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/Commit.java
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/Commit.java
index bbd6cfc9432..78e74896ef1 100644
---
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/Commit.java
+++
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/Commit.java
@@ -66,6 +66,10 @@ public class Commit {
return computationState().getComputationId();
}
+ public final String systemName() {
+ return computationState().getSystemName();
+ }
+
public @Nullable WorkItemCommitRequest singleKeyRequest() {
return singleKeyRequest;
};
@@ -92,8 +96,8 @@ public class Commit {
@Override
public String toString() {
Work work = workBatch.get(0);
- return "[computationId="
- + computationId()
+ return "[systemName="
+ + systemName()
+ ", shardingKey="
+ work.getShardedKey()
+ ", workId="
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 9e8265e509a..7c65c3326c9 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
@@ -175,8 +175,8 @@ public class StreamingWorkScheduler {
.setCacheToken(workItem.getCacheToken());
}
- private static void setLoggingContextComputation(@Nullable String
computationId) {
- DataflowWorkerLoggingMDC.setStageName(computationId);
+ private static void setLoggingContextComputation(@Nullable String
systemStageName) {
+ DataflowWorkerLoggingMDC.setSystemStageName(systemStageName);
}
private static void setLoggingContextWorkId(@Nullable String
workLatencyTrackingId) {
@@ -228,7 +228,7 @@ public class StreamingWorkScheduler {
Windmill.WorkItem workItem = work.getWorkItem();
String computationId = computationState.getComputationId();
LOG.debug("Starting processing for {}:\n{}", computationId, work);
- setLoggingContextComputation(computationId);
+ setLoggingContextComputation(computationState.getSystemName());
KeyTransitionListener keyTransitionListener =
createKeyTransitionListener();
keyTransitionListener.onKeyTransition(null, work);
diff --git
a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowWorkUnitClientTest.java
b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowWorkUnitClientTest.java
index 85d79e6be3c..b4b70a3a0ab 100644
---
a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowWorkUnitClientTest.java
+++
b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowWorkUnitClientTest.java
@@ -121,7 +121,7 @@ public class DataflowWorkUnitClientTest {
// Publish and acquire a map task work item, and verify we're now
processing that stage.
final String stageName = "test_stage_name";
MapTask mapTask = new MapTask();
- mapTask.setStageName(stageName);
+ mapTask.setSystemName(stageName);
WorkItem workItem = createWorkItem(PROJECT_ID, JOB_ID);
workItem.setMapTask(mapTask);
@@ -133,7 +133,7 @@ public class DataflowWorkUnitClientTest {
WorkUnitClient client = new DataflowWorkUnitClient(pipelineOptions, LOG);
assertEquals(Optional.of(workItem), client.getWorkItem());
- assertEquals(stageName, DataflowWorkerLoggingMDC.getStageName());
+ assertEquals(stageName, DataflowWorkerLoggingMDC.getSystemStageName());
}
@Test
@@ -141,7 +141,7 @@ public class DataflowWorkUnitClientTest {
// Publish and acquire a seq map task work item, and verify we're now
processing that stage.
final String stageName = "test_stage_name";
SeqMapTask seqMapTask = new SeqMapTask();
- seqMapTask.setStageName(stageName);
+ seqMapTask.setSystemName(stageName);
WorkItem workItem = createWorkItem(PROJECT_ID, JOB_ID);
workItem.setSeqMapTask(seqMapTask);
@@ -153,7 +153,7 @@ public class DataflowWorkUnitClientTest {
WorkUnitClient client = new DataflowWorkUnitClient(pipelineOptions, LOG);
assertEquals(Optional.of(workItem), client.getWorkItem());
- assertEquals(stageName, DataflowWorkerLoggingMDC.getStageName());
+ assertEquals(stageName, DataflowWorkerLoggingMDC.getSystemStageName());
}
@Test
diff --git
a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingHandlerTest.java
b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingHandlerTest.java
index c6a8581cf50..9572f404362 100644
---
a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingHandlerTest.java
+++
b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingHandlerTest.java
@@ -227,7 +227,7 @@ public class DataflowWorkerLoggingHandlerTest {
String testWorkId = "testWorkId";
DataflowWorkerLoggingMDC.setJobId(testJobId);
- DataflowWorkerLoggingMDC.setStageName(testStage);
+ DataflowWorkerLoggingMDC.setSystemStageName(testStage);
DataflowWorkerLoggingMDC.setWorkerId(testWorkerId);
DataflowWorkerLoggingMDC.setWorkId(testWorkId);
@@ -514,7 +514,7 @@ public class DataflowWorkerLoggingHandlerTest {
String testWorkId = "testWorkId";
String testJobId = "testJobId";
- DataflowWorkerLoggingMDC.setStageName(testStage);
+ DataflowWorkerLoggingMDC.setSystemStageName(testStage);
DataflowWorkerLoggingMDC.setWorkerId(testWorkerId);
DataflowWorkerLoggingMDC.setWorkId(testWorkId);
DataflowWorkerLoggingMDC.setJobId(testJobId);
diff --git
a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/testing/RestoreDataflowLoggingMDC.java
b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/testing/RestoreDataflowLoggingMDC.java
index 0bd5ceea1de..1b1226cb7dd 100644
---
a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/testing/RestoreDataflowLoggingMDC.java
+++
b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/testing/RestoreDataflowLoggingMDC.java
@@ -23,7 +23,7 @@ import org.junit.rules.ExternalResource;
/** Saves, clears and restores the current thread-local logging parameters for
tests. */
public class RestoreDataflowLoggingMDC extends ExternalResource {
private String previousJobId;
- private String previousStageName;
+ private String previousSystemStageName;
private String previousWorkerId;
private String previousWorkId;
@@ -32,11 +32,11 @@ public class RestoreDataflowLoggingMDC extends
ExternalResource {
@Override
protected void before() throws Throwable {
previousJobId = DataflowWorkerLoggingMDC.getJobId();
- previousStageName = DataflowWorkerLoggingMDC.getStageName();
+ previousSystemStageName = DataflowWorkerLoggingMDC.getSystemStageName();
previousWorkerId = DataflowWorkerLoggingMDC.getWorkerId();
previousWorkId = DataflowWorkerLoggingMDC.getWorkId();
DataflowWorkerLoggingMDC.setJobId(null);
- DataflowWorkerLoggingMDC.setStageName(null);
+ DataflowWorkerLoggingMDC.setSystemStageName(null);
DataflowWorkerLoggingMDC.setWorkerId(null);
DataflowWorkerLoggingMDC.setWorkId(null);
}
@@ -44,7 +44,7 @@ public class RestoreDataflowLoggingMDC extends
ExternalResource {
@Override
protected void after() {
DataflowWorkerLoggingMDC.setJobId(previousJobId);
- DataflowWorkerLoggingMDC.setStageName(previousStageName);
+ DataflowWorkerLoggingMDC.setSystemStageName(previousSystemStageName);
DataflowWorkerLoggingMDC.setWorkerId(previousWorkerId);
DataflowWorkerLoggingMDC.setWorkId(previousWorkId);
}
diff --git
a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/testing/RestoreDataflowLoggingMDCTest.java
b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/testing/RestoreDataflowLoggingMDCTest.java
index 3b78e93cce3..15dfd4ede9d 100644
---
a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/testing/RestoreDataflowLoggingMDCTest.java
+++
b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/testing/RestoreDataflowLoggingMDCTest.java
@@ -42,7 +42,7 @@ public class RestoreDataflowLoggingMDCTest {
final boolean[] evaluateRan = new boolean[1];
DataflowWorkerLoggingMDC.setJobId("oldJobId");
- DataflowWorkerLoggingMDC.setStageName("oldStageName");
+ DataflowWorkerLoggingMDC.setSystemStageName("oldStageName");
DataflowWorkerLoggingMDC.setWorkerId("oldWorkerId");
DataflowWorkerLoggingMDC.setWorkId("oldWorkId");
@@ -54,19 +54,19 @@ public class RestoreDataflowLoggingMDCTest {
evaluateRan[0] = true;
// Ensure parameters are cleared before the test runs
assertNull("null JobId", DataflowWorkerLoggingMDC.getJobId());
- assertNull("null StageName",
DataflowWorkerLoggingMDC.getStageName());
+ assertNull("null StageName",
DataflowWorkerLoggingMDC.getSystemStageName());
assertNull("null WorkerId",
DataflowWorkerLoggingMDC.getWorkerId());
assertNull("null WorkId",
DataflowWorkerLoggingMDC.getWorkId());
// Simulate updating parameters for the test
DataflowWorkerLoggingMDC.setJobId("newJobId");
- DataflowWorkerLoggingMDC.setStageName("newStageName");
+ DataflowWorkerLoggingMDC.setSystemStageName("newStageName");
DataflowWorkerLoggingMDC.setWorkerId("newWorkerId");
DataflowWorkerLoggingMDC.setWorkId("newWorkId");
// Ensure that the values changed
assertEquals("newJobId", DataflowWorkerLoggingMDC.getJobId());
- assertEquals("newStageName",
DataflowWorkerLoggingMDC.getStageName());
+ assertEquals("newStageName",
DataflowWorkerLoggingMDC.getSystemStageName());
assertEquals("newWorkerId",
DataflowWorkerLoggingMDC.getWorkerId());
assertEquals("newWorkId",
DataflowWorkerLoggingMDC.getWorkId());
}
@@ -77,7 +77,7 @@ public class RestoreDataflowLoggingMDCTest {
// Validate that the statement ran and that the values were reverted
assertTrue(evaluateRan[0]);
assertEquals("oldJobId", DataflowWorkerLoggingMDC.getJobId());
- assertEquals("oldStageName", DataflowWorkerLoggingMDC.getStageName());
+ assertEquals("oldStageName",
DataflowWorkerLoggingMDC.getSystemStageName());
assertEquals("oldWorkerId", DataflowWorkerLoggingMDC.getWorkerId());
assertEquals("oldWorkId", DataflowWorkerLoggingMDC.getWorkId());
}