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

Confluent Schema Registry未对JSON数据执行校验的问题求助

问题描述

尝试向关联Confluent Schema Registry的Kafka Topic发送JSON数据时,发现Topic会接收所有不符合Schema规范的内容。

client.properties

bootstrap.servers=pkc-xxx.northamerica-northeast1.gcp.confluent.cloud:9092
security.protocol=SASL_SSL
sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required username='xxx' password='Y0LiwuvfvDts/iQEXntX/yyy';
sasl.mechanism=PLAIN
key.serializer=org.apache.kafka.common.serialization.StringSerializer
#value.serializer=org.apache.kafka.common.serialization.StringSerializer

value.serializer=io.confluent.kafka.serializers.json.KafkaJsonSchemaSerializer
confluent.value.schema.validation=true
latest.compatibility.strict=false
auto.register.schemas=false
use.latest.version=true


# Required connection configs for Confluent Cloud Schema Registry
schema.registry.url=https://xxx.us-east1.gcp.confluent.cloud
basic.auth.credentials.source=USER_INFO
basic.auth.user.info=xxx:YrE7Oztqrx7FvGKF6Ofk+yyy+YySi

Topic Schema

{
  "$id": "http://example.com/myURI.schema.json",
  "$schema": "http://json-schema.org/draft-07/schema#",
  "additionalProperties": false,
  "description": "Sample schema to help you get started.",
  "properties": {
    "age": {
      "description": "The string type is used for strings of text.",
      "type": "string"
    }
  },
  "required": [
    "age"
  ],
  "title": "SampleRecord",
  "type": "object"
}

测试代码(Map类型)

@Test
public void testMapUser() throws IOException {
    String topic = "gcp.test_user_schema";
    final Properties props = new Properties();
    InputStream inputStream = Files.newInputStream(Paths.get("src/main/resources/client.properties"));
    props.load(inputStream);
    Producer<String, Map> producer = new KafkaProducer<>(props);

    Gson gson = new Gson();
    File initialFile = new File("src/test/resources/user.json");
    InputStream in = Files.newInputStream(initialFile.toPath());
    BufferedReader reader = new BufferedReader(new InputStreamReader(in, StandardCharsets.UTF_8));
    HashMap event = gson.fromJson(reader, HashMap.class);


    ProducerRecord<String, Map> kafkaRecord = new ProducerRecord<>(topic, event);

    // Send the record to Kafka topic
    producer.send(kafkaRecord);
    producer.send(kafkaRecord, (metadata, exception) -> {
        if (exception != null) {
            exception.printStackTrace();
        } else {
            System.out.println("Message sent successfully!");
            System.out.println("Partition: " + metadata.partition());
            System.out.println("Offset: " + metadata.offset());
        }
    });

    producer.flush();
}

user.json

{
  "age": 30,
  "name": "gaurang"
}

Topic详情

Topic详情截图

更新:改用JsonNode对象测试

问题依旧,Topic仍接收所有不符合规范的消息。

@Test
public void TestJsonNodeObject() throws IOException, ExecutionException, InterruptedException {
    String topic = "gcp.test_user_schema";
    final Properties props = new Properties();
    InputStream inputStream = Files.newInputStream(Paths.get("src/main/resources/client.properties"));
    props.load(inputStream);
    Producer<String, JsonNode> producer = new KafkaProducer<>(props);
    
    File initialFile = new File("src/test/resources/user.json");
    InputStream in = Files.newInputStream(initialFile.toPath());
    BufferedReader reader = new BufferedReader(new InputStreamReader(in, StandardCharsets.UTF_8));

    ObjectMapper objectMapper = new ObjectMapper();
    JsonNode jsonNode = objectMapper.readTree(reader);
    ProducerRecord<String, JsonNode> kafkaRecord = new ProducerRecord<>(topic, jsonNode);

    // Send the record to Kafka topic
    producer.send(kafkaRecord).get();
    producer.flush();
}

内容的提问来源于stack exchange,提问作者Gaurang Shah

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 03:45:05