如何在Dataflow作业Pipeline执行完成后运行指定代码逻辑
根因说明
你遇到的并行执行问题本质是对Dataflow Pipeline的执行机制理解有偏差:直接调用pipeline.run()时,方法只会完成作业定义到Dataflow服务端的异步提交,提交动作完成后当前线程就会继续向下执行,不会阻塞等待服务端的作业实际运行完成。这时候不管你把后置逻辑写在finally块还是try块的后续代码里,都会和服务端正在运行的Pipeline作业并行,自然无法保证执行顺序。
可行实现方案
批作业场景(作业会运行到自然结束):主动阻塞等待作业终态后再执行后置逻辑
调用pipeline.run()拿到PipelineResult实例后,调用它的waitUntilFinish()方法,该方法会持续阻塞当前线程,直到服务端的Pipeline作业完全进入终态(成功/失败/被取消)才会返回最终的作业状态,此时再执行数据表更新逻辑,就能保证一定是在Pipeline全部执行完成后触发。
参考实现代码:Pipeline pipeline = Pipeline.create(options); // 此处编写Pipeline的全部数据处理逻辑,例如各类apply转换 PipelineResult runResult = pipeline.run(); // 阻塞等待作业执行完成,获取最终状态 State jobFinalState = runResult.waitUntilFinish(); // 代码走到此处时,Pipeline已经完全执行结束 if (jobFinalState == State.DONE) { // 作业运行成功,在此处编写两张数据表的更新逻辑 } else { // 作业失败/被取消,按需编写异常处理逻辑 }注意:不要把数据表更新逻辑放在finally代码块中,直接写在
waitUntilFinish()返回后的代码段即可,避免作业提交阶段就触发异常时误执行更新操作。流作业场景(作业长期运行无自动终态):如果你的Pipeline是流式作业,不存在自然执行结束的节点,就不能用客户端阻塞等待的方案,可以改用两种方式实现:
- 配置Dataflow作业状态通知,将作业状态变更事件推送到Pub/Sub,消费到作业终止/达到指定处理进度的事件时再触发表更新
- 在Pipeline内部通过全局窗口+收尾触发器、
Wait变换等内置组件,在数据处理完成的节点触发外部表更新操作
内容的提问来源于stack exchange,提问作者raj
相关产品推荐
相关产品推荐

