ElasticsearchSinkConnector Avro反序列化失败问题求助
解决Kafka Elasticsearch Sink Connector的Avro反序列化失败问题
这个问题我之前碰到过,核心原因是你的Connector默认使用了Avro消息转换器,但你的resd topic里存的是纯JSON格式消息,两者不匹配导致反序列化失败。下面是具体的解决方法:
问题根源拆解
Confluent Connect默认会用io.confluent.connect.avro.AvroConverter作为消息转换器,它会尝试把消息解析成Avro格式。但你的topic里是无Schema的JSON数据,Avro转换器自然无法识别,就会抛出反序列化错误。
修改Connector配置,指定JSON转换器
你需要在配置中显式指定JSON转换器,同时关闭Schema校验(和你已设置的schema.ignore=true匹配)。修改后的完整配置如下:
{ "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector", "type.name": "test-type", "tasks.max": "1", "topics": "resd", // 注意:你原配置写的是dialogs,实际使用的topic是resd,这里要保持一致 "name": "elasticsearch-sink", "key.ignore": "true", "connection.url": "http://localhost:9200", "schema.ignore": "true", "key.converter": "org.apache.kafka.connect.json.JsonConverter", "key.converter.schemas.enable": "false", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "false" }
关键配置说明
key.converter/value.converter:指定用JSON转换器处理消息的Key和Value,替代默认的Avro转换器key.converter.schemas.enable/value.converter.schemas.enable:设为false,因为你的JSON消息没有附带Schema信息,和schema.ignore=true的配置逻辑一致- 修正
topics字段:确保监听的是你实际存储JSON消息的resdtopic,避免Connector监听错误的数据源
验证修复
重新部署修改后的Connector,它会用JSON转换器解析resd topic里的消息,就不会再出现Avro反序列化失败的问题了。如果你的Connect Worker全局配置默认是Avro转换器也没关系,Connector级别的配置会覆盖全局设置。
内容的提问来源于stack exchange,提问作者Mikhail
相关产品推荐
相关产品推荐

