You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.15 13:27:05