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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 10:09:03