This is an automated email from the ASF dual-hosted git repository.

stankiewicz 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 faf253a519f [Java] Reduce logging severity for state requests 
cancelled by the runner (#40302)
faf253a519f is described below

commit faf253a519f46581fee499d0938e05dcd1a54fe0
Author: parveensania <[email protected]>
AuthorDate: Fri Oct 2 01:48:26 2026 -0700

    [Java] Reduce logging severity for state requests cancelled by the runner 
(#40302)
    
    * Reduce logging severity for state requests cancelled by runner in Java Fn 
Harness
---
 .../fn/harness/control/BeamFnControlClient.java    | 19 +++--
 .../harness/data/PCollectionConsumerRegistry.java  | 17 +++--
 .../harness/state/BeamFnStateGrpcClientCache.java  |  2 +
 .../fn/harness/state/WorkCancelledException.java   | 44 ++++++++++++
 .../harness/control/BeamFnControlClientTest.java   | 22 ++++++
 .../data/PCollectionConsumerRegistryTest.java      | 84 ++++++++++++++++++++++
 .../state/BeamFnStateGrpcClientCacheTest.java      | 21 +++++-
 7 files changed, 196 insertions(+), 13 deletions(-)

diff --git 
a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/control/BeamFnControlClient.java
 
b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/control/BeamFnControlClient.java
index 119bdd73d94..74e7d7cf1a4 100644
--- 
a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/control/BeamFnControlClient.java
+++ 
b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/control/BeamFnControlClient.java
@@ -23,6 +23,7 @@ import java.util.EnumMap;
 import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.Executor;
 import org.apache.beam.fn.harness.logging.BeamFnLoggingMDC;
+import org.apache.beam.fn.harness.state.WorkCancelledException;
 import org.apache.beam.model.fnexecution.v1.BeamFnApi;
 import org.apache.beam.model.fnexecution.v1.BeamFnControlGrpc;
 import org.apache.beam.model.pipeline.v1.Endpoints.ApiServiceDescriptor;
@@ -151,11 +152,19 @@ public class BeamFnControlClient {
           .setInstructionId(value.getInstructionId())
           .build();
     } catch (Exception e) {
-      LOG.error(
-          "Exception while trying to handle {} {}",
-          BeamFnApi.InstructionRequest.class.getSimpleName(),
-          value.getInstructionId(),
-          e);
+      if (WorkCancelledException.isWorkCancelledException(e)) {
+        LOG.info(
+            "Instruction {} {} cancelled by runner",
+            BeamFnApi.InstructionRequest.class.getSimpleName(),
+            value.getInstructionId(),
+            e);
+      } else {
+        LOG.error(
+            "Exception while trying to handle {} {}",
+            BeamFnApi.InstructionRequest.class.getSimpleName(),
+            value.getInstructionId(),
+            e);
+      }
       return BeamFnApi.InstructionResponse.newBuilder()
           .setInstructionId(value.getInstructionId())
           .setError(getStackTraceAsString(e))
diff --git 
a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/data/PCollectionConsumerRegistry.java
 
b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/data/PCollectionConsumerRegistry.java
index 3ba8b4e76c3..58c06bac310 100644
--- 
a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/data/PCollectionConsumerRegistry.java
+++ 
b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/data/PCollectionConsumerRegistry.java
@@ -36,6 +36,7 @@ import 
org.apache.beam.fn.harness.control.Metrics.BundleDistribution;
 import org.apache.beam.fn.harness.debug.DataSampler;
 import org.apache.beam.fn.harness.debug.ElementSample;
 import org.apache.beam.fn.harness.debug.OutputSampler;
+import org.apache.beam.fn.harness.state.WorkCancelledException;
 import org.apache.beam.model.fnexecution.v1.BeamFnApi.ProcessBundleDescriptor;
 import org.apache.beam.model.pipeline.v1.MetricsApi.MonitoringInfo;
 import org.apache.beam.model.pipeline.v1.RunnerApi;
@@ -281,14 +282,16 @@ public class PCollectionConsumerRegistry {
       @Nullable OutputSampler<T> outputSampler,
       @Nullable ElementSample<T> elementSample)
       throws Exception {
-    ExecutionStateSampler.ExecutionStateTrackerStatus status = 
executionStateTracker.getStatus();
-    String processBundleId = status == null ? null : 
status.getProcessBundleId();
-    if (outputSampler != null) {
-      outputSampler.exception(elementSample, e, ptransformId, processBundleId);
-    }
+    if (!WorkCancelledException.isWorkCancelledException(e)) {
+      ExecutionStateSampler.ExecutionStateTrackerStatus status = 
executionStateTracker.getStatus();
+      String processBundleId = status == null ? null : 
status.getProcessBundleId();
+      if (outputSampler != null) {
+        outputSampler.exception(elementSample, e, ptransformId, 
processBundleId);
+      }
 
-    if (executionState.error()) {
-      LOG.error("Failed to process element for bundle \"{}\"", 
processBundleId, e);
+      if (executionState.error()) {
+        LOG.error("Failed to process element for bundle \"{}\"", 
processBundleId, e);
+      }
     }
     throw e;
   }
diff --git 
a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/state/BeamFnStateGrpcClientCache.java
 
b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/state/BeamFnStateGrpcClientCache.java
index 453dd31e0d4..c9adb9d7d14 100644
--- 
a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/state/BeamFnStateGrpcClientCache.java
+++ 
b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/state/BeamFnStateGrpcClientCache.java
@@ -187,6 +187,8 @@ public class BeamFnStateGrpcClientCache {
         }
         if (value.getError().isEmpty()) {
           responseFuture.complete(value);
+        } else if (value.getErrorReason() == 
StateResponse.ErrorReason.CANCELLED) {
+          responseFuture.completeExceptionally(new 
WorkCancelledException(value.getError()));
         } else {
           responseFuture.completeExceptionally(new 
IllegalStateException(value.getError()));
         }
diff --git 
a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/state/WorkCancelledException.java
 
b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/state/WorkCancelledException.java
new file mode 100644
index 00000000000..587f5962f03
--- /dev/null
+++ 
b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/state/WorkCancelledException.java
@@ -0,0 +1,44 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.fn.harness.state;
+
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Throwables;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables;
+import org.checkerframework.checker.nullness.qual.Nullable;
+
+/**
+ * Indicates that the work item is no longer valid on the runner and should be 
cancelled without
+ * logging an error.
+ */
+public class WorkCancelledException extends RuntimeException {
+
+  public WorkCancelledException(String message) {
+    super(message);
+  }
+
+  public WorkCancelledException(Throwable cause) {
+    super(cause);
+  }
+
+  /** Returns whether an exception was caused by a {@link 
WorkCancelledException}. */
+  public static boolean isWorkCancelledException(@Nullable Throwable t) {
+    return t != null
+        && !Iterables.isEmpty(
+            Iterables.filter(Throwables.getCausalChain(t), 
WorkCancelledException.class));
+  }
+}
diff --git 
a/sdks/java/harness/src/test/java/org/apache/beam/fn/harness/control/BeamFnControlClientTest.java
 
b/sdks/java/harness/src/test/java/org/apache/beam/fn/harness/control/BeamFnControlClientTest.java
index 5afd2c475f6..c54c35c46e6 100644
--- 
a/sdks/java/harness/src/test/java/org/apache/beam/fn/harness/control/BeamFnControlClientTest.java
+++ 
b/sdks/java/harness/src/test/java/org/apache/beam/fn/harness/control/BeamFnControlClientTest.java
@@ -88,6 +88,18 @@ public class BeamFnControlClientTest {
           .setInstructionId("3L")
           .setError(getStackTraceAsString(FAILURE))
           .build();
+  private static final org.apache.beam.fn.harness.state.WorkCancelledException 
CANCELLED =
+      new org.apache.beam.fn.harness.state.WorkCancelledException("Work item 
cancelled");
+  private static final BeamFnApi.InstructionRequest CANCELLED_REQUEST =
+      BeamFnApi.InstructionRequest.newBuilder()
+          .setInstructionId("4L")
+          .setSampleData(BeamFnApi.SampleDataRequest.getDefaultInstance())
+          .build();
+  private static final BeamFnApi.InstructionResponse CANCELLED_RESPONSE =
+      BeamFnApi.InstructionResponse.newBuilder()
+          .setInstructionId("4L")
+          .setError(getStackTraceAsString(CANCELLED))
+          .build();
 
   @Rule public TestRule restoreMDCAfterTest = new RestoreBeamFnLoggingMDC();
 
@@ -137,6 +149,12 @@ public class BeamFnControlClientTest {
             assertEquals(value.getInstructionId(), 
BeamFnLoggingMDC.getInstructionId());
             throw FAILURE;
           });
+      handlers.put(
+          BeamFnApi.InstructionRequest.RequestCase.SAMPLE_DATA,
+          value -> {
+            assertEquals(value.getInstructionId(), 
BeamFnLoggingMDC.getInstructionId());
+            throw CANCELLED;
+          });
 
       ExecutorService executor = Executors.newCachedThreadPool();
       BeamFnControlClient client =
@@ -163,6 +181,10 @@ public class BeamFnControlClientTest {
       outboundServerObserver.onNext(FAILURE_REQUEST);
       assertEquals(FAILURE_RESPONSE, values.take());
 
+      // Ensure that WorkCancelledException is caught and translated to a 
failure response
+      outboundServerObserver.onNext(CANCELLED_REQUEST);
+      assertEquals(CANCELLED_RESPONSE, values.take());
+
       // Ensure that the server completing the stream translates to the 
completable future
       // being completed allowing for a successful shutdown of the client.
       outboundServerObserver.onCompleted();
diff --git 
a/sdks/java/harness/src/test/java/org/apache/beam/fn/harness/data/PCollectionConsumerRegistryTest.java
 
b/sdks/java/harness/src/test/java/org/apache/beam/fn/harness/data/PCollectionConsumerRegistryTest.java
index 7ba2f8921be..804fd695e34 100644
--- 
a/sdks/java/harness/src/test/java/org/apache/beam/fn/harness/data/PCollectionConsumerRegistryTest.java
+++ 
b/sdks/java/harness/src/test/java/org/apache/beam/fn/harness/data/PCollectionConsumerRegistryTest.java
@@ -20,6 +20,7 @@ package org.apache.beam.fn.harness.data;
 import static org.apache.beam.sdk.values.WindowedValues.valueInGlobalWindow;
 import static org.hamcrest.MatcherAssert.assertThat;
 import static org.hamcrest.Matchers.containsInAnyOrder;
+import static org.hamcrest.Matchers.is;
 import static org.junit.Assert.assertEquals;
 import static org.junit.Assert.assertFalse;
 import static org.junit.Assert.assertNotNull;
@@ -50,6 +51,7 @@ import org.apache.beam.fn.harness.debug.DataSampler;
 import org.apache.beam.fn.harness.logging.BeamFnLoggingMDC;
 import org.apache.beam.fn.harness.logging.LoggingClient;
 import org.apache.beam.fn.harness.logging.LoggingClientFactory;
+import org.apache.beam.fn.harness.state.WorkCancelledException;
 import org.apache.beam.model.fnexecution.v1.BeamFnApi;
 import org.apache.beam.model.fnexecution.v1.BeamFnApi.ProcessBundleDescriptor;
 import org.apache.beam.model.fnexecution.v1.BeamFnLoggingGrpc;
@@ -703,6 +705,88 @@ public class PCollectionConsumerRegistryTest {
     }
   }
 
+  @Test
+  public void doesNotLogWorkCancelledException() throws Exception {
+    final String pTransformId = "pTransformId";
+    final String message = "Work item cancelled";
+    final String instructionId = "instruction";
+    final Exception thrownException = new RuntimeException(new 
WorkCancelledException(message));
+
+    AtomicBoolean clientClosedStream = new AtomicBoolean();
+    Collection<BeamFnApi.LogEntry> values = new ConcurrentLinkedQueue<>();
+    AtomicReference<StreamObserver<BeamFnApi.LogControl>> 
outboundServerObserver =
+        new AtomicReference<>();
+    CallStreamObserver<BeamFnApi.LogEntry.List> inboundServerObserver =
+        TestStreams.withOnNext(
+                (BeamFnApi.LogEntry.List logEntries) ->
+                    values.addAll(logEntries.getLogEntriesList()))
+            .withOnCompleted(
+                () -> {
+                  clientClosedStream.set(true);
+                  outboundServerObserver.get().onCompleted();
+                })
+            .build();
+
+    Endpoints.ApiServiceDescriptor apiServiceDescriptor =
+        Endpoints.ApiServiceDescriptor.newBuilder()
+            .setUrl(this.getClass().getName() + "-" + 
UUID.randomUUID().toString())
+            .build();
+    Server server =
+        InProcessServerBuilder.forName(apiServiceDescriptor.getUrl())
+            .addService(
+                new BeamFnLoggingGrpc.BeamFnLoggingImplBase() {
+                  @Override
+                  public StreamObserver<BeamFnApi.LogEntry.List> logging(
+                      StreamObserver<BeamFnApi.LogControl> outboundObserver) {
+                    outboundServerObserver.set(outboundObserver);
+                    return inboundServerObserver;
+                  }
+                })
+            .build();
+    server.start();
+    ManagedChannel channel = 
InProcessChannelBuilder.forName(apiServiceDescriptor.getUrl()).build();
+
+    ExecutionStateSampler sampler =
+        new ExecutionStateSampler(PipelineOptionsFactory.create(), 
System::currentTimeMillis, null);
+    ExecutionStateSampler.ExecutionStateTracker stateTracker = 
sampler.create();
+    stateTracker.start("process-bundle");
+    ExecutionStateSampler.ExecutionState state =
+        stateTracker.create("shortId", pTransformId, pTransformId, "process");
+    state.activate();
+
+    BeamFnLoggingMDC.setInstructionId(instructionId);
+    BeamFnLoggingMDC.setStateTracker(stateTracker);
+
+    try (LoggingClient ignored =
+        LoggingClientFactory.createAndStart(
+            PipelineOptionsFactory.create(),
+            apiServiceDescriptor,
+            (Endpoints.ApiServiceDescriptor descriptor) -> channel)) {
+
+      ShortIdMap shortIds = new ShortIdMap();
+      BundleProgressReporter.InMemory reporterAndRegistrar = new 
BundleProgressReporter.InMemory();
+      PCollectionConsumerRegistry consumers =
+          new PCollectionConsumerRegistry(
+              stateTracker, shortIds, reporterAndRegistrar, TEST_DESCRIPTOR);
+      FnDataReceiver<WindowedValue<String>> consumer = 
mock(FnDataReceiver.class);
+
+      consumers.register(P_COLLECTION_A, pTransformId, pTransformId + "Name", 
consumer);
+
+      FnDataReceiver<WindowedValue<String>> wrapperConsumer =
+          (FnDataReceiver<WindowedValue<String>>)
+              (FnDataReceiver) 
consumers.getMultiplexingConsumer(P_COLLECTION_A);
+
+      doThrow(thrownException).when(consumer).accept(any());
+      expectedException.expect(is(thrownException));
+
+      wrapperConsumer.accept(valueInGlobalWindow("elem"));
+
+    } finally {
+      assertTrue(values.isEmpty());
+      server.shutdownNow();
+    }
+  }
+
   private static class TestElementByteSizeObservableIterable<T>
       extends ElementByteSizeObservableIterable<T, 
ElementByteSizeObservableIterator<T>> {
     private List<T> elements;
diff --git 
a/sdks/java/harness/src/test/java/org/apache/beam/fn/harness/state/BeamFnStateGrpcClientCacheTest.java
 
b/sdks/java/harness/src/test/java/org/apache/beam/fn/harness/state/BeamFnStateGrpcClientCacheTest.java
index 2d7de3489b2..6360c459511 100644
--- 
a/sdks/java/harness/src/test/java/org/apache/beam/fn/harness/state/BeamFnStateGrpcClientCacheTest.java
+++ 
b/sdks/java/harness/src/test/java/org/apache/beam/fn/harness/state/BeamFnStateGrpcClientCacheTest.java
@@ -22,6 +22,7 @@ import static org.hamcrest.Matchers.containsString;
 import static org.junit.Assert.assertNotNull;
 import static org.junit.Assert.assertNotSame;
 import static org.junit.Assert.assertSame;
+import static org.junit.Assert.assertTrue;
 import static org.junit.Assert.fail;
 
 import java.util.UUID;
@@ -148,13 +149,16 @@ public class BeamFnStateGrpcClientCacheTest {
         client.handle(StateRequest.newBuilder().setInstructionId(SUCCESS));
     CompletableFuture<StateResponse> unsuccessfulResponse =
         client.handle(StateRequest.newBuilder().setInstructionId(FAIL));
+    CompletableFuture<StateResponse> cancelledResponse =
+        client.handle(StateRequest.newBuilder().setInstructionId("CANCELLED"));
 
     // Wait for the client to connect.
     StreamObserver<StateResponse> outboundServerObserver = 
outboundServerObservers.take();
     // Ensure the client doesn't break when sent garbage.
     outboundServerObserver.onNext(StateResponse.newBuilder().setId("UNKNOWN 
ID").build());
 
-    // We expect to receive and handle two requests
+    // We expect to receive and handle three requests
+    handleServerRequest(outboundServerObserver, values.take());
     handleServerRequest(outboundServerObserver, values.take());
     handleServerRequest(outboundServerObserver, values.take());
 
@@ -166,6 +170,13 @@ public class BeamFnStateGrpcClientCacheTest {
     } catch (ExecutionException e) {
       assertThat(e.toString(), containsString(TEST_ERROR));
     }
+    try {
+      cancelledResponse.get();
+      fail("Expected cancelled response");
+    } catch (ExecutionException e) {
+      assertThat(e.toString(), containsString(TEST_ERROR));
+      assertTrue(WorkCancelledException.isWorkCancelledException(e));
+    }
   }
 
   @Test
@@ -237,6 +248,14 @@ public class BeamFnStateGrpcClientCacheTest {
         outboundObserver.onNext(
             
StateResponse.newBuilder().setId(value.getId()).setError(TEST_ERROR).build());
         return;
+      case "CANCELLED":
+        outboundObserver.onNext(
+            StateResponse.newBuilder()
+                .setId(value.getId())
+                .setError(TEST_ERROR)
+                .setErrorReason(StateResponse.ErrorReason.CANCELLED)
+                .build());
+        return;
       default:
         
outboundObserver.onNext(StateResponse.newBuilder().setId(value.getId()).build());
         return;

Reply via email to