使用单消息转换(SMT)创建Kafka Redis Sink配置报错如何排查
1. 配置文件格式错误
你贴出的Redis Sink配置缺失JSON开头的{,属于基础格式错误,Kafka Connect解析配置JSON时会直接报错。
2. SMT类路径配置错误
配置中transforms.invalidaterediskeys.type的取值是你参考的insert-uuid项目的包路径,不是你自己编写的InvalidateRedisKeys类的实际全限定类名,Kafka Connect加载SMT时找不到对应类,直接报配置无效。你需要将该配置的值修改为你自己代码中InvalidateRedisKeys类的完整包名+类名。
3. 缺失Value转换器配置
你仅配置了key.converter,未配置value.converter,Kafka Connect会使用默认转换器处理消息值,你调用r.value().toString()拿到的大概率不是正常的JSON字符串,会触发后续反序列化失败。如果上游消息的value是JSON字符串,需要补充配置:
"value.converter": "org.apache.kafka.connect.storage.StringConverter"
4. POJO类缺少序列化/反序列化必要方法
你定义的类A的user_id、round_id都是私有属性,且没有对应的Getter、Setter方法,Jackson反序列化JSON时无法读写属性值,会抛出异常。你需要给类A补充对应的Getter、Setter方法,或者给属性添加public修饰符,也可以用相关注解简化开发。
5. 缺失Redis Sink删除行为配置
你使用的jcustenborder版RedisSinkConnector默认不会将null值映射为删除操作,需要补充配置开启该逻辑:
"redis.null.value.behavior": "DELETE"
6. SMT包未部署到Kafka Connect插件路径
你自己编写的SMT需要打成可被加载的jar包,放到Kafka Connect进程plugin.path配置对应的目录下,且要重启Kafka Connect进程才能识别到新的插件,否则会报类不存在的错误。
7. Jackson依赖版本不兼容
你代码中用到的DeserializationConfig.Feature.FAIL_ON_UNKNOWN_PROPERTIES属于codehaus版本的老Jackson,而Kafka Connect内置的是fasterxml版本的Jackson,依赖包不匹配会触发类/方法不存在的运行时错误,你需要修改ObjectMapper的相关配置,使用fasterxml包下的对应枚举。
内容的提问来源于stack exchange,提问作者mjennet

