【发布时间】:2019-03-06 07:43:41
【问题描述】:
尝试在 Cloud Dataflow Job 中启用流式传输,这需要从一个 BigQuery 表中读取数据并以附加模式将其写入另一个 BigQuery 表。
为此,我在 Java 代码中启用了options.setStreaming(true);。
应用窗口概念 - 固定窗口选项(如下代码),
PCollection<TableRow> fixedWindowedItems = finalRecords.apply(Window.<TableRow>into(FixedWindows.of(Duration.standardMinutes(1))));
最后使用 BigQueryIO(如下代码)将数据写入 BigQuery 表,
fixedWindowedItems.apply(BigQueryIO.writeTableRows()
.withSchema(schema1)
.to(options.getTargetTable())
.withMethod(BigQueryIO.Write.Method.STREAMING_INSERTS)
.withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_IF_NEEDED)
.withFailedInsertRetryPolicy(InsertRetryPolicy.alwaysRetry())
.withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_APPEND));
代码运行良好。没有错误。第一次将数据从一个表移动到另一个表。但是,如果您在第一个表中插入新数据,则第二个表不会得到反映。尽管 Job 类型为 Streaming,但 Job 似乎以 Succeeded 状态完成。
如果我在代码/配置级别错过了启用流媒体模式的内容,能否告诉我。
【问题讨论】:
-
您能解释一下保持两个 BigQuery 表同步的动机吗? - 不支持从 BigQuery 读取作为流式源。它仅用作批处理源。发生的事情是您正在从一个表读取一批数据到另一个表。 - 所以我倾向于问 agian:您为什么对不断地将数据从一个 BQ 表移动到另一个表感兴趣?
标签: google-bigquery google-cloud-dataflow apache-beam