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,threeSTRING>,旧版本Flink对ROW类型的Parquet支持度远高于MAP,1.12及以上版本基本都支持固定结构ROW的写入。 - 方案4:自定义序列化逻辑
若必须保留MAP结构且无法升级版本,可自定义Parquet序列化Schema,重写ParquetWriteSupport的类型转换逻辑适配MAP类型,该方案改造成本较高,不优先推荐。
内容的提问来源于stack exchange,提问作者success malla
相关产品推荐
相关产品推荐

