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

shunping 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 6317b1942f2 Reduce severity for state requests cancelled by runner 
(#40254)
6317b1942f2 is described below

commit 6317b1942f2aaf4de6afad341c895488af2abdb4
Author: parveensania <[email protected]>
AuthorDate: Fri Sep 25 07:38:19 2026 -0700

    Reduce severity for state requests cancelled by runner (#40254)
---
 .../beam/model/fn_execution/v1/beam_fn_api.proto   | 13 ++++++++++
 sdks/python/apache_beam/runners/common.py          |  7 ++++++
 .../apache_beam/runners/worker/sdk_worker.py       | 12 ++++++++++
 .../apache_beam/runners/worker/sdk_worker_test.py  | 28 ++++++++++++++++++++++
 4 files changed, 60 insertions(+)

diff --git 
a/model/fn-execution/src/main/proto/org/apache/beam/model/fn_execution/v1/beam_fn_api.proto
 
b/model/fn-execution/src/main/proto/org/apache/beam/model/fn_execution/v1/beam_fn_api.proto
index 9b8493f1e49..0c3633ce247 100644
--- 
a/model/fn-execution/src/main/proto/org/apache/beam/model/fn_execution/v1/beam_fn_api.proto
+++ 
b/model/fn-execution/src/main/proto/org/apache/beam/model/fn_execution/v1/beam_fn_api.proto
@@ -930,6 +930,19 @@ message StateResponse {
   // failed.
   string error = 2;
 
+  enum ErrorReason {
+    ERROR_REASON_UNSPECIFIED = 0;
+
+    // Indicates that the state request failed because the runner cancelled or
+    // invalidated the work item (e.g., due to autoscaling, rebalancing, or
+    // hedging). SDKs may choose to reduce logging severity for this error.
+    CANCELLED = 1;
+  }
+
+  // (Optional) If error is specified, provides a structured reason for why the
+  // request failed.
+  ErrorReason error_reason = 3;
+
   // A corresponding response matching the request will be populated.
   oneof response {
     // A response to getting state.
diff --git a/sdks/python/apache_beam/runners/common.py 
b/sdks/python/apache_beam/runners/common.py
index 87c3dd4f2f4..b1ae4d4dd34 100644
--- a/sdks/python/apache_beam/runners/common.py
+++ b/sdks/python/apache_beam/runners/common.py
@@ -1618,11 +1618,18 @@ class DoFnRunner:
     _, _, tb = exc_info
 
     new_exn = new_exn.with_traceback(tb)
+    if isinstance(exn, WorkCancelledException):
+      raise new_exn
     self._maybe_sample_exception(exc_info, windowed_value)
     _LOGGER.exception(new_exn)
     raise new_exn
 
 
+class WorkCancelledException(RuntimeError):
+  """Indicates that a state request or work item was cancelled by the 
runner."""
+  pass
+
+
 class OutputHandler(object):
   def handle_process_outputs(
       self, windowed_input_element, results, watermark_estimator=None):
diff --git a/sdks/python/apache_beam/runners/worker/sdk_worker.py 
b/sdks/python/apache_beam/runners/worker/sdk_worker.py
index 84b4ae1bf23..bc7e80c1edc 100644
--- a/sdks/python/apache_beam/runners/worker/sdk_worker.py
+++ b/sdks/python/apache_beam/runners/worker/sdk_worker.py
@@ -52,6 +52,7 @@ from apache_beam.metrics.execution import MetricsEnvironment
 from apache_beam.portability.api import beam_fn_api_pb2
 from apache_beam.portability.api import beam_fn_api_pb2_grpc
 from apache_beam.portability.api import metrics_pb2
+from apache_beam.runners.common import WorkCancelledException
 from apache_beam.runners.worker import bundle_processor
 from apache_beam.runners.worker import data_plane
 from apache_beam.runners.worker import statesampler
@@ -308,6 +309,14 @@ class SdkHarness(object):
     with statesampler.instruction_id(request.instruction_id):
       try:
         response = task()
+      except WorkCancelledException:
+        traceback_string = traceback.format_exc()
+        _LOGGER.info(
+            'Instruction %s cancelled by runner. Original traceback is\n%s\n',
+            request.instruction_id,
+            traceback_string)
+        response = beam_fn_api_pb2.InstructionResponse(
+            instruction_id=request.instruction_id, error=traceback_string)
       except:  # pylint: disable=bare-except
         traceback_string = traceback.format_exc()
         print(traceback_string, file=sys.stderr)
@@ -1153,6 +1162,9 @@ class GrpcStateHandler(StateHandler):
         raise RuntimeError()
     response = req_future.get()
     if response.error:
+      if (response.error_reason ==
+          beam_fn_api_pb2.StateResponse.ErrorReason.CANCELLED):
+        raise WorkCancelledException(response.error)
       raise RuntimeError(response.error)
     else:
       return response
diff --git a/sdks/python/apache_beam/runners/worker/sdk_worker_test.py 
b/sdks/python/apache_beam/runners/worker/sdk_worker_test.py
index 4b00051a337..1f8ac774c21 100644
--- a/sdks/python/apache_beam/runners/worker/sdk_worker_test.py
+++ b/sdks/python/apache_beam/runners/worker/sdk_worker_test.py
@@ -432,6 +432,34 @@ class SdkWorkerTest(unittest.TestCase):
     mock_bundle_processor.process_bundle.assert_called_once_with(
         instruction_id, 'stream_xyz')
 
+  def test_cancelled_state_request_logs_info(self):
+    state_handler = sdk_worker.GrpcStateHandler(mock.MagicMock())
+    state_handler._context.process_instruction_id = 'bundle_1'
+    cancelled_response = beam_fn_api_pb2.StateResponse(
+        id='1',
+        error='Work item cancelled by runner',
+        error_reason=beam_fn_api_pb2.StateResponse.ErrorReason.CANCELLED)
+    state_handler._request = mock.MagicMock(
+        return_value=sdk_worker._Future().set(cancelled_response))
+
+    with self.assertRaises(sdk_worker.WorkCancelledException):
+      state_handler._blocking_request(beam_fn_api_pb2.StateRequest())
+
+    with mock.patch('grpc.channel_ready_future'):
+      harness = sdk_worker.SdkHarness('localhost:0')
+    request = beam_fn_api_pb2.InstructionRequest(instruction_id='bundle_1')
+    with mock.patch.object(sdk_worker._LOGGER, 'info') as mock_info, \
+         mock.patch.object(sdk_worker._LOGGER, 'error') as mock_error:
+      harness._execute(
+          lambda: state_handler._blocking_request(
+              beam_fn_api_pb2.StateRequest()),
+          request)
+      mock_info.assert_called_once()
+      mock_error.assert_not_called()
+    response = harness._responses.get_nowait()
+    self.assertEqual(response.instruction_id, 'bundle_1')
+    self.assertIn('Work item cancelled by runner', response.error)
+
 
 class CachingStateHandlerTest(unittest.TestCase):
   def test_caching(self):

Reply via email to