【问题标题】:Is it possible to import in Dataflow streaming pipeline written in Python the Java method `wrapBigQueryInsertError`?是否可以在用 Python 编写的 Dataflow 流式管道中导入 Java 方法 `wrapBigQueryInsertError`?
【发布时间】: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


    【解决方案1】:

    在 Beam Python 中,当执行流式 BigQuery 写入时,在 BigQuery 写入期间失败的行由转换返回。见https://github.com/apache/beam/blob/master/sdks/python/apache_beam/io/gcp/bigquery.py#L1248

    因此您可以像处理 Java 模板一样处理这些内容。

    【讨论】:

      猜你喜欢
      • 2021-05-15
      • 1970-01-01
      • 1970-01-01
      • 2020-11-12
      • 2019-05-01
      • 2010-11-24
      • 2014-01-08
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多