【发布时间】:2019-11-22 13:37:48
【问题描述】:
我正在尝试使用 Python3 创建一个 Dataflow 流式传输管道,该管道从 Pub/Sub 主题读取消息,最终“从头开始”将它们写入 BigQuery 表。我在名为 PubSubToBigQuery.java 的 Dataflow Java 模板中看到了第三步中的一段代码,用于处理那些转换为表行的 Pub/Sub 消息,当您尝试插入时这些行会失败将它们放入 BigQuery 表中。最后,在第 4 步和第 5 步的代码片段中,将它们展平并插入错误表中:
- 第三步:
PCollection<FailsafeElement<String, String>> failedInserts =
writeResult
.getFailedInsertsWithErr()
.apply(
"WrapInsertionErrors",
MapElements.into(FAILSAFE_ELEMENT_CODER.getEncodedTypeDescriptor())
.via((BigQueryInsertError e) -> wrapBigQueryInsertError(e)))
.setCoder(FAILSAFE_ELEMENT_CODER);
- 第 4 步和第 5 步
PCollectionList.of(
ImmutableList.of(
convertedTableRows.get(UDF_DEADLETTER_OUT),
convertedTableRows.get(TRANSFORM_DEADLETTER_OUT)))
.apply("Flatten", Flatten.pCollections())
.apply(
"WriteFailedRecords",
ErrorConverters.WritePubsubMessageErrors.newBuilder()
.setErrorRecordsTable(
ValueProviderUtils.maybeUseDefaultDeadletterTable(
options.getOutputDeadletterTable(),
options.getOutputTableSpec(),
DEFAULT_DEADLETTER_TABLE_SUFFIX))
.setErrorRecordsTableSchema(ResourceUtils.getDeadletterTableSchemaJson())
.build());
failedInserts.apply(
"WriteFailedRecords",
ErrorConverters.WriteStringMessageErrors.newBuilder()
.setErrorRecordsTable(
ValueProviderUtils.maybeUseDefaultDeadletterTable(
options.getOutputDeadletterTable(),
options.getOutputTableSpec(),
DEFAULT_DEADLETTER_TABLE_SUFFIX))
.setErrorRecordsTableSchema(ResourceUtils.getDeadletterTableSchemaJson())
.build());
为了做到这一点,我怀疑使这成为可能的关键在于模板中的第一个导入库:
package com.google.cloud.teleport.templates;
import static com.google.cloud.teleport.templates.TextToBigQueryStreaming.wrapBigQueryInsertError;
此方法在 Python 中可用吗?
如果不是,有一些方法可以在 Python 中执行相同的操作,即不检查应插入的记录的字段的结构和数据类型是否符合 BigQuery 表的预期?
这种解决方法会大大减慢我的流式传输管道。
【问题讨论】:
标签: java python google-bigquery google-cloud-dataflow