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):