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
相关产品推荐
相关产品推荐

