Unexpected type: BINARY

Viewed 38

I am trying to read parquet files via the Flink table, and it throws the error when I select one of the timestamps.

My parquet table is something like this.

enter image description here

I create a table with this SQL :

    CREATE TABLE MyDummyTable (
              `id` INT,
              ts BIGINT,
              ts_ltz AS TO_TIMESTAMP_LTZ(ts, 3),
              ts2 TIMESTAMP,
              ts3 TIMESTAMP,
              ts4 TIMESTAMP,
              ts5 TIMESTAMP
            )

It throws an error when I select one of the ts2, ts3, ts4, ts5

The error stack is this.

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

My Current appproach

I am currently using this approach to jump over the problem but it does not seem a legit solution since I create two columns in the table.

        CREATE TABLE MyDummyTable (
          `id` INT,
          ts2 STRING,
          ts2_ts AS TO_TIMESTAMP(ts2)
        )

Flink Version: 1.13.2 Scala Version: 2.11.12

0 Answers
Related