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

如何配置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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 09:17:13