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

如何让Kafka JDBC源连接器输出嵌套JSON格式的消息载荷?

Kafka消息载荷格式转换问题

当前Kafka Topic中的消息载荷格式如下:

{
  "id": 73,
  "message": "{\"id\":3,\"firstName\":\"firstname\",\"lastName\":\"lastname\",\"emailAddress\":\"firstname.lastname@gmail.com\"}"
}

希望将载荷转换为如下格式(message字段从JSON字符串变为嵌套JSON对象):

{
  "id": 73,
  "message": {"id":3,"firstName":"firstname","lastName":"lastname","emailAddress":"firstname.lastname@gmail.com"}
}

请问应该使用转换器(Converter)还是单消息转换(SMT)来实现?
已尝试更换JSON、Avro类型的转换器,但没有效果;也尝试自定义SMT,同样未能成功。


解决方案:使用单消息转换(SMT)

为什么不选转换器?

转换器负责的是整个消息的序列化/反序列化(比如字节与JSON/AVRO对象的互转),它不会针对消息内部单个字段做结构转换。你的场景里,原始消息本身是合法的JSON格式,message字段的类型就是字符串,转换器只会按原样解析这个字段,不会自动把内部的JSON字符串解析成嵌套对象,所以更换转换器无法解决问题。

自定义SMT的实现要点

你需要编写自定义SMT来完成message字段的解析替换,核心逻辑如下:

  1. 继承Kafka Connect的AbstractTransformation<Map<String, Object>>(无Schema模式)或AbstractTransformation<Struct>(Schema模式)。
  2. 在apply方法中,提取消息里的message字段值(JSON字符串)。
  3. 用JSON解析库(如Jackson)把该字符串解析为JSON对象(Map或JsonNode)。
  4. 替换原消息中的message字段为解析后的JSON对象。
  5. 添加错误处理:如果message不是合法JSON字符串,可选择记录日志、保留原字段值或标记消息失败,避免影响整个任务。

示例代码片段(无Schema模式):

import org.apache.kafka.connect.transforms.AbstractTransformation;
import org.apache.kafka.connect.connector.ConnectRecord;
import com.fasterxml.jackson.databind.ObjectMapper;
import java.util.Map;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

public class ParseNestedJsonSMT extends AbstractTransformation<Map<String, Object>> {
    private static final Logger log = LoggerFactory.getLogger(ParseNestedJsonSMT.class);
    private final ObjectMapper objectMapper = new ObjectMapper();

    @Override
    public org.apache.kafka.common.config.ConfigDef config() {
        return new org.apache.kafka.common.config.ConfigDef(); // 可扩展添加配置项,比如指定要解析的字段名
    }

    @Override
    public void configure(Map<String, ?> configs) {
        // 初始化配置(如果有)
    }

    @Override
    public Map<String, Object> apply(ConnectRecord<Map<String, Object>> record) {
        Map<String, Object> value = record.value();
        if (value == null || !value.containsKey("message")) {
            return value;
        }

        Object messageValue = value.get("message");
        if (!(messageValue instanceof String)) {
            log.warn("Message field is not a string, skip parsing");
            return value;
        }

        String messageStr = (String) messageValue;
        try {
            Map<String, Object> messageObj = objectMapper.readValue(messageStr, Map.class);
            value.put("message", messageObj);
        } catch (Exception e) {
            log.error("Failed to parse message content: {}", messageStr, e);
            // 可选择保留原字段值,或抛出异常终止消息处理
        }
        return value;
    }

    @Override
    public void close() {
        // 资源清理操作
    }
}

部署注意事项

  • 把自定义SMT打包成JAR文件,放到Kafka Connect的插件目录下。
  • 在Connect任务配置中添加该SMT,示例配置:
transforms=parseNestedMessage
transforms.parseNestedMessage.type=com.your.package.ParseNestedJsonSMT

内容的提问来源于stack exchange,提问作者Ckey

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 04:50:14