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());
   }

Reply via email to