如何通过Kafka将JSON消息正确发送至SQL Server的JDBC Sink连接器
为无Schema的JSON消息添加Schema&Payload结构以适配Kafka JDBC Sink连接器
针对你的场景,优先推荐无需修改Java生产者代码的方案,通过Kafka Connect的转换器或Transforms功能实现;如果必须在生产者端处理,再考虑修改代码。以下是具体实现方式:
方案1:无需修改生产者代码,通过Kafka Connect配置处理
方式1:JSONConverter + Transforms(配置简单,推荐)
JDBC Sink连接器支持带Schema信息的JSON格式,你可以通过Kafka Connect的Transforms功能,将纯JSON消息自动包装成schema+payload的结构,配合JSONConverter使用。
Sink连接器配置示例
name=sqlserver-jdbc-sink connector.class=io.confluent.connect.jdbc.JdbcSinkConnector tasks.max=1 # SQL Server连接信息 connection.url=jdbc:sqlserver://<你的服务器地址>:1433;databaseName=<你的数据库名> connection.user=<数据库用户名> connection.password=<数据库密码> # 要消费的Kafka Topic topics=<你的Topic名称> # 自动创建/更新表结构 auto.create=true auto.evolve=true # Key转换器(不需要Schema) key.converter=org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable=false # Value转换器(开启Schema支持) value.converter=org.apache.kafka.connect.json.JsonConverter value.converter.schemas.enable=true # 添加Transforms:将原始消息包装为payload,并插入静态Schema transforms=wrapPayload,addSchema # 第一步:将原始JSON提升为payload字段 transforms.wrapPayload.type=org.apache.kafka.connect.transforms.HoistField$Value transforms.wrapPayload.field=payload # 第二步:插入静态Schema字段(匹配你的消息结构) transforms.addSchema.type=org.apache.kafka.connect.transforms.InsertField$Value transforms.addSchema.static.field=schema transforms.addSchema.static.value={"type":"struct","fields":[{"field":"DagId","type":"string"},{"field":"RunId","type":"string"},{"field":"ChatKey","type":"string"},{"field":"ConversationId","type":"string"},{"field":"EventTimestamp","type":["null","string"]},{"field":"EventType","type":"string"},{"field":"MessageType","type":"string"},{"field":"LastUpdateDatetime","type":"string"}]}
配置生效后,原始纯JSON会被自动转换为JDBC Sink能识别的结构:
{ "schema": {"type":"struct","fields":[{"field":"DagId","type":"string"},...]} "payload": { "DagId": "chat-bot-process-v1.0", "RunId": "scheduled__2021-07-25T10:00:00+00:00", ... } }
方式2:JSONSchemaConverter + Schema Registry(适合长期Schema管理)
如果你需要集中管理消息Schema,可以使用Confluent的JSONSchemaConverter,搭配Schema Registry服务自动处理Schema注册和消息包装。
Sink连接器配置示例
name=sqlserver-jdbc-sink connector.class=io.confluent.connect.jdbc.JdbcSinkConnector tasks.max=1 # SQL Server连接信息 connection.url=jdbc:sqlserver://<你的服务器地址>:1433;databaseName=<你的数据库名> connection.user=<数据库用户名> connection.password=<数据库密码> topics=<你的Topic名称> auto.create=true auto.evolve=true # Key转换器 key.converter=org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable=false # Value转换器(搭配Schema Registry) value.converter=io.confluent.connect.json.JsonSchemaConverter value.converter.schema.registry.url=http://<你的Schema Registry地址>:8081
这种方式下,Schema Registry会自动为你的消息生成并注册Schema,消息会被包装为带Schema ID和Payload的格式,无需手动维护Schema内容。
方案2:修改Java生产者代码,直接发送带Schema&Payload的消息
如果必须在生产者端处理消息结构,可以直接构造schema+payload格式的JSON发送:
Java生产者代码示例
import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerRecord; import com.fasterxml.jackson.databind.ObjectMapper; import java.util.HashMap; import java.util.Map; import java.util.Properties; public class SchemaWrappedJsonProducer { public static void main(String[] args) { Properties props = new Properties(); props.put("bootstrap.servers", "<你的Kafka Broker地址>"); props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer"); KafkaProducer<String, String> producer = new KafkaProducer<>(props); ObjectMapper mapper = new ObjectMapper(); // 构造原始消息Payload Map<String, Object> payload = new HashMap<>(); payload.put("DagId", "chat-bot-process-v1.0"); payload.put("RunId", "scheduled__2021-07-25T10:00:00+00:00"); payload.put("ChatKey", "82a4daf8-c1be-4524-bb80-ec252b38c020"); payload.put("ConversationId", "2158db2e-0bcc-48e6-a96e-3347e156a90a"); payload.put("EventTimestamp", null); payload.put("EventType", "ASYNC"); payload.put("MessageType", "EPORTWEB"); payload.put("LastUpdateDatetime", "2022-09-21T17:05:51.473-04:00"); // 构造匹配消息结构的Schema Map<String, Object> schema = new HashMap<>(); schema.put("type", "struct"); schema.put("fields", new Object[]{ Map.of("field", "DagId", "type", "string"), Map.of("field", "RunId", "type", "string"), Map.of("field", "ChatKey", "type", "string"), Map.of("field", "ConversationId", "type", "string"), Map.of("field", "EventTimestamp", "type", new String[]{"null", "string"}), Map.of("field", "EventType", "type", "string"), Map.of("field", "MessageType", "type", "string"), Map.of("field", "LastUpdateDatetime", "type", "string") }); // 包装为最终消息结构 Map<String, Object> finalMsg = new HashMap<>(); finalMsg.put("schema", schema); finalMsg.put("payload", payload); try { String jsonStr = mapper.writeValueAsString(finalMsg); ProducerRecord<String, String> record = new ProducerRecord<>("<你的Topic名称>", jsonStr); producer.send(record); producer.flush(); } catch (Exception e) { e.printStackTrace(); } finally { producer.close(); } } }
总结
- 优先选择方案1的JSONConverter+Transforms,无需修改生产者代码,配置成本低;
- 若需要长期维护和管理消息Schema,选择JSONSchemaConverter+Schema Registry;
- 仅当业务要求必须在生产者端处理时,再采用方案2修改Java代码。
内容的提问来源于stack exchange,提问作者cluis92
相关产品推荐
相关产品推荐

