Apache Beam多表读取管道:条件跳步与步骤依赖配置问题
问题解答
问题1:步骤1无数据时跳过步骤2
要实现步骤1无返回数据时跳过步骤2的逻辑,你可以借助Beam的If变换结合PCollectionView来完成,核心是先判断步骤1的结果是否为空,再决定是否执行步骤2的后续处理。
首先把步骤1的结果转换成可作为侧输入的PCollectionView:
// 将步骤1的结果转为View,用于后续判断和步骤2的侧输入 PCollectionView<List<Trigger>> triggerList = readyOrInProgressTriggers .apply(View.asList());
然后修改步骤2的代码,用If变换做条件判断:
// 仅当triggerList不为空时,才执行步骤2的过滤、转换逻辑 PCollectionView<List<Event>> srcAggrEvents = pipeline .apply(BigtableHelper.getBigtableIORead("event")) .apply(If.of(Input.isNotEmpty(triggerList)) .then( p -> p.apply(ParDo.of(new FilterSrcAggrEventRowsFn()) .withSideInput("triggerList", triggerList)) .apply(Event.bigTableRowToPojo()) .apply(View.asList())) .otherwise( p -> p.apply(Create.empty(Event.class)) .apply(View.asList())));
如果步骤1没有返回数据,triggerList为空,步骤2会直接生成一个空的PCollectionView,跳过实际的业务处理逻辑。
问题2:让步骤3依赖步骤2完成
当前步骤3和步骤2是并行执行的,要建立依赖关系,只需要让步骤3的处理逻辑关联步骤2的输出即可——哪怕只是形式上的依赖,Beam执行引擎会自动保证步骤2完成后再启动步骤3。
最简单的方式是在步骤3的ParDo中添加步骤2的输出作为侧输入(不需要实际使用,仅用于建立依赖):
// 先保留步骤2的原始PCollection引用(如果之前转成View,也可以直接用View) PCollection<Event> eventPCollection = pipeline .apply(BigtableHelper.getBigtableIORead("event")) .apply(If.of(Input.isNotEmpty(triggerList)) .then( p -> p.apply(ParDo.of(new FilterSrcAggrEventRowsFn()) .withSideInput("triggerList", triggerList)) .apply(Event.bigTableRowToPojo())) .otherwise(p -> p.apply(Create.empty(Event.class)))); // 转换为View(如果业务需要) PCollectionView<List<Event>> srcAggrEvents = eventPCollection.apply(View.asList()); // 步骤3:添加步骤2的侧输入,建立依赖关系 readyOrInProgressTriggers .apply("update status", ParDo.of( new TriggerStatus.updateTriggerStatus(BNCConstant.COMPLETE_STATUS)) .withSideInput("dummy_dependency", srcAggrEvents)) // 仅用于建立依赖,无需在DoFn中使用 .apply(ParDo.of(new TriggerStatus.pojoToMutation())) .apply(BigtableHelper.writeToBigtable("trgr_sta"));
也可以通过添加一个空的ParDo变换来明确建立依赖:
readyOrInProgressTriggers .apply("wait_for_step2", ParDo.of(new DoFn<Trigger, Trigger>() { @ProcessElement public void processElement(ProcessContext c) { c.output(c.element()); // 原样输出,仅用于绑定依赖 } }).withSideInput("dummy", srcAggrEvents)) .apply("update status", ParDo.of( new TriggerStatus.updateTriggerStatus(BNCConstant.COMPLETE_STATUS))) .apply(ParDo.of(new TriggerStatus.pojoToMutation())) .apply(BigtableHelper.writeToBigtable("trgr_sta"));
内容的提问来源于stack exchange,提问作者Learner
相关产品推荐
相关产品推荐

