【发布时间】: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