Apachebeam、GoogleDataflow中的“finish_bundle”方法执行多次。
创始人
2024-09-05 12:30:33
0

这个问题通常是由于worker或pipeline在运行时出现异常而导致的。针对这种情况,可以使用try except块来捕捉这些异常并确保“finish_bundle”仅被执行一次。

以下是示例代码:

class MyDoFn(beam.DoFn):
    def process(self, element):
        try:
            # do some processing here
            yield BeamRecord
        except:
            # handle exception here
        finally:
            # ensure that finish_bundle is only executed once
            with self._lock:
                if self._should_finish_bundle:
                    self._should_finish_bundle = False
                    yield beam.pvalue.TaggedOutput('finished', None)
                        
    def finish_bundle(self):
        with self._lock:
            self._should_finish_bundle = True

pipeline = beam.Pipeline(options=options)
(pipeline
    | "Read input" >> beam.io.ReadFromText(input_file)
    | "Process elements" >> beam.ParDo(MyDoFn()).with_outputs('finished))
result = pipeline.run()

在这个示例中,我们定义了一个名为“MyDoFn”的DoFn类,并使用属性“self._should_finish_bundle”来跟踪一个worker是否应该执行“finish_bundle”。如果出现异常,则将其处理并不执行“finish_bundle”。在任何情况下,如果_worker_存在于_finish_bundle中,则使用一个_lock_确保它只执行一次。由于_finish_bundle是耗时操作,因此在确保它只被执行一次的同时有助于保持Beam管道的高效性。

相关内容

热门资讯

一分钟了解!人海大厅挂件怎么买... 一分钟了解!人海大厅挂件怎么买!一直一直都是有辅助技巧(竟然有挂)-哔哩哔哩小薇(辅助器软件下载)致...
第六分钟了解!游戏黑科技夫追求... 第六分钟了解!游戏黑科技夫追求!切实一直总是有辅助工具(确实有挂)-哔哩哔哩1、完成游戏黑科技夫追求...
5分钟了解!雀姬无限钻石辅助!... 5分钟了解!雀姬无限钻石辅助!果然一直都是有辅助app(有挂助手)-哔哩哔哩1、雀姬无限钻石辅助辅助...
第1分钟了解!德州来玩辅助器!... 第1分钟了解!德州来玩辅助器!一直真的是有辅助教程(有挂秘诀)-哔哩哔哩德州来玩辅助器破解侠是真的助...
十分钟了解!乐乐川南茶馆辅助!... 十分钟了解!乐乐川南茶馆辅助!总是一直总是有辅助方法(果真有挂)-哔哩哔哩1、乐乐川南茶馆辅助透视辅...
第4分钟了解!拼三张自建房软件... 第4分钟了解!拼三张自建房软件!果然真的有辅助软件(发现有挂)-哔哩哔哩1、第4分钟了解!拼三张自建...
第7分钟了解!亿游十三道脚本插... 第7分钟了解!亿游十三道脚本插件!确实存在有辅助攻略(有挂教程)-哔哩哔哩1、这是跨平台的亿游十三道...
第6分钟了解!花城牌舍辅助系统... 第6分钟了解!花城牌舍辅助系统下载!确实是真的有辅助技巧(有挂秘诀)-哔哩哔哩1、用户打开应用后不用...
第五分钟了解!决胜麻架胡易辅助... 第五分钟了解!决胜麻架胡易辅助!本来存在有辅助app(有挂技巧)-哔哩哔哩1.决胜麻架胡易辅助 选牌...
八分钟了解!奇迹脚本辅助器手机... 八分钟了解!奇迹脚本辅助器手机版!其实真的是有辅助技巧(有挂规律)-哔哩哔哩暗藏猫腻,小编详细说明奇...