Flink SQL使用unnest函数时输出格式异常问题求助
Flink SQL拆分国家JSON数组为标准单个JSON对象的解决方案
我正在处理api.countrylayer.com返回的国家JSON列表数据,需要将列表拆分为单个国家对象做后续处理。使用Flink SQL的unnest函数实现拆分后,输出结果不符合预期:每个国家对象被方括号包裹,且是key=value的非标准格式,而预期是单个标准JSON格式的国家对象。
原始输入数据(payload)
[{"name":"Afghanistan","topLevelDomain":[".af"],"alpha2Code":"AF","alpha3Code":"AFG","callingCodes":["93"],"capital":"Kabul"...}, {"name":"Åland Islands","topLevelDomain":[".ax"],"alpha2Code":"AX","alpha3Code":"ALA","callingCodes":["358"],"capital":"Mariehamn"...}, {"name":"Albania","topLevelDomain":[".al"],"alpha2Code":"AL","alpha3Code":"ALB","callingCodes":["355"],"capital":"Tirana"...}]
当前处理代码
tableEnv.executeSql("create view countryview as" + " select " + " country " + " from kafkaCountryLayer " + " cross join unnest(json_query(countries, '$' returning array<string>)) as countrylayer (country))"); Table table = tableEnv.sqlQuery("select " + " country " + " from countryview "); table.executeInsert("print_output");
实际输出
[{name=Afghanistan, topLevelDomain=[.af], alpha2Code=AF, alpha3Code=AFG, callingCodes=[93], capital=Kabul...}] [{name=Åland Islands, topLevelDomain=[.ax], alpha2Code=AX, alpha3Code=ALA, callingCodes=[358], capital=Mariehamn...}] [{name=Albania, topLevelDomain=[.al], alpha2Code=AL, alpha3Code=ALB, callingCodes=[355], capital=Tirana...}]
预期输出
{"name":"Afghanistan","topLevelDomain":[".af"],"alpha2Code":"AF","alpha3Code":"AFG","callingCodes":["93"],"capital":"Kabul"...} {"name":"Åland Islands","topLevelDomain":[".ax"],"alpha2Code":"AX","alpha3Code":"ALA","callingCodes":["358"],"capital":"Mariehamn"...} {"name":"Albania","topLevelDomain":[".al"],"alpha2Code":"AL","alpha3Code":"ALB","callingCodes":["355"],"capital":"Tirana"...}
问题原因
json_query(countries, '$' returning array<string>)用法错误:该函数会将整个JSON数组转为包含单个字符串元素的数组(字符串内容为原JSON数组),而非将数组中的每个对象拆分为独立元素。unnest拆分后得到的是原数组的字符串形式,不是单个国家对象。- 输出的
key=value格式是Flink对复杂类型的默认toString输出,并非标准JSON格式。
解决方案
推荐使用JSON_TABLE函数直接解析JSON数组并提取对象,再用TO_JSON转换为标准JSON字符串:
修正后的代码(显式定义字段)
// 创建视图,解析JSON数组为单个国家对象,并转为标准JSON字符串 tableEnv.executeSql("CREATE VIEW countryview AS " + "SELECT TO_JSON(country_row) AS country " + "FROM kafkaCountryLayer, " + "JSON_TABLE( " + " countries, " + " '$[*]' COLUMNS( " + " name STRING PATH '$.name', " + " topLevelDomain ARRAY<STRING> PATH '$.topLevelDomain', " + " alpha2Code STRING PATH '$.alpha2Code', " + " alpha3Code STRING PATH '$.alpha3Code', " + " callingCodes ARRAY<STRING> PATH '$.callingCodes', " + " capital STRING PATH '$.capital' " + " -- 其他字段按实际需求添加 " + " ) " + "AS country_row)"); // 查询并输出 Table table = tableEnv.sqlQuery("SELECT country FROM countryview"); table.executeInsert("print_output");
简化方案(保留完整JSON对象)
如果不需要单独提取字段,直接拆分数组并保留完整对象:
tableEnv.executeSql("CREATE VIEW countryview AS " + "SELECT TO_JSON(country_obj) AS country " + "FROM kafkaCountryLayer, " + "JSON_TABLE( " + " countries, " + " '$[*]' COLUMNS( " + " country_obj OBJECT PATH '$' " + " ) " + "AS t)");
处理后,输出的每个country字段即为标准JSON字符串,完全符合预期格式。
内容的提问来源于stack exchange,提问作者Eugene
相关产品推荐
相关产品推荐

