Apache Beam窗口功能失效,抛出IntervalWindowBase相关Cython错误
问题描述
- 使用Apache Beam版本:2.37.0、2.43.0
- Runner:Dataflow
- Pipeline流程:从Pub/Sub读取流数据 → 经过两个30分钟的FixedWindow聚合元素 → 基于消息分配的键分组
- 错误现象:Pipeline正常运行约1小时后(期间窗口输出及后续处理均正常),突然抛出与IntervalWindowBase相关的Cython错误
- 窗口处理输入类型:
Tuple[str, Dict[str, str]]
错误栈信息
Error message from worker: generic::unknown: Traceback (most recent call last): File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/runners/worker/sdk_worker.py", line 287, in _execute response = task() File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/runners/worker/sdk_worker.py", line 360, in <lambda> lambda: self.create_worker().do_instruction(request), request) File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/runners/worker/sdk_worker.py", line 596, in do_instruction return getattr(self, request_type)( File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/runners/worker/sdk_worker.py", line 634, in process_bundle bundle_processor.process_bundle(instruction_id)) File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/runners/worker/bundle_processor.py", line 1003, in process_bundle input_op_by_transform_id[element.transform_id].process_encoded( File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/runners/worker/bundle_processor.py", line 225, in process_encoded decoded_value = self.windowed_coder_impl.decode_from_stream( File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/coders/coder_impl.py", line 1465, in decode_from_stream value = self._value_coder.decode_from_stream(in_stream, nested) File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/coders/coder_impl.py", line 1008, in decode_from_stream return self._construct_from_components([ File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/coders/coder_impl.py", line 1009, in <listcomp> c.decode_from_stream( File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/coders/coder_impl.py", line 1194, in decode_from_stream elements = [ File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/coders/coder_impl.py", line 1195, in <listcomp> self._elem_coder.decode_from_stream(in_stream, True) File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/coders/coder_impl.py", line 1557, in decode_from_stream return self._value_coder.decode(in_stream.read(value_length)) File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/coders/coder_impl.py", line 240, in decode return self.decode_from_stream(create_InputStream(encoded), False) File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/coders/coder_impl.py", line 583, in decode_from_stream return self.fallback_coder_impl.decode_from_stream(stream, nested) File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/coders/coder_impl.py", line 273, in decode_from_stream return self._decoder(stream.read_all(nested)) AttributeError: Can't get attribute '__pyx_unpickle__IntervalWindowBase' on <module 'apache_beam.utils.windowed_value' from '/pipeline/.venv/lib/python3.9/site-packages/apache_beam/utils/windowed_value.py'> generic::unknown: Traceback (most recent call last): File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/runners/worker/sdk_worker.py", line 287, in _execute response = task() File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/runners/worker/sdk_worker.py", line 360, in <lambda> lambda: self.create_worker().do_instruction(request), request) File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/runners/worker/sdk_worker.py", line 596, in do_instruction return getattr(self, request_type)( File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/runners/worker/sdk_worker.py", line 634, in process_bundle bundle_processor.process_bundle(instruction_id)) File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/runners/worker/bundle_processor.py", line 1003, in process_bundle input_op_by_transform_id[element.transform_id].process_encoded( File "/pipeline/.venv/lib/python3.9/site-packages/apache_beam/runners/worker/bundle_processor.py", line 225, in process_encoded decoded_value = self.windowed_coder_impl.decode_from_stream( File "apache_beam/coders/coder_impl.py", line 1446, in apache_beam.coders.coder_impl.WindowedValueCoderImpl.decode_from_stream File "apache_beam/coders/coder_impl.py", line 1465, in apache_beam.coders.coder_impl.WindowedValueCoderImpl.decode_from_stream File "apache_beam/coders/coder_impl.py", line 1009, in apache_beam.coders.coder_impl.AbstractComponentCoderImpl.decode_from_stream File "apache_beam/coders/coder_impl.py", line 1195, in apache_beam.coders.coder_impl.SequenceCoderImpl.decode_from_stream File "apache_beam/coders/coder_impl.py", line 1557, in apache_beam.coders.coder_impl.LengthPrefixCoderImpl.decode_from_stream File "apache_beam/coders/coder_impl.py", line 240, in apache_beam.coders.coder_impl.StreamCoderImpl.decode File "apache_beam/coders/coder_impl.py", line 583, in apache_beam.coders.coder_impl.FastPrimitivesCoderImpl.decode_from_stream File "apache_beam/coders/coder_impl.py", line 273, in apache_beam.coders.coder_impl.CallbackCoderImpl.decode_from_stream File "stringsource", line 17, in apache_beam.utils.windowed_value._IntervalWindowBase.__setstate_cython__ TypeError: Expected tuple, got dict
已尝试的操作
- 缩小窗口大小
- 使用更严格的输入输出类型
上述操作均未解决问题,窗口功能仍会在运行一段时间后突然报错。
内容的提问来源于stack exchange,提问作者RMCP
相关产品推荐
相关产品推荐

