如何配置Kafka-MongoDB Sink Connector处理非JSON格式消息
Kafka Connect MongoDB Sink: Handling Non-JSON Messages
我来帮你搞定这两个处理非JSON消息的问题,都是Kafka Connect生产环境里常碰到的场景~
1. 忽略非JSON消息,避免Sink任务崩溃
要实现这个,核心是利用Kafka Connect的错误容忍机制,让连接器遇到解析失败的消息时,要么跳过要么转存到死信队列(DLQ),而不是直接崩溃。你只需要在现有配置里添加以下几个参数:
# 允许连接器继续处理消息,即使遇到错误 errors.tolerance=all # 开启死信队列,把无法处理的消息转存到这个topic(需要提前创建好) errors.deadletterqueue.topic.name=test-dlq # 可选:设置死信队列的保留时间,比如7天 errors.deadletterqueue.retention.ms=604800000
errors.tolerance=all:告诉连接器遇到错误时不要停止,继续处理后续消息errors.deadletterqueue.topic.name:指定死信队列的topic,所有解析失败的消息会被发送到这里,方便你后续排查问题- 记得提前用
kafka-topics.sh创建好这个死信队列topic哦
修改后的完整配置片段(新增部分已标注):
name=mongo-sink topics=test connector.class=com.mongodb.kafka.connect.MongoSinkConnector tasks.max=1 key.ignore=true connection.uri=mongodb://localhost:27017 database=test_kafka collection=transaction max.num.retries=3 retries.defer.timeout=5000 type.name=kafka-connect key.converter=org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable=false value.converter=org.apache.kafka.connect.json.JsonConverter value.converter.schemas.enable=false # 新增错误处理配置 errors.tolerance=all errors.deadletterqueue.topic.name=test-dlq errors.deadletterqueue.retention.ms=604800000
2. 将非JSON消息转为带默认键的JSON格式
如果不想丢弃这些消息,而是要转成{"non-json": "abc"}这种格式插入MongoDB,你需要用到Kafka Connect的Single Message Transform(SMT),具体用ScriptTransform来写一段简单的处理脚本,自动修复非JSON消息。
步骤1:添加SMT配置到你的连接器属性里
# 启用SMT,命名为fixNonJson transforms=fixNonJson # 指定SMT类型为ScriptTransform(处理消息值) transforms.fixNonJson.type=org.apache.kafka.connect.transforms.ScriptTransform$Value # 指定处理脚本的路径(比如放在/opt/kafka/scripts/fix-non-json.groovy) transforms.fixNonJson.script.path=/opt/kafka/scripts/fix-non-json.groovy
步骤2:编写Groovy处理脚本fix-non-json.groovy
这个脚本会尝试把消息值解析成JSON,如果失败就把原始字符串包装成{"non-json": 原始内容}的格式:
import groovy.json.JsonSlurper def jsonSlurper = new JsonSlurper() try { // 尝试解析JSON,如果成功就直接返回原内容 return jsonSlurper.parseText(value) } catch (Exception e) { // 解析失败,包装成指定格式的JSON return [ "non-json": value ] }
步骤3:整合后的完整配置
把SMT配置和之前的错误处理配置结合起来,最终的MongoSinkConnector.properties如下:
name=mongo-sink topics=test connector.class=com.mongodb.kafka.connect.MongoSinkConnector tasks.max=1 key.ignore=true connection.uri=mongodb://localhost:27017 database=test_kafka collection=transaction max.num.retries=3 retries.defer.timeout=5000 type.name=kafka-connect key.converter=org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable=false value.converter=org.apache.kafka.connect.json.JsonConverter value.converter.schemas.enable=false # 错误处理配置 errors.tolerance=all errors.deadletterqueue.topic.name=test-dlq errors.deadletterqueue.retention.ms=604800000 # SMT转换配置 transforms=fixNonJson transforms.fixNonJson.type=org.apache.kafka.connect.transforms.ScriptTransform$Value transforms.fixNonJson.script.path=/opt/kafka/scripts/fix-non-json.groovy
这样配置后,不管是合法的JSON还是非JSON消息,你的Mongo Sink任务都能稳定运行:合法消息正常插入,非JSON消息会被自动转成指定格式存入MongoDB~
内容的提问来源于stack exchange,提问作者Vu Le Anh
相关产品推荐
相关产品推荐

