如何对BigQuery中RECORD/REPEATED类型数据去标识?Dataflow构建DLP Table遇阻
解决BigQuery嵌套/重复字段转为DLP Table的问题
一、递归获取所有字段(含嵌套RECORD和REPEATED)
BigQuery的Schema是层级结构,需要递归遍历字段提取所有嵌套字段,同时保留完整路径(如user.name),示例代码如下:
private static List<Table.Field> extractAllFields(List<Field> bigQueryFields, String parentPath) { List<Table.Field> dlpFields = new ArrayList<>(); for (Field field : bigQueryFields) { String fullPath = parentPath.isEmpty() ? field.getName() : parentPath + "." + field.getName(); Table.Field.Builder dlpFieldBuilder = Table.Field.newBuilder().setName(fullPath); // 处理嵌套RECORD类型 if (field.getType() == StandardSQLTypeName.STRUCT && field.getMode() != Mode.REPEATED) { dlpFields.add(dlpFieldBuilder.build()); dlpFields.addAll(extractAllFields(field.getSubFields(), fullPath)); } // 处理REPEATED类型(含重复基本类型和重复STRUCT) else if (field.getMode() == Mode.REPEATED) { dlpFields.add(dlpFieldBuilder.build()); if (field.getType() == StandardSQLTypeName.STRUCT) { dlpFields.addAll(extractAllFields(field.getSubFields(), fullPath)); } } // 基本类型字段 else { dlpFields.add(dlpFieldBuilder.build()); } } return dlpFields; }
调用方式:
List<Table.Field> allDlpFields = extractAllFields( bigquery.getTable(table).getDefinition().getSchema().getFields(), "" );
二、递归转换BigQuery行数据为DLP Table.Row
递归处理嵌套和重复类型的字段值,将BigQuery的FieldValueList和FieldValue转为DLP的Table.Row.Cell:
private static List<Table.Row.Cell> convertBigQueryValuesToDlpCells(List<FieldValue> values, List<Field> fields, String parentPath) { List<Table.Row.Cell> cells = new ArrayList<>(); for (int i = 0; i < fields.size(); i++) { Field field = fields.get(i); FieldValue value = values.get(i); String fullPath = parentPath.isEmpty() ? field.getName() : parentPath + "." + field.getName(); // 处理NULL值 if (value.isNull()) { cells.add(Table.Row.Cell.newBuilder().setValue("").build()); continue; } // 处理REPEATED类型 if (field.getMode() == Mode.REPEATED) { List<FieldValue> repeatedValues = value.getRepeatedValue(); if (field.getType() == StandardSQLTypeName.STRUCT) { // 递归处理重复STRUCT的每个元素 for (FieldValue structVal : repeatedValues) { cells.addAll(convertBigQueryValuesToDlpCells( structVal.getRecordValue(), field.getSubFields(), fullPath )); } } else { // 重复基本类型直接拼接值 String repeatedStr = repeatedValues.stream() .map(FieldValue::getValueAsString) .collect(Collectors.joining(",")); cells.add(Table.Row.Cell.newBuilder().setValue(repeatedStr).build()); } } // 处理STRUCT类型 else if (field.getType() == StandardSQLTypeName.STRUCT) { cells.add(Table.Row.Cell.newBuilder().setValue("").build()); // 父STRUCT字段占位 cells.addAll(convertBigQueryValuesToDlpCells( value.getRecordValue(), field.getSubFields(), fullPath )); } // 基本类型 else { cells.add(Table.Row.Cell.newBuilder().setValue(value.getValueAsString()).build()); } } return cells; }
修改原有行转换逻辑:
List<Table.Row> rows = new ArrayList<>(); TableSchema schema = bigquery.getTable(table).getDefinition().getSchema(); for (FieldValueList bigQueryRowItem : bigquery.listTableData(table).getValues()) { List<Table.Row.Cell> cells = convertBigQueryValuesToDlpCells( bigQueryRowItem, schema.getFields(), "" ); Table.Row row = Table.Row.newBuilder().addAllCells(cells).build(); rows.add(row); }
三、构建完整的DLP Table对象
将提取的字段和转换后的行数据组合:
Table dlpTable = Table.newBuilder() .addAllHeaders(allDlpFields) .addAllRows(rows) .build();
内容的提问来源于stack exchange,提问作者fna
相关产品推荐
相关产品推荐

