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

Apache Beam使用Timers报错:组件数量与编码器数量不匹配

Apache Beam Timer使用时的Coder不匹配错误分析

问题场景

尝试在Apache Beam中添加Timers延迟处理部分数据时,触发ValueError: Number of components does not match number of coders错误。但将WaitUntilDevicesExist移至GroupMessagesByShardedKey之后,程序可正常运行。

错误日志

Error message from worker: generic::unknown: Traceback (most recent call last):
File "apache_beam/runners/common.py", line 1435, in apache_beam.runners.common.DoFnRunner.process
  File "apache_beam/runners/common.py", line 636, in apache_beam.runners.common.SimpleInvoker.invoke_process
  File "apache_beam/runners/common.py", line 1621, in apache_beam.runners.common._OutputHandler.handle_process_outputs
  File "apache_beam/runners/common.py", line 1734, in apache_beam.runners.common._OutputHandler._write_value_to_tag
  File "apache_beam/runners/worker/operations.py", line 266, in apache_beam.runners.worker.operations.SingletonElementConsumerSet.receive
  File "apache_beam/runners/worker/operations.py", line 528, in apache_beam.runners.worker.operations.Operation.process
  File "/usr/local/lib/python3.9/site-packages/apache_beam/runners/worker/bundle_processor.py", line 158, in process
    self.windowed_coder_impl.encode_to_stream(
  File "apache_beam/coders/coder_impl.py", line 1448, in apache_beam.coders.coder_impl.WindowedValueCoderImpl.encode_to_stream
  File "apache_beam/coders/coder_impl.py", line 1467, in apache_beam.coders.coder_impl.WindowedValueCoderImpl.encode_to_stream
  File "apache_beam/coders/coder_impl.py", line 1023, in apache_beam.coders.coder_impl.AbstractComponentCoderImpl.encode_to_stream
ValueError: Number of components does not match number of coders.

涉及代码

WaitUntilDevicesExist DoFn代码

class WaitUntilDevicesExist(beam.DoFn):
    BUFFER_STATE = beam.transforms.userstate.BagStateSpec('buffer', beam.coders.StrUtf8Coder())
    TIMER = beam.transforms.userstate.TimerSpec('timer', beam.TimeDomain.REAL_TIME)

    BUFFER_TIMER = 15  # seconds

    ...

    def process(self, key_value, timer=beam.DoFn.TimerParam(TIMER), buffer=beam.DoFn.StateParam(BUFFER_STATE)):
        shard_id, batch = key_value

        for message in batch:
            logging.info(f"Checking = {message}")
            
            ...

            if (...):
                timer.set(timestamp.Timestamp.now() + timestamp.Duration(seconds=self.BUFFER_TIMER)) 
                buffer.add(DeviceCheckHelper(message).to_string())
            else:
                yield message

    @beam.transforms.userstate.on_timer(TIMER)
    def expiry_callback(self, timer=beam.DoFn.TimerParam(TIMER), buffer=beam.DoFn.StateParam(BUFFER_STATE)):
        events = buffer.read()
        logging.info("Timer")
        new_buffer = []

        for row in events:
            message = DeviceCheckHelper.from_string(row)
            logging.info(message)
            
            ....

            if (...):
                if retry == 3:
                    logging.info(f"Waited 3 times, yielding ")
                    yield message.message
                else:
                    message.increase_retry()
                    new_buffer.append(message.to_string())
                    logging.info(f"retry = {message}")

        buffer.clear()
        timer.clear()

        logging.info(f"New buffer = {new_buffer}")
        if new_buffer:
            for row in new_buffer:
                logging.info(f"Adding {row}")
                buffer.add(row)

            timer.set(timestamp.Timestamp.now() + timestamp.Duration(seconds=self.BUFFER_TIMER))

Pipeline代码

# 1 filter messages
filtered_messages = (
    transformed_messages[TransformData.TAG_OK]
    | f"Clean Devices {tenant}" >> beam.ParDo(FilterMessages()).with_outputs(FilterMessages.DEVICE_TAG, FilterMessages.OBSERVATION_TAG)
)

# 2 Write observations
observation_results = (
    filtered_messages[FilterMessages.OBSERVATION_TAG]
    | f"{tenant} Check Devices" >> beam.ParDo(WaitUntilDevicesExist(...))
    | f"{tenant} Window Observations messages" >> GroupMessagesByShardedKey(max_messages=200, max_waiting_time=10, shard_key="obs", num_shards=10)
    | f"{tenant} Write Observations" >> beam.ParDo(Write(...)).with_outputs(FAILED_TAG)
)

问题原因与修复方案

核心原因

当WaitUntilDevicesExist放在GroupMessagesByShardedKey之前时,输入到DoFn的是未分组的单条消息,但DoFn的process方法却按(shard_id, batch)键值对结构解析输入,导致数据结构与预期不匹配,进而触发Coder组件数量不匹配的错误。

而将DoFn移到分组之后时,输入正好是GroupMessagesByShardedKey输出的(shard_id, batch)键值对,与DoFn的输入解析逻辑完全匹配,因此运行正常。

修复方案

方案1:调整DoFn适配单条消息输入

如果需要在分组前使用该DoFn,需修改process方法,直接处理单条消息而非键值对:

def process(self, message, timer=beam.DoFn.TimerParam(TIMER), buffer=beam.DoFn.StateParam(BUFFER_STATE)):
    logging.info(f"Checking = {message}")
    
    ...

    if (...):
        timer.set(timestamp.Timestamp.now() + timestamp.Duration(seconds=self.BUFFER_TIMER)) 
        buffer.add(DeviceCheckHelper(message).to_string())
    else:
        yield message

同时需确保数据已按业务所需的键完成分区(如通过beam.GroupByKey),否则状态无法正确关联到对应消息组。

方案2:保持DoFn在分组之后执行

更合理的方式是维持当前可运行的结构,将带状态的延迟处理放在分组之后。因为分组后的数据已按shard聚合,状态管理更高效,也更符合延迟处理批量消息的业务逻辑。

内容的提问来源于stack exchange,提问作者Alex Fragotsis

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 18:42:02