Flink读取Parquet文件时Timestamp字段报Unexpected type: BINARY错误
Flink Table读取Parquet Timestamp字段报错问题
问题场景
使用Flink Table读取Parquet文件时,查询ts2、ts3、ts4、ts5这些定义为Timestamp类型的字段时触发错误。
Parquet表结构
实际Parquet文件中,上述时间字段的存储类型为BINARY。
Flink建表SQL
CREATE TABLE MyDummyTable ( `id` INT, ts BIGINT, ts_ltz AS TO_TIMESTAMP_LTZ(ts, 3), ts2 TIMESTAMP, ts3 TIMESTAMP, ts4 TIMESTAMP, ts5 TIMESTAMP )
错误堆栈
Caused by: java.lang.IllegalArgumentException: Unexpected type: BINARY at org.apache.parquet.Preconditions.checkArgument(Preconditions.java:77) at org.apache.flink.formats.parquet.vector.ParquetSplitReaderUtil.createWritableColumnVector(ParquetSplitReaderUtil.java:369) at org.apache.flink.formats.parquet.ParquetVectorizedInputFormat.createWritableVectors(ParquetVectorizedInputFormat.java:264) at org.apache.flink.formats.parquet.ParquetVectorizedInputFormat.createReaderBatch(ParquetVectorizedInputFormat.java:254) at org.apache.flink.formats.parquet.ParquetVectorizedInputFormat.createPoolOfBatches(ParquetVectorizedInputFormat.java:244) at org.apache.flink.formats.parquet.ParquetVectorizedInputFormat.createReader(ParquetVectorizedInputFormat.java:137) at org.apache.flink.formats.parquet.ParquetVectorizedInputFormat.createReader(ParquetVectorizedInputFormat.java:73) at org.apache.flink.connector.file.src.impl.FileSourceSplitReader.checkSplitOrStartNext(FileSourceSplitReader.java:112) at org.apache.flink.connector.file.src.impl.FileSourceSplitReader.fetch(FileSourceSplitReader.java:65) at org.apache.flink.connector.base.source.reader.fetcher.FetchTask.run(FetchTask.java:56) at org.apache.flink.connector.base.source.reader.fetcher.SplitFetcher.runOnce(SplitFetcher.java:138) ... 7 more
临时解决方案
通过将字段定义为STRING类型再转换的方式绕过,但需要额外定义转换字段,并非最优方案:
CREATE TABLE MyDummyTable ( `id` INT, ts2 STRING, ts2_ts AS TO_TIMESTAMP(ts2) )
版本信息
- Flink 1.13.2
- Scala 2.11.12
内容的提问来源于stack exchange,提问作者None
相关产品推荐
相关产品推荐

