You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Apache Beam Java SDK中查找上传BigQuery时被拒的文件及记录

Apache Beam Java SDK写入BigQuery时定位失败记录及对应源文件的实现方案

针对你使用通配符路径C:\\users\\*\\*批量读取文件写入BigQuery的场景,按以下步骤实现即可准确拿到被BQ拒绝的记录内容、所属源文件名及错误原因:

核心实现逻辑

  • 读取阶段必须给每条记录绑定所属源文件的元信息,不要直接使用无上下文的批量读取API
  • 写入BQ时开启扩展错误信息收集,不要让管道遇到错误直接终止
  • 从写入返回的失败结果集中解析溯源信息,按需求输出到日志/错误表/错误文件即可

具体实现步骤

1. 带文件元信息读取通配符路径下的所有文件

不要直接用TextIO.read()读取通配符路径,这种方式会丢失文件来源信息,改用FileIO匹配路径后逐文件读取,把文件名和每条记录绑定后再往下游传输:

// 读取后返回KV结构:key是源文件全路径,value是单条原始记录内容
PCollection<KV<String, String>> fileBoundRecords = p
  .apply(FileIO.match().filepattern("C:\\users\\*\\*"))
  .apply(FileIO.readMatches())
  .apply(ParDo.of(new DoFn<ReadableFile, KV<String, String>>() {
    @ProcessElement
    public void process(@Element ReadableFile readableFile, OutputReceiver<KV<String, String>> out) throws IOException {
      String sourcePath = readableFile.getMetadata().resourceId().toString();
      // 按自身文件格式替换解析逻辑,此处以普通逐行文本读取为例
      try (BufferedReader br = new BufferedReader(
              Channels.newReader(readableFile.open(), StandardCharsets.UTF_8))) {
        String line;
        while ((line = br.readLine()) != null) {
          // 空行过滤逻辑可自行添加
          out.output(KV.of(sourcePath, line));
        }
      }
    }
  }));

2. 转换BQ行时保留溯源字段

把原始记录转换成BQ要求的TableRow结构时,新增临时字段存储源文件路径,建议同时存一份原始记录内容方便排查:

PCollection<TableRow> bqTableRows = fileBoundRecords.apply(ParDo.of(new DoFn<KV<String, String>, TableRow>() {
  @ProcessElement
  public void process(@Element KV<String, String> fileAndRecord, OutputReceiver<TableRow> out) {
    String sourceFile = fileAndRecord.getKey();
    String rawContent = fileAndRecord.getValue();
    TableRow row = new TableRow();
    // 写入溯源信息
    row.set("source_file_path", sourceFile);
    row.set("raw_record_content", rawContent);
    // 此处替换为自身业务的字段映射逻辑,把rawContent解析成BQ表对应的业务字段
    // row.set("biz_field1", xxx);
    out.output(row);
  }
}));

3. 配置BQ写入参数,开启失败收集

写入BQ时必须开启扩展错误信息开关,配置重试策略处理临时网络类错误,永久性错误(字段不匹配、格式错误、长度超限等)会自动进入失败集合:

WriteResult bqWriteResult = bqTableRows.apply(BigQueryIO.writeTableRows()
  .to("你的项目ID:你的数据集名.你的目标表名")
  .withSchema(预先定义好的BQ表Schema对象)
  .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_IF_NEEDED)
  .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_APPEND)
  // 仅重试临时网络类、服务端过载类错误
  .withFailedInsertRetryPolicy(InsertRetryPolicy.retryTransientErrors())
  // 必须开启该配置,否则拿不到完整错误上下文
  .withExtendedErrorInfo()
);

4. 解析失败记录输出排查信息

从写入结果中拿到失败插入的错误集合,即可提取出源文件路径、原始记录、BQ返回的具体错误原因:

PCollection<BigQueryInsertError> insertErrors = bqWriteResult.getFailedInsertsWithErr();
insertErrors.apply(ParDo.of(new DoFn<BigQueryInsertError, Void>() {
  @ProcessElement
  public void process(@Element BigQueryInsertError err) {
    TableRow failedRow = err.getRow();
    String sourceFile = (String) failedRow.get("source_file_path");
    String rawRecord = (String) failedRow.get("raw_record_content");
    String errorReason = err.getError().toString();
    // 自定义错误处理逻辑:打印日志、写入专门的BQ错误表、写入本地错误文件均可
    System.out.printf("BQ写入拒绝 | 源文件:%s | 错误原因:%s | 失败记录:%s%n",
            sourceFile, errorReason, rawRecord);
  }
}));

注意事项

  • 如果你使用FILE_LOADS模式(批量加载)而非STREAMING_INSERTS模式写入BQ,失败记录会输出到你配置的临时路径下,错误文件内同样会携带提前写入的source_file_path和raw_record_content字段,直接读取解析即可。
  • 如果使用BigQuery Storage Write API写入,逻辑完全一致:提前绑定文件元信息,从WriteResult中获取失败集合解析即可。
  • 不要在转换TableRow的步骤中删除溯源字段,否则拿到失败记录后也无法定位来源。
  • 临时类错误(比如配额超限、网络波动)会被自动重试,不会进入失败集合,只有重试多次仍失败的永久性错误才会被收集。

内容的提问来源于stack exchange,提问作者raj

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.03 03:09:42