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

