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详情

更新:改用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
相关产品推荐
相关产品推荐

