You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flink SQL使用unnest函数时输出格式异常问题求助

我正在处理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"...}

问题原因

  1. json_query(countries, '$' returning array<string>)用法错误:该函数会将整个JSON数组转为包含单个字符串元素的数组(字符串内容为原JSON数组),而非将数组中的每个对象拆分为独立元素。unnest拆分后得到的是原数组的字符串形式,不是单个国家对象。
  2. 输出的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.13 14:25:16