flink-issues mailing list archives

Site index · List index
Message view « Date » · « Thread »
Top « Date » · « Thread »
From "Jun Zhang (JIRA)" <j...@apache.org>
Subject [jira] [Updated] (FLINK-9444) KafkaAvroTableSource failed to work for map fields
Date Sat, 26 May 2018 13:18:00 GMT

     [ https://issues.apache.org/jira/browse/FLINK-9444?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]

Jun Zhang updated FLINK-9444:
-----------------------------
    Description: 
When some Avro schema has map fields and the corresponding TableSchema declares *MapTypeInfo* for
these fields, an exception will be thrown when registering the *KafkaAvroTableSource*, complaining
like:

Exception in thread "main" org.apache.flink.table.api.ValidationException: Type Map<String,
Integer> of table field 'event' does not match with type GenericType<java.util.Map>
of the field 'event' of the TableSource return type.
 at org.apache.flink.table.api.ValidationException$.apply(exceptions.scala:74)
 at org.apache.flink.table.sources.TableSourceUtil$$anonfun$validateTableSource$1.apply(TableSourceUtil.scala:92)
 at org.apache.flink.table.sources.TableSourceUtil$$anonfun$validateTableSource$1.apply(TableSourceUtil.scala:71)
 at scala.collection.IndexedSeqOptimized$class.foreach(IndexedSeqOptimized.scala:33)
 at scala.collection.mutable.ArrayOps$ofRef.foreach(ArrayOps.scala:186)
 at org.apache.flink.table.sources.TableSourceUtil$.validateTableSource(TableSourceUtil.scala:71)
 at org.apache.flink.table.plan.schema.StreamTableSourceTable.<init>(StreamTableSourceTable.scala:33)
 at org.apache.flink.table.api.StreamTableEnvironment.registerTableSourceInternal(StreamTableEnvironment.scala:124)
 at org.apache.flink.table.api.TableEnvironment.registerTableSource(TableEnvironment.scala:438)

  was:
When some Avro schema has map fields and the corresponding TableSchema declares *MapTypeInfo* for
these fields, an exception will be thrown when registering the *KafkaAvroTableSource*, complaining
like:

Exception in thread "main" org.apache.flink.table.api.ValidationException: Type Map<String,
String> of table field 'event' does not match with type GenericType<java.util.Map>
of the field 'event' of the TableSource return type.
 at org.apache.flink.table.api.ValidationException$.apply(exceptions.scala:74)
 at org.apache.flink.table.sources.TableSourceUtil$$anonfun$validateTableSource$1.apply(TableSourceUtil.scala:92)
 at org.apache.flink.table.sources.TableSourceUtil$$anonfun$validateTableSource$1.apply(TableSourceUtil.scala:71)
 at scala.collection.IndexedSeqOptimized$class.foreach(IndexedSeqOptimized.scala:33)
 at scala.collection.mutable.ArrayOps$ofRef.foreach(ArrayOps.scala:186)
 at org.apache.flink.table.sources.TableSourceUtil$.validateTableSource(TableSourceUtil.scala:71)
 at org.apache.flink.table.plan.schema.StreamTableSourceTable.<init>(StreamTableSourceTable.scala:33)
 at org.apache.flink.table.api.StreamTableEnvironment.registerTableSourceInternal(StreamTableEnvironment.scala:124)
 at org.apache.flink.table.api.TableEnvironment.registerTableSource(TableEnvironment.scala:438)


> KafkaAvroTableSource failed to work for map fields
> --------------------------------------------------
>
>                 Key: FLINK-9444
>                 URL: https://issues.apache.org/jira/browse/FLINK-9444
>             Project: Flink
>          Issue Type: Bug
>          Components: Table API &amp; SQL
>    Affects Versions: 1.6.0
>            Reporter: Jun Zhang
>            Priority: Blocker
>             Fix For: 1.6.0
>
>         Attachments: flink-9444
>
>
> When some Avro schema has map fields and the corresponding TableSchema declares *MapTypeInfo* for
these fields, an exception will be thrown when registering the *KafkaAvroTableSource*, complaining
like:
> Exception in thread "main" org.apache.flink.table.api.ValidationException: Type Map<String,
Integer> of table field 'event' does not match with type GenericType<java.util.Map>
of the field 'event' of the TableSource return type.
>  at org.apache.flink.table.api.ValidationException$.apply(exceptions.scala:74)
>  at org.apache.flink.table.sources.TableSourceUtil$$anonfun$validateTableSource$1.apply(TableSourceUtil.scala:92)
>  at org.apache.flink.table.sources.TableSourceUtil$$anonfun$validateTableSource$1.apply(TableSourceUtil.scala:71)
>  at scala.collection.IndexedSeqOptimized$class.foreach(IndexedSeqOptimized.scala:33)
>  at scala.collection.mutable.ArrayOps$ofRef.foreach(ArrayOps.scala:186)
>  at org.apache.flink.table.sources.TableSourceUtil$.validateTableSource(TableSourceUtil.scala:71)
>  at org.apache.flink.table.plan.schema.StreamTableSourceTable.<init>(StreamTableSourceTable.scala:33)
>  at org.apache.flink.table.api.StreamTableEnvironment.registerTableSourceInternal(StreamTableEnvironment.scala:124)
>  at org.apache.flink.table.api.TableEnvironment.registerTableSource(TableEnvironment.scala:438)



--
This message was sent by Atlassian JIRA
(v7.6.3#76005)

Mime
View raw message