Kafka连接器开发中JSON字符串转Kafka Connect Struct方案咨询
问题解答
1 转换的高效实现方案
不需要自行实现JSON到Struct的转换逻辑,直接复用Kafka Connect生态已有的官方Converter组件是最高效的方案,这类官方实现已经完成了性能优化、边界场景兼容、Schema合法性校验,比自行编码的实现稳定性和执行效率更高。
实现逻辑:在连接器初始化阶段实例化Converter对象,传入你现有配置中的参数,后续在poll()方法中直接调用Converter的转换接口即可拿到可用于构造SourceRecord的Struct对象。
2 现成工具库及调用示例
你已经在Docker Compose中配置的io.confluent.connect.json.JsonSchemaConverter就是可以直接在代码中调用的转换工具,完全符合你给出的理想调用逻辑,示例代码如下:
import org.apache.kafka.connect.data.Struct; import org.apache.kafka.connect.data.Schema; import org.apache.kafka.connect.source.SourceRecord; import io.confluent.connect.json.JsonSchemaConverter; import org.apache.kafka.connect.storage.ConverterConfig; import java.util.HashMap; import java.util.Map; import org.json.JSONObject; // 连接器初始化阶段执行,只需实例化一次 JsonSchemaConverter keyConverter = new JsonSchemaConverter(); JsonSchemaConverter valueConverter = new JsonSchemaConverter(); // 传入你在Docker Compose中配置的参数 Map<String, Object> configs = new HashMap<>(); configs.put(ConverterConfig.SCHEMAS_ENABLE_CONFIG, true); configs.put("schema.registry.url", "http://schema-registry:8081"); // 分别配置key和value的转换器,isKey参数标记是否为消息key keyConverter.configure(configs, true); valueConverter.configure(configs, false); // poll()方法中执行转换逻辑 // 若输入为JSONObject,调用toString()转成字符串即可使用 String keyJsonStr = keyJSONObject.toString(); String valueJsonStr = valueJSONObject.toString(); // 转换key org.apache.kafka.connect.data.SchemaAndValue keySchemaAndValue = keyConverter.toConnectData(topicName, keyJsonStr.getBytes()); Schema keySchema = keySchemaAndValue.schema(); Struct keyStruct = (Struct) keySchemaAndValue.value(); // 转换value org.apache.kafka.connect.data.SchemaAndValue valueSchemaAndValue = valueConverter.toConnectData(topicName, valueJsonStr.getBytes()); Schema valSchema = valueSchemaAndValue.schema(); Struct valStruct = (Struct) valueSchemaAndValue.value(); // 构造SourceRecord SourceRecord sr = new SourceRecord( sourcePartition, sourceOffset, topicName, partition, keySchema, keyStruct, valSchema, valStruct);
如果你不需要对接Schema Registry,也可以使用Apache Kafka官方自带的无Registry版本转换器org.apache.kafka.connect.json.JsonConverter,使用逻辑完全一致,仅需调整配置参数即可。
内容的提问来源于stack exchange,提问作者Eric Broda
相关产品推荐
相关产品推荐

