如何在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
相关产品推荐
相关产品推荐

