如何将任意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
相关产品推荐
相关产品推荐

