【问题标题】:Execute read operations in sequence - Apache Beam按顺序执行读取操作 - Apache Beam
【发布时间】:2017-07-17 12:53:15
【问题描述】:

我需要按照给定的顺序执行以下操作:-

 PCollection<String> read = p.apply("Read Lines",TextIO.read().from(options.getInputFile())) 

      .apply("Get fileName",ParDo.of(new DoFn<String,String>(){
          ValueProvider<String> fileReceived = options.getfilename();
          @ProcessElement
          public void procesElement(ProcessContext c)
          {
              fileName = fileReceived.get().toString();
              LOG.info("File: "+fileName);
          }
      }));

      PCollection<TableRow> rows = p.apply("Read from BigQuery",
              BigQueryIO.read()
                  .fromQuery("SELECT table,schema FROM `DatasetID.TableID` WHERE file='" + fileName +"'")
              .usingStandardSql());

如何在 Apache Beam/Dataflow 中实现这一点?

【问题讨论】:

  • 您能详细介绍一下您的用例吗?这些似乎是简单的读取操作,没有任何副作用,所以我不明白为什么它们是按顺序执行还是并行执行,或者外部观察者如何甚至能够检测到其中的情况。
  • 好的...您可能已经注意到,我使用在以下查询中的“获取文件名”操作中派生的“文件名”中的变量值从 BigQuery 表中读取。但是正在发生的事情是“从 BigQuery 读取”操作发生在“获取文件名”之前,因此它得到一个空值。因此,操作必须按顺序进行。我猜上述情况正在发生,因为我在从 BigQuery 读取数据时再次使用 p.apply ...如何解决这种情况?
  • 哦,我明白了,我错过了那部分。在我解决这个问题之前 - 我对您代码中的其他内容感到困惑。您的 DoFn 始终输出相同的值,来自您的 PipelineOptions,并忽略其输入 PCollection 的内容(即您的 TextIO.read() 的结果被有效地丢弃)。这是故意的吗?
  • 是的...我的意思是我这样做只是为了访问 filaName 的值... TextIO.read() 的结果稍后在程序中使用...所以它不会被丢弃这样……
  • 如代码 sn-p 中所写,它被丢弃了 - 我认为您的实际程序是不同的。此外,在这个 sn-p 中,您不是一次输出 options.getfilename() 的值,而是输出它的 N 个副本,其中 N 是与模式“getInputFile()”匹配的所有文件中的行数 - 即 PCollection“读取" 包含许多相同的 options.getfilename() 副本。我想我可以建议如何做你真正想做的事情;将发布答案。

标签: google-cloud-dataflow apache-beam


【解决方案1】:

您似乎想将BigQueryIO.read().fromQuery() 应用于一个查询,该查询取决于通过您的ValueProvider&lt;String&gt; 中的ValueProvider&lt;String&gt; 类型的属性可用的值,并且在管道构建时无法访问提供程序 - 即您是通过模板调用您的作业。

在这种情况下,正确的解决方案是使用NestedValueProvider

PCollection<TableRow> tableRows = p.apply(BigQueryIO.read().fromQuery(
    NestedValueProvider.of(
      options.getfilename(),
      new SerializableFunction<String, String>() {
        @Override
        public String apply(String filename) {
          return "SELECT table,schema FROM `DatasetID.TableID` WHERE file='" + fileName +"'";
        }
      })));

【讨论】:

  • 如果您不介意的话,我还有另一个与之相关的问题...有没有办法使用架构和表将“读取行”中读取的文件中的数据直接插入 BigQuery从上述函数中检索。我尝试参考 Apache Beam 中的 DynamicDestinations 和您建议的帖子 - stackoverflow.com/a/43505535/278042 但不知道该怎么做。谢谢。
  • 您是否介意为此发布一个单独的问题,包括您尝试过的具体内容以及无效的内容的更多详细信息?这样,更多的人将能够看到它并从答案中学习。
  • 当然@jkff...!!
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-09-27
相关资源
最近更新 更多