如何将JSON Schema转换为Flink SQL DDL语句?
如何将JSON Schema转换为Flink SQL DDL语句?
给定包含嵌套对象的JSON Schema:
{ "type": "object", "properties": { "name": { "type": "string" }, "id": { "type": "string" }, "i": { "type": "integer" }, "p": { "type": "object", "properties": { "name": { "type": "string" }, "id": { "type": "string" }, "i": { "type": "integer" }, "p1": { "type": "object", "properties": { "name": { "type": "string" }, "id": { "type": "string" }, "i": { "type": "integer" } } } } } } }
需要将其转换为使用ROW类型的Flink SQL DDL,示例目标结果如下:
CREATE TABLE test_nested ( name STRING, id STRING, i INTEGER, p ROW(name STRING, id STRING, i INTEGER, p1 ROW(name STRING, id STRING, i INTEGER)) ) WITH ( 'connector'='print' );
我尝试遍历properties字段,但遇到object类型时不知道如何处理,猜测需要递归但不清楚具体步骤,现有代码如下:
private static String traverse(JsonNode rootNode) { JsonNode properties = rootNode.get("properties"); Iterator<Map.Entry<String, JsonNode>> fields = properties.fields(); // columnNameType will contain the "columnName dataType", for example // "name STRING" // "id STRING" // "i INTEGER" // "p ROW(name STRING, id STRING, i INTEGER, p1 ROW(name STRING, id STRING, i INTEGER))" List<String> columnNameType = new ArrayList<>(); while (fields.hasNext()) { Map.Entry<String, JsonNode> next = fields.next(); String key = next.getKey(); JsonNode value = next.getValue(); String type = value.get("type").asText(); if(type.equals("object")) { // recurse??? } else { columnNameType.add(key+" "+type); } } // this should return a formatted string like /* name STRING, id STRING, i INTEGER, p ROW(name STRING, id STRING, i INTEGER, p1 ROW(name STRING, id STRING, i INTEGER)) */ return String.join(",\n", columnNameType); }
解决思路:递归处理嵌套对象
遇到object类型时,递归调用traverse方法处理该对象的properties,将结果包裹在ROW(...)中后和字段名组合,加入字段列表。修改后的代码如下:
private static String traverse(JsonNode rootNode) { JsonNode properties = rootNode.get("properties"); if (properties == null || !properties.isObject()) { return ""; } Iterator<Map.Entry<String, JsonNode>> fields = properties.fields(); List<String> columnNameType = new ArrayList<>(); while (fields.hasNext()) { Map.Entry<String, JsonNode> next = fields.next(); String key = next.getKey(); JsonNode value = next.getValue(); String type = value.get("type").asText(); if ("object".equals(type)) { // 递归处理嵌套对象,生成ROW类型的内容 String nestedFields = traverse(value); columnNameType.add(key + " ROW(" + nestedFields + ")"); } else { // 基础类型直接拼接字段名和大写类型(符合Flink SQL语法) columnNameType.add(key + " " + type.toUpperCase()); } } return String.join(", ", columnNameType); }
完整DDL生成示例
可以封装一个方法生成完整的CREATE TABLE语句,保证格式可读性:
public static String generateFlinkDDL(String tableName, JsonNode jsonSchema) { String columns = traverse(jsonSchema); // 替换分隔符,添加换行和缩进 String formattedColumns = columns.replace(", ", ",\n "); return String.format("CREATE TABLE %s (\n %s\n) WITH (\n 'connector'='print'\n);", tableName, formattedColumns); }
调用generateFlinkDDL("test_nested", jsonSchemaNode)即可得到目标DDL语句。
代码说明
- 递归逻辑:字段类型为
object时,递归解析其内部属性,将结果嵌入ROW()结构,实现嵌套类型的转换。 - 类型规范:将JSON Schema的小写类型转为Flink SQL要求的大写形式,避免语法错误。
- 格式优化:生成完整DDL时调整字段分隔符,添加换行和缩进,让最终SQL语句更易读。
内容的提问来源于stack exchange,提问作者slowratatoskr
相关产品推荐
相关产品推荐

