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

如何将任意JSON字符串转换为Kafka Schema并生成Kafka Connect可用的SourceRecord

通用JSON转Kafka Connect Schema与SourceRecord实现方案

实现基于Kafka Connect原生API和Jackson JsonNode完成,支持任意复杂度JSON的Schema自动推导、Struct自动构建,无需提前定义固定Schema,可直接在poll()方法中调用。

依赖引入

Kafka Connect运行环境默认自带以下依赖,打包时可设为provided scope:

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>connect-json</artifactId>
    <version>你的Kafka集群对应版本号</version>
    <scope>provided</scope>
</dependency>
<dependency>
    <groupId>com.fasterxml.jackson.core</groupId>
    <artifactId>jackson-databind</artifactId>
    <version>匹配Kafka对应Jackson版本</version>
    <scope>provided</scope>
</dependency>

核心转换工具实现

封装通用转换逻辑,递归处理嵌套对象、数组等所有JSON结构:

import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.kafka.connect.data.*;
import java.util.ArrayList;
import java.util.Iterator;
import java.util.List;
import java.util.Map;

public class GenericJsonConverter {
    private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();

    // 入参为任意JSON字符串,返回推导好的Schema和对应的Struct对象
    public static Map.Entry<Schema, Struct> convert(String jsonStr) throws Exception {
        JsonNode jsonNode = OBJECT_MAPPER.readTree(jsonStr);
        Schema schema = buildSchema(jsonNode, "AutoGeneratedJsonSchema");
        Struct struct = buildStruct(schema, jsonNode);
        return Map.entry(schema, struct);
    }

    // 递归推导JSON对应的Schema
    private static Schema buildSchema(JsonNode jsonNode, String schemaName) {
        SchemaBuilder builder;
        switch (jsonNode.getNodeType()) {
            case STRING:
                builder = SchemaBuilder.string();
                break;
            case BOOLEAN:
                builder = SchemaBuilder.bool();
                break;
            case NUMBER:
                builder = jsonNode.isIntegralNumber() ? 
                    (jsonNode.isLong() ? SchemaBuilder.int64() : SchemaBuilder.int32()) : 
                    SchemaBuilder.float64();
                break;
            case ARRAY:
                JsonNode firstElement = jsonNode.elements().next();
                Schema elementSchema = buildSchema(firstElement, schemaName + ".ArrayItem");
                builder = SchemaBuilder.array(elementSchema);
                break;
            case OBJECT:
                builder = SchemaBuilder.struct().name(schemaName).version(1);
                Iterator<Map.Entry<String, JsonNode>> fields = jsonNode.fields();
                while (fields.hasNext()) {
                    Map.Entry<String, JsonNode> field = fields.next();
                    Schema fieldSchema = buildSchema(field.getValue(), schemaName + "." + field.getKey());
                    // 默认所有字段设为可选,可根据业务需求调整
                    builder.field(field.getKey(), fieldSchema.optional().build());
                }
                break;
            case NULL:
            default:
                builder = SchemaBuilder.string().optional();
        }
        return jsonNode.isNull() ? builder.optional().build() : builder.build();
    }

    // 递归构建Kafka Connect兼容的Struct对象
    private static Object buildStruct(Schema schema, JsonNode jsonNode) {
        if (jsonNode.isObject()) {
            Struct struct = new Struct(schema);
            for (Field field : schema.fields()) {
                JsonNode fieldNode = jsonNode.get(field.name());
                if (fieldNode != null && !fieldNode.isNull()) {
                    struct.put(field.name(), buildStruct(field.schema(), fieldNode));
                }
            }
            return struct;
        } else if (jsonNode.isArray()) {
            List<Object> list = new ArrayList<>();
            Schema elementSchema = schema.valueSchema();
            for (JsonNode element : jsonNode) {
                list.add(buildStruct(elementSchema, element));
            }
            return list;
        } else if (jsonNode.isTextual()) {
            return jsonNode.asText();
        } else if (jsonNode.isBoolean()) {
            return jsonNode.asBoolean();
        } else if (jsonNode.isIntegralNumber()) {
            return jsonNode.isLong() ? jsonNode.asLong() : jsonNode.asInt();
        } else if (jsonNode.isFloatingPointNumber()) {
            return jsonNode.asDouble();
        }
        return null;
    }
}

poll()方法调用示例

直接替换原有固定Schema的逻辑即可,也可兼容你原有字段结构:

@Override
public List<SourceRecord> poll() throws InterruptedException {
    List<SourceRecord> records = new ArrayList<>();
    // 拉取到的任意JSON字符串
    String rawJson = fetchFromHttpApi();
    long id = generateId();
    Object key = buildKey(id);
    Schema keySchema = HttpSourceSchemas.KEY_SCHEMA;
    String timestampStr = getCurrentTimestamp();

    try {
        Map.Entry<Schema, Struct> valueEntry = GenericJsonConverter.convert(rawJson);
        // 如果需要保留原有timestamp、data的外层结构,可做一层封装
        Schema wrappedValueSchema = SchemaBuilder.struct()
                .field(HttpSourceSchemas.TIMESTAMP_FIELD, Schema.STRING_SCHEMA)
                .field(HttpSourceSchemas.DATA_FIELD, valueEntry.getKey())
                .build();
        Struct wrappedValue = new Struct(wrappedValueSchema)
                .put(HttpSourceSchemas.TIMESTAMP_FIELD, timestampStr)
                .put(HttpSourceSchemas.DATA_FIELD, valueEntry.getValue());

        records.add(new SourceRecord(
                sourcePartition, sourceOffset, topic, partition,
                keySchema, key, wrappedValueSchema, wrappedValue));
    } catch (Exception e) {
        // 自行补充异常处理逻辑
        e.printStackTrace();
    }
    return records;
}

注意事项

  • 若业务中存在同名字段类型不固定的情况,可在buildSchema方法中增加类型兼容逻辑,比如所有数值统一用INT64或者STRING类型兜底
  • 相同结构的JSON生成的Schema会被Kafka Connect自动缓存,不会重复生成,性能损耗极低
  • 支持任意层级的嵌套对象、嵌套数组自动解析,无需额外适配

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 03:15:03