创建Elasticsearch Sink Connector时遭遇序列化错误的技术求助
解决Elasticsearch Sink Connector的反序列化错误
先给你拆解下这个报错的核心问题:从堆栈信息里能明确看到,Kafka Connect一直在尝试用JSON反序列化器去解析你topic里的AVRO格式消息,这肯定会失败!你配置里写了"input.data.format"= 'AVRO',但这个参数根本不是Elasticsearch Sink Connector的标准配置项,连接器完全不会识别它,所以还是默认用了JsonConverter去解析,自然就抛出序列化异常了。
直接修改方案
把你的Connector创建语句改成下面这样,重点添加AVRO转换器的配置:
CREATE SOURCE CONNECTOR `testconnector` WITH( "type.name"= '_doc', "connector.class"= 'io.confluent.connect.elasticsearch.ElasticsearchSinkConnector', "tasks.max"= '1', "transforms"= 'Dealership', "topics"= 'es.contact.model', "transforms.Dealership.type"= 'io.confluent.connect.transforms.ExtractTopic$Value', "transforms.Dealership.field"= 'indexTopicName', "transforms.Dealership.skip.missing.or.null"= 'true', "connection.url"= 'http://192.168.1.5:19200', "key.ignore"= 'true', "schema.ignore"= 'true', -- 新增AVRO转换器配置,替换默认的JsonConverter "value.converter"= 'io.confluent.connect.avro.AvroConverter', "value.converter.schema.registry.url"= 'http://<你的Schema Registry地址>:8081' );
关键修改说明
- 删掉无效的
input.data.format:这个参数是给部分特定Connector用的,Elasticsearch Sink完全不认,留着没用。 - 指定Avro转换器:
value.converter设置成io.confluent.connect.avro.AvroConverter,明确告诉连接器用AVRO的规则去解析消息内容。 - 配置Schema Registry地址:AVRO格式依赖Schema Registry来获取消息的结构信息,所以必须填写你实际的Schema Registry地址和端口(默认端口是8081)。
额外要检查的点
- 确认你的Schema Registry服务处于正常运行状态,并且Connector所在的机器能正常访问这个服务地址。
- 确认
es.contact.model这个topic里的消息确实是用AVRO序列化的,而且Schema Registry中已经存在对应的消息schema。 - 如果你们没用到Schema Registry,而是采用无schema的AVRO编码,那需要调整转换器的配置,或者改成带schema的AVRO序列化方式。
内容的提问来源于stack exchange,提问作者bala n
相关产品推荐
相关产品推荐

