【问题标题】:How to get failed insert record for file load insertion in BigQuery如何在 BigQuery 中获取文件加载插入失败的插入记录
【发布时间】:2021-07-18 09:47:31
【问题描述】:

我正在使用 Apache Beam (Java SDK) 通过批量加载方法(文件加载)在 BigQuery 中插入记录。我想检索那些在插入过程中失败的记录。

是否可以对失败的记录制定重试策略?

下面是我的代码:

public static void insertToBigQueryDataLake(
        final PCollectionTuple dataStoresCollectionTuple,
        final TupleTag<KV<DataLake, PayloadSpecs>> dataLakeValidTag,
        final Long loadJobTriggerFrequency,
        final Integer loadJobNumShard) {


    WriteResult writeResult = dataStoresCollectionTuple
            .get(dataLakeValidTag)
            .apply(TRANSFORMATION_NAME, DataLakeTableProcessor.dataLakeTableProcessorTransform())
            .apply(
                    WRITING_EVENTS_NAME,
                    BigQueryIO.<KV<DataLake, TableRowSpecs>>write()
                            .withMethod(BigQueryIO.Write.Method.FILE_LOADS)
                            .withTriggeringFrequency(Duration.standardMinutes(loadJobTriggerFrequency))
                            .withNumFileShards(loadJobNumShard)
                            .to(new DynamicTableRowDestinations<>(IS_DATA_LAKE))
                            .withFormatFunction(BigQueryServiceImpl::dataLakeTableRow));

    writeResult.getFailedInserts().apply(ParDo.of(new DoFn<TableRow, Void>() {
        @ProcessElement
        public void processElement(final ProcessContext processContext) throws IOException {
            System.out.println("Table Row : " + processContext.element().toPrettyString());
        }
    }));

}

【问题讨论】:

    标签: java google-bigquery apache-beam apache-beam-io spring-cloud-gcp-bigquery


    【解决方案1】:

    使用 getFailedInsertsWithErr() 方法,我们可以将失败的插入推送到另一个表以执行根本原因分析 (RCA),请查看here 了解更多详细信息。

    Example:
    // write failed rows with their error to error table                
    writeResult
            .getFailedInsertsWithErr()
            .apply(Window.into(FixedWindows.of(Duration.standardMinutes(5))))
            .apply("BQ-insert-error-extract", ParDo.of(new BigQueryInsertErrorExtractFn(tableRowToInsertView)).withSideInputs(tableRowToInsertView))
            .apply("BQ-insert-error-write", BigQueryIO.writeTableRows()
                    .to(errTableSpec)
                    .withJsonSchema(errSchema)
                    .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_IF_NEEDED)
                    .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_APPEND));
    

    【讨论】:

    • 您好兄弟,我们能够获得插入错误和记录,但只有在重试 1000 次后才能获得错误,这需要 7 个多小时。有没有办法在这个中设置重试策略。我正在加载文件(批量加载)。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-11-26
    • 2021-02-02
    • 2019-12-05
    • 1970-01-01
    • 1970-01-01
    • 2018-11-04
    • 1970-01-01
    相关资源
    最近更新 更多