You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用单消息转换(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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.26 14:36:01