Couchbase Kafka Sink连接器遇空键报错停滞,如何忽略错误继续运行?
解决方案
1. 调整Couchbase Sink文档ID生成策略
默认Couchbase Sink会将Kafka消息Key作为文档ID,空Key会触发Id cannot be null or empty错误。可通过couchbase.document.id参数指定其他生成逻辑,摆脱对Key的依赖:
- 生成随机UUID作为文档ID:
"couchbase.document.id": "${random.uuid()}" - 使用消息Value中的字段作为文档ID(假设Value为JSON格式,包含
id字段):"couchbase.document.id": "${value.id}"
2. 用Kafka Connect Transform过滤空Key消息
若无需保留空Key消息,可添加Filter转换规则直接丢弃此类消息,避免触发错误:
"transforms": "FilterNullKeys", "transforms.FilterNullKeys.type": "org.apache.kafka.connect.transforms.Filter$Value", "transforms.FilterNullKeys.filter.condition": "$key == null || $key == ''", "transforms.FilterNullKeys.filter.type": "exclude"
3. 升级Couchbase Sink连接器版本
部分旧版Couchbase Sink未适配Kafka Connect的错误容忍机制,即便配置errors.tolerance: all也无法捕获异常。建议升级至最新稳定版(至少2.0.0以上),确保连接器将空Key异常标记为可恢复,让错误容忍参数生效。
4. 完善错误容忍参数配置
确保连接器配置包含以下参数,无拼写错误:
"errors.tolerance": "all", "errors.log.enable": true, "errors.log.include.messages": true, "errors.deadletterqueue.topic.name": "your-dlq-topic", // 可选,将错误消息转发至死信队列 "errors.deadletterqueue.topic.replication.factor": 1
添加死信队列可存储处理失败的消息,便于后续排查,同时不中断正常消息的处理流程。
内容的提问来源于stack exchange,提问作者javac
相关产品推荐
相关产品推荐

