Python版Dataflow流处理使用MatchContinuously运行异常报错求助
我尝试使用Dataflow Python库以流处理方式读取存储桶中的数据,参考了官方2.33.0版本的
apache_beam.io.fileio.MatchContinuously接口文档,使用以下代码片段定期轮询存储桶:(pipeline | 'Match Files' >> fileio.MatchContinuously(file_pattern="gs://xyz/abc/*.txt", interval=10.0, has_deduplication=True) | 'Read Matches' >> fileio.ReadMatches() ...... )代码运行很短时间后就会失败,报错栈如下:
<PCollection[Read Matches/ParDo(_ReadMatchesFn).None] at 0x7fad0ccb7e80> WARNING:root:Make sure that locally built Python SDK docker image has Python 3.8 interpreter. Traceback (most recent call last): File "<stdin>", line 3, in <module> File "/usr/local/Caskroom/miniconda/base/lib/python3.8/site-packages/apache_beam/pipeline.py", line 586, in __exit__ self.result = self.run() File "/usr/local/Caskroom/miniconda/base/lib/python3.8/site-packages/apache_beam/pipeline.py", line 565, in run return self.runner.run_pipeline(self, self._options) File "/usr/local/Caskroom/miniconda/base/lib/python3.8/site-packages/apache_beam/runners/direct/direct_runner.py", line 131, in run_pipeline return runner.run_pipeline(pipeline, options) File "/usr/local/Caskroom/miniconda/base/lib/python3.8/site-packages/apache_beam/runners/portability/fn_api_runner/fn_runner.py", line 195, in run_pipeline self._latest_run_result = self.run_via_runner_api( File "/usr/local/Caskroom/miniconda/base/lib/python3.8/site-packages/apache_beam/runners/portability/fn_api_runner/fn_runner.py", line 206, in run_via_runner_api return self.run_stages(stage_context, stages) File "/usr/local/Caskroom/miniconda/base/lib/python3.8/site-packages/apache_beam/runners/portability/fn_api_runner/fn_runner.py", line 384, in run_stages stage_results = self._run_stage( File "/usr/local/Caskroom/miniconda/base/lib/python3.8/site-packages/apache_beam/runners/portability/fn_api_runner/fn_runner.py", line 663, in _run_stage assert (runner_execution_context.watermark_manager.get_stage_node( AssertionError: wrong timestamp for StageNode<inputs=['ref_PCollection_PCollection_3_split'],side_inputs=[].
错误原因
这个报错是Apache Beam 2.33.0版本的已知问题,仅出现在本地测试用的DirectRunner运行器上,是MatchContinuously接口输出的元素时间戳没有被DirectRunner的水印管理器正确识别导致的断言错误,和业务代码逻辑无关。
解决方案
- 方案1:升级Apache Beam版本
该bug在2.34.0及后续稳定版本已经被官方修复,直接升级依赖即可,原有代码无需调整。 - 方案2:无法升级版本时的兼容处理
如果暂时不能升级Beam版本,可以选择两种绕过方法:
- 直接使用Dataflow运行器部署任务,该bug仅存在于本地DirectRunner,Dataflow服务运行时不会触发
- 必须本地测试时,手动给
MatchContinuously输出的元素绑定时间戳,调整后的代码如下:
import apache_beam as beam from apache_beam.utils.timestamp import Timestamp class BindTimestamp(beam.DoFn): def process(self, element): yield beam.window.TimestampedValue(element, Timestamp.now()) # 管道逻辑调整 (pipeline | 'Match Files' >> fileio.MatchContinuously(file_pattern="gs://xyz/abc/*.txt", interval=10.0, has_deduplication=True) | 'Bind Process Timestamp' >> beam.ParDo(BindTimestamp()) | 'Read Matches' >> fileio.ReadMatches() # 后续业务逻辑保持不变 )
额外说明
日志中的Python 3.8解释器警告和本次报错无关,如果你需要将任务部署到Dataflow服务,只需要保证本地开发使用的Python版本和Dataflow运行时版本一致即可避免依赖兼容性问题。
内容的提问来源于stack exchange,提问作者Aman Vaishya
相关产品推荐
相关产品推荐

