From 065dd773a14f6fb69fe9edece055de2a48a5bdbb Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Fri, 31 Jul 2026 00:13:48 +0000 Subject: [PATCH] Part 1: Log systemName in DataflowWorkUnitClient, Commit, and core worker states --- .../dataflow/worker/DataflowWorkUnitClient.java | 6 +++--- .../worker/StreamingModeExecutionContext.java | 6 +++++- .../beam/runners/dataflow/worker/WindmillSink.java | 14 ++++++++++---- .../worker/streaming/ComputationState.java | 4 ++++ .../streaming/harness/MetricsDataProvider.java | 2 +- .../worker/windmill/client/commits/Commit.java | 8 ++++++-- .../worker/DataflowWorkUnitClientTest.java | 4 ++-- 7 files changed, 31 insertions(+), 13 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 af8e7dd50c95..810e9d20ed77 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,13 +135,13 @@ public Optional getWorkItem() throws IOException { 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; 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 6894ac20ef97..d577b8614078 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 final long getBacklogBytes() { return backlogBytes; } + public String getSystemName() { + return systemName; + } + public long getMaxOutputKeyBytes() { return operationalLimits.getMaxOutputKeyBytes(); } @@ -585,7 +589,7 @@ public void invalidateCache() { } 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 abe5f96bb7f4..9d8a0f0da309 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 @@ public long add(WindowedValue data) throws IOException { } 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/streaming/ComputationState.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationState.java index 8020eda1b25d..5e850d4312ea 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 String getComputationId() { 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 901e2d235f85..0580b7a0b05b 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 @@ -59,7 +59,7 @@ public void appendSummaryHtml(PrintWriter writer) { writer.println("Active Keys:
"); for (ComputationState computationState : allComputationStates.get()) { - writer.print(computationState.getComputationId()); + writer.print(computationState.getSystemName()); writer.print(":
"); computationState.printActiveWork(writer); writer.println("
"); 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 bbd6cfc9432b..78e74896ef18 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 final String computationId() { return computationState().getComputationId(); } + public final String systemName() { + return computationState().getSystemName(); + } + public @Nullable WorkItemCommitRequest singleKeyRequest() { return singleKeyRequest; }; @@ -92,8 +96,8 @@ public final int getSerializedByteSize() { @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/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 85d79e6be3c1..e5f606e061f8 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 void testCloudServiceCallMapTaskStagePropagation() throws Exception { // 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); @@ -141,7 +141,7 @@ public void testCloudServiceCallSeqMapTaskStagePropagation() throws Exception { // 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);