【问题标题】:Conditional skip in google cloud-dataflow java pipelinegoogle cloud-dataflow java管道中的条件跳过
【发布时间】:2018-07-14 00:17:31
【问题描述】:

我有代码:

maxHotelEntry.apply("Convert to string", ToString.elements()).apply("Write to file", TextIO.write().to( config.getString( "gcs.checkPointLocation")).withoutSharding());

这个PCollection可以为空,-2147483648(Integer.minvalue), 124526 (+value)

如果 checkPointLocation 为空或其值小于 0,我不想写入。

一个选项是在 DoFn 中写入 GCS,但我不知道该怎么做。

【问题讨论】:

    标签: google-cloud-dataflow


    【解决方案1】:

    添加一个额外的步骤来过滤掉你不想写的元素

    maxHotelEntry
        .apply("Convert to string", ToString.elements())
        .apply("Filter out elements", ParDo
              .of(new DoFn<String, String>() {
                  public void processElement(ProcessContext c) {
                    String element = c.element();
                    if ( /*perform your filter here*/ ) {
                      c.output(element);
                    }
                  }
              }))
        .apply("Write to file", TextIO.write().to( config.getString( "gcs.checkPointLocation")).withoutSharding());
    

    【讨论】:

    • 我的问题是 maxHotelEntry 只包含一个值,这是生成的代码 - maxHotelEntry = roomsTransformedData .apply("Get max HotelId", ParDo.of(new FindMaxHotelFn())).apply(Max.整数全局());应用过滤器后,它仍然向文件写入空文件。我只想跳过文件写操作
    • 作为替代解决方案,我从 checkpointlocation 读取值并将其传递给带有侧输入的 maxHotelEntry 集合,然后比较两个值并将最大值写入 checkPointLocation
    猜你喜欢
    • 2022-11-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-01-25
    • 2020-12-10
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多