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

Apache Beam中Deduplicate函数结合文件读取时的报错问题

Apache Beam Deduplicate转换读取磁盘文件时异常问题

问题现象

使用apache_beam.transforms.deduplicate.Deduplicate转换时,通过beam.Create手动构造数据集可正常完成去重;但通过beam.io.ReadFromText(或读取Avro等磁盘文件)加载数据时,转换会抛出异常,无法正常工作。该问题仅出现在Deduplicate和DeduplicatePerKey转换中,ParDo、Map等其他转换运行完全正常。

代码示例

from apache_beam import Pipeline, Create, io, Map
from apache_beam.transforms.deduplicate import Deduplicate
from typing import AnyStr

with Pipeline() as pipeline:
    (
        pipeline
        # | 'Load' >> Create(['a', 'b', 'b']) ## <- 此分支正常运行
        | 'Load' >> io.ReadFromText('./input.txt') ## <- 此分支触发异常
        | 'Dedup' >> Deduplicate(processing_time_duration=1000).with_input_types(AnyStr)
        | 'Print' >> Map(print)
    )

抛出的异常信息

Traceback (most recent call last):
  File "/Users/ds/.pyenv/versions/3.9.14/lib/python3.9/site-packages/apache_beam/runners/direct/executor.py", line 370, in call
    self.attempt_call(
  File "/Users/ds/.pyenv/versions/3.9.14/lib/python3.9/site-packages/apache_beam/runners/direct/executor.py", line 404, in attempt_call
    evaluator.start_bundle()
  File "/Users/ds/.pyenv/versions/3.9.14/lib/python3.9/site-packages/apache_beam/runners/direct/transform_evaluator.py", line 867, in start_bundle
    self.runner.start()
  File "/Users/ds/.pyenv/versions/3.9.14/lib/python3.9/site-packages/apache_beam/runners/common.py", line 1475, in start
    self._invoke_bundle_method(self.do_fn_invoker.invoke_start_bundle)
  File "/Users/ds/.pyenv/versions/3.9.14/lib/python3.9/site-packages/apache_beam/runners/common.py", line 1460, in _invoke_bundle_method
    self._reraise_augmented(exn)
  File "/Users/ds/.pyenv/versions/3.9.14/lib/python3.9/site-packages/apache_beam/runners/common.py", line 1507, in _reraise_augmented
    raise new_exn.with_traceback(tb)
  File "/Users/ds/.pyenv/versions/3.9.14/lib/python3.9/site-packages/apache_beam/runners/common.py", line 1458, in _invoke_bundle_method
    bundle_method()
  File "/Users/ds/.pyenv/versions/3.9.14/lib/python3.9/site-packages/apache_beam/runners/common.py", line 559, in invoke_start_bundle
    self.signature.start_bundle_method.method_value())
  File "/Users/ds/.pyenv/versions/3.9.14/lib/python3.9/site-packages/apache_beam/runners/direct/sdf_direct_runner.py", line 122, in start_bundle
    self._invoker = DoFnInvoker.create_invoker(
TypeError: create_invoker() got an unexpected keyword argument 'output_processor' [while running 'Load/Read/SDFBoundedSourceReader/ParDo(SDFBoundedSourceDoFn)/pair']

环境信息

  • Python版本:3.9.14
  • Apache Beam版本:2.41.0
  • 运行平台:Apple M1 ARM

内容的提问来源于stack exchange,提问作者Dmitry Stropalov

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 16:05:59