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

如何通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 03:20:16