如何在Flink SQL中解决Kafka Avro枚举类型反序列化报错问题
解决方案
你遇到的报错由字段类型定义错误、缺少枚举映射配置两个问题导致,按以下步骤修改即可:
- 修正字段类型定义
Avro枚举类型不需要用ROW嵌套定义,直接映射为Flink SQL的STRING类型即可。 - 新增枚举转字符串配置
Flink 1.13版本的avro-confluent格式默认不会自动将Avro枚举转换为字符串,需要在WITH参数中显式开启value.avro-confluent.enum-as-string配置。
正确的DDL如下:
CREATE TABLE my_table ( `entry` ROW( `myEntryType` STRING ) ) WITH ( 'connector' = 'kafka', 'topic' = 'my-topic', 'properties.bootstrap.servers' = '...:9092', 'scan.startup.mode' = 'latest-offset', 'value.format' = 'avro-confluent', 'value.avro-confluent.schema-registry.url' = 'http://...:8081', 'value.avro-confluent.enum-as-string' = 'true' );
问题原因说明
- 原报错
Found my.type.avro.MyEnumType, expecting union是因为你错误将枚举类型的myEntryType声明为ROW类型,导致Flink生成的期望Avro Schema和Schema Registry中注册的实际Schema不匹配,反序列化时类型校验失败。 - 尝试RAW类型失败的原因是:avro-confluent格式默认反序列化枚举时生成的是Avro自带的
GenericEnumSymbol实例,和你指定的自定义生成类my.type.avro.MyEnumType不兼容,不存在匹配的序列化器,无需使用RAW类型。
内容的提问来源于stack exchange,提问作者Kiran Ashraf
相关产品推荐
相关产品推荐

