Apache Beam中BigQuery TableSchema序列化问题排查(继承方案)
问题:Apache Beam中BigQuery TableSchema序列化问题(继承方案未生效)
在Apache Beam实践中,遇到com.google.api.services.bigquery.model.TableSchema无法序列化的问题,尝试通过继承父类DoFn的方案规避,但未生效,希望保留父类设计思路,排查遗漏环节。
现有代码
父类TransformBasic
public class TransformBasic extends DoFn<String, TableRow> { private static final Logger LOG = LoggerFactory.getLogger(TransformBasic.class); protected final TupleTag<TableRow> outputTag; protected final TupleTag<TableRow> failureTag; private transient com.google.api.services.bigquery.model.TableSchema tableSchema; private String whoAmI; private String header; public TransformBasic(String whoAmI, String header, TupleTag<TableRow> outputTag, TupleTag<TableRow> failureTag) { this.whoAmI = whoAmI; this.header = header; this.outputTag = outputTag; this.failureTag = failureTag; } public TransformBasic() { this.outputTag = null; this.failureTag = null; } @ProcessElement public void processElement(ProcessContext c) { tableSchema = getTableSchemaByName(); String line = c.element(); TableRow row = new TableRow(); CSVParser csvParser = new CSVParserBuilder() .withSeparator(',') .build(); try { if (!line.equalsIgnoreCase(header)) { String[] csvColumns = csvParser.parseLine(line); for (int i = 0; i < csvColumns.length; i++) { TableFieldSchema bigqueryColumn = tableSchema.getFields().get(i); String newBigqueryColumns = replaceCharacters(bigqueryColumn.getName()); row.set(newBigqueryColumns, csvColumns[i]); } c.output(outputTag,row); } } catch (Exception e) { LOG.error("FAILURE in " + whoAmI.toUpperCase() + e); Failure failure = new Failure( LocalDate.now().toString(), whoAmI, line, e.toString()); c.output(failureTag, failure.getAsTableRow()); } } protected TableSchema getTableSchemaByName() { if ("L".equals(whoAmI)) { return getTableSchemaL(); } else if ("P".equals(whoAmI)) { return getTableSchemaP(); } else { throw new RuntimeException("Unknown class: " + whoAmI); } } }
子类TransformL
public class TransformL extends TransformBasic { public static final TupleTag<TableRow> OUTPUT_TAGS_L=new TupleTag<TableRow>() {}; public static final TupleTag<TableRow> FAILURE_TAGS_L = new TupleTag<TableRow>() {}; private static final String WHO_AM_I = "l"; private static final String HEADER_L = "metric,scen,sec,year,value,id"; public TransformL() { super(WHO_AM_I, HEADER_L,OUTPUT_TAGS_L,FAILURE_TAGS_L); } }
TableSchema定义类
public class TableSchema { public static com.google.api.services.bigquery.model.TableSchema getTableSchemaL() { List<TableFieldSchema> fields = new ArrayList<>(); fields.add(new TableFieldSchema().setName("metric")); fields.add(new TableFieldSchema().setName("scen")); fields.add(new TableFieldSchema().setName("sec")); fields.add(new TableFieldSchema().setName("year")); fields.add(new TableFieldSchema().setName("value")); fields.add(new TableFieldSchema().setName("id")); return new com.google.api.services.bigquery.model.TableSchema().setFields(fields); } }
测试代码TransformLitUwTest
public class TransformLitUwTest { private static final Logger LOG = LoggerFactory.getLogger(TransformLitUwTest.class); @Rule public final transient TestPipeline pipeline = TestPipeline.create(); public void testTransformL() { String fakeInput = "metric,scen,sec,year,value,id\n" + "metric1,scen1,10,2022,100,id1\n" + "metric2,scen2,20,2022,200,id2"; PCollection<String> inputCollection = pipeline .apply("CreateFakeInput", org.apache.beam.sdk.transforms.Create.of(fakeInput)); PCollectionTuple outputTuple = inputCollection.apply("TransformL", ParDo.of(new TransformL()) .withOutputTags(TransformLitUw.OUTPUT_TAGS_L, TupleTagList.of(TransformLitUw.FAILURE_TAGS_L))); PCollection<TableRow> outputCollection = outputTuple.get(TransformLitUw.OUTPUT_TAGS_L); outputCollection.apply(ParDo.of(new DoFn<TableRow, Void>() { @ProcessElement public void processElement(@Element TableRow c) { LOG.info("TABLE ROW -------------------------------------> " + c.toString()); } })); pipeline.run().waitUntilFinish(); } }
问题排查与修复方案
1. 大小写不匹配导致Schema获取失败
父类getTableSchemaByName中判断条件是"L".equals(whoAmI),但子类TransformL传入的WHO_AM_I是小写的"l",导致匹配失败,抛出RuntimeException("Unknown class: l")。
修复:统一大小写,推荐父类修改为不区分大小写的判断:
if ("L".equalsIgnoreCase(whoAmI)) { return TableSchema.getTableSchemaL(); }
或子类修改WHO_AM_I为大写:
private static final String WHO_AM_I = "L";
2. 未正确引用TableSchema类的静态方法
父类getTableSchemaByName中直接调用getTableSchemaL(),但父类中没有该方法,需明确调用自定义TableSchema类的静态方法TableSchema.getTableSchemaL(),同时修正方法返回值类型(原代码返回自定义TableSchema,实际需要返回BigQuery官方的com.google.api.services.bigquery.model.TableSchema)。
修复:修改父类方法:
protected com.google.api.services.bigquery.model.TableSchema getTableSchemaByName() { if ("L".equalsIgnoreCase(whoAmI)) { return TableSchema.getTableSchemaL(); } else if ("P".equalsIgnoreCase(whoAmI)) { return TableSchema.getTableSchemaP(); // 假设存在对应方法 } else { throw new RuntimeException("Unknown class: " + whoAmI); } }
3. 测试代码中TupleTag引用错误
测试代码中使用TransformLitUw.OUTPUT_TAGS_L,但子类是TransformL,应引用TransformL类的静态标签,否则会找不到对应输出分支。
修复:修改测试代码中的标签引用:
PCollectionTuple outputTuple = inputCollection.apply("TransformL", ParDo.of(new TransformL()) .withOutputTags(TransformL.OUTPUT_TAGS_L, TupleTagList.of(TransformL.FAILURE_TAGS_L))); PCollection<TableRow> outputCollection = outputTuple.get(TransformL.OUTPUT_TAGS_L);
4. 其他潜在问题
- 确保
replaceCharacters方法已正确实现,否则会影响字段名替换逻辑。 - 测试代码中
fakeInput的换行符需使用\n转义,避免编译错误。
内容的提问来源于stack exchange,提问作者Testaastopp
相关产品推荐
相关产品推荐

