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

Flink filesystem sink存parquet时嵌套MAP类型数据报错

问题描述

尝试将JSON数据转换为Parquet格式,供Trino或Presto查询使用。
示例JSON数据:

{"name": "success","message": "test","id": 1, "test1":  {"one": 1, "two":  2, "three":  "t3"}, "test2":  [1,2,3], "test3": [{"a": "a"},{"a": "aa"}], "test4": [{"a": "a"},{"a": "aa"}]}

现有Flink SQL代码:

// 建JSON源表
tEnv.executeSql("create TEMPORARY table test (" +
            "name string," +
            "message string," +
            "id int," +
            "test1 map<string,string>," +
            "test2 array< int >," +
            "test3 array< map<string,string>>," +
            "test4 string" +
            ")" +
            "with (" +
            "'connector' = 'filesystem'," +
            "'path' = 'file:///Users/successmalla/big_data/flink/src/main/resources/test.json'," +
            "'format' = 'json'" +
            ")");

// 建Parquet目标表
tEnv.executeSql("create table test2 (" +
            "name string," +
            "message string," +
            "id int," +
            "test1 map<string,string>," +
            "test4 string" +
            ")" +
            "with (" +
            "'connector' = 'filesystem'," +
            "'path' = 'file:///Users/successmalla/big_data/flink/src/main/resources/testresilt'," +
            "'format' = 'parquet'" +
            ")");

// 写入数据
tEnv.executeSql("insert into test2 " +
            "select name, message, id, test1, test4 " +
            "from test ");

运行时抛出的异常:

Caused by: java.lang.UnsupportedOperationException: Unsupported type: MAP<STRING, STRING>
at org.apache.flink.formats.parquet.utils.ParquetSchemaConverter.convertToParquetType(ParquetSchemaConverter.java:105)
at org.apache.flink.formats.parquet.utils.ParquetSchemaConverter.convertToParquetType(ParquetSchemaConverter.java:43)
at org.apache.flink.formats.parquet.utils.ParquetSchemaConverter.convertToParquetMessageType(ParquetSchemaConverter.java:37)
at org.apache.flink.formats.parquet.row.ParquetRowDataBuilder$ParquetWriteSupport.<init>(ParquetRowDataBuilder.java:72)
at org.apache.flink.formats.parquet.row.ParquetRowDataBuilder$ParquetWriteSupport.<init>(ParquetRowDataBuilder.java:70)
at org.apache.flink.formats.parquet.row.ParquetRowDataBuilder.getWriteSupport(ParquetRowDataBuilder.java:67)
at org.apache.parquet.hadoop.ParquetWriter$Builder.build(ParquetWriter.java:652)
at org.apache.flink.formats.parquet.row.ParquetRowDataBuilder$FlinkParquetBuilder.createWriter(ParquetRowDataBuilder.java:135)
at org.apache.flink.formats.parquet.ParquetWriterFactory.create(ParquetWriterFactory.java:56)

当前map、array、row类型都可以正常查询打印,但无法写入Parquet格式。


报错原因

当前使用的Flink版本较低,旧版本Flink的Parquet序列化组件未实现MAP、嵌套ARRAY等复合类型的转换逻辑,因此读取JSON时可以正常解析复合类型做计算,但写入Parquet时会抛出不支持类型的异常。


解决方案
  • 方案1(优先推荐):升级Flink版本
    Flink 1.14及以上版本已经完善了Parquet对复合类型的支持,升级到1.14+版本后,原生支持MAP、ROW、嵌套ARRAY类型的写入,现有代码无需修改即可正常运行。
  • 方案2:转换为JSON字符串存储
    如果无法升级Flink版本,可将目标表中的复合类型字段定义为STRING类型,写入时用JSON_FORMAT()函数将map/row转成JSON字符串存储,后续Trino/Presto查询时再用JSON_EXTRACT类函数解析即可,适合查询时不需要高频访问嵌套字段的场景。
  • 方案3:替换MAP为固定结构ROW类型
    若嵌套字段结构固定,可将源表和目标表的对应字段定义为明确结构的ROW类型,例如示例中的test1可定义为ROW<oneINT,twoINT,three STRING>,旧版本Flink对ROW类型的Parquet支持度远高于MAP,1.12及以上版本基本都支持固定结构ROW的写入。
  • 方案4:自定义序列化逻辑
    若必须保留MAP结构且无法升级版本,可自定义Parquet序列化Schema,重写ParquetWriteSupport的类型转换逻辑适配MAP类型,该方案改造成本较高,不优先推荐。

内容的提问来源于stack exchange,提问作者success malla

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 19:09:01