自定义Kafka Connect SMT使用Avro序列化时出现Leader not known异常
自定义Kafka Connect SMT搭配Debezium使用时的Avro序列化异常问题
我开发了一个用于特定字段清洗的自定义Kafka Connect SMT,代码结构简洁,已成功编译并添加至plugin.path。使用Debezium MySQL Connector创建如下配置的连接器时,日志中出现Avro序列化异常,核心报错为Leader not known.; error code: 50004。但将连接器配置改为使用StringConverter作为key转换器、JsonConverter作为value转换器(关闭schema)后,一切运行正常,不清楚问题所在,寻求帮助。
连接器配置
{ "name": "sql-to-kafka", "config": { "connector.class": "io.debezium.connector.mysql.MySqlConnector", "database.hostname": "127.0.0.1", "database.port": "3306", "database.user": "username", "database.password": "password", "database.server.id": "11111", "database.include.list": "test", "schema.history.internal.kafka.bootstrap.servers": "localhost:9092", "schema.history.internal.kafka.topic": "schemahistory.localdb", "include.schema.changes": "false", "database.encrypt": false, "table.include.list": "test.bins", "topic.prefix":"localdb", "transforms":"unwrap,MyCustomSMT", "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState", "transforms.unwrap.drop.tombstones": false, "transforms.unwrap.delete.handling.mode": "drop", "transforms.MyCustomSMT.type": "MyCustomSMT$Value", "transforms.MyCustomSMT.field": "segment" } }
异常日志
org.apache.kafka.connect.errors.ConnectException: Tolerance exceeded in error handler at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:220) ...(省略部分栈信息) Caused by: io.confluent.kafka.schemaregistry.client.rest.exceptions.RestClientException: Leader not known.; error code: 50004 ...(省略部分栈信息)
正常运行的转换器配置
"key.converter": "org.apache.kafka.connect.storage.StringConverter", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable":"false"
内容的提问来源于stack exchange,提问作者Karim Tawfik
相关产品推荐
相关产品推荐

