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;