Apache Kafka & Stream Processing Fundamentals · Lección

Integración de Schema Registry con Kafka

Implemente Schema Registry en sus aplicaciones de Kafka para gestionar y aplicar esquemas de datos automáticamente.

Lección 3 de 411 pasos

Integración de Schema Registry con Kafka es una lección gratuita de Apache Kafka & Stream Processing Fundamentals en CoddyKit. Esta es la lección 3 de 4. Puedes leer la lección completa abajo gratuitamente — luego la practicas en el navegador con un editor de código integrado y un tutor de IA 24/7. Forma parte de la ruta de aprendizaje de Apache Kafka & Stream Processing Fundamentals, y tu progreso se sincroniza en la web y la app de CoddyKit. El curso de Apache Kafka & Stream Processing Fundamentals incluye 4 lecciones en total.

Partes de esta lección aún no han sido traducidas y se muestran en inglés.

Why Integrate Schema Registry?

You've learned about Kafka and Schema Registry. Now, let's connect them! Integrating Schema Registry into your Kafka applications is vital for ensuring data quality and compatibility.

It acts as a central repository for schemas, allowing producers and consumers to validate and evolve data formats safely.

How it Works: Serializers

To integrate, Kafka clients use special serializers and deserializers that communicate with the Schema Registry.

When a producer sends data, the KafkaAvroSerializer (or Protobuf/JSON Schema equivalent) takes your data, registers its schema (if new), and then prefixes the data with a schema ID before sending it to Kafka.

Producer Configuration Essentials

To make your Kafka producer work with Schema Registry, you need to set specific properties. These tell the producer where the Schema Registry is and which serializer to use.

  • key.serializer: Often StringSerializer or KafkaAvroSerializer.
  • value.serializer: Set this to io.confluent.kafka.serializers.KafkaAvroSerializer.
  • schema.registry.url: The URL of your Schema Registry instance (e.g., http://localhost:8081).

Producer Code: Defining an Avro Schema

Before we send data, we need to define its structure using an Avro schema. For simplicity, we'll create a basic 'User' schema with a name and age field.

This schema will be used to create a GenericRecord.

import org.apache.avro.Schema;

public class AvroSchemaDef {
  public static final String USER_SCHEMA_JSON = 
    "{\"namespace\": \"com.coddykit\", " +
    "\"type\": \"record\", " +
    "\"name\": \"User\", " +
    "\"fields\": [" +
    "{\"name\": \"name\", \"type\": \"string\"}," +
    "{\"name\": \"age\", \"type\": \"int\"}]}";

  public static final Schema USER_SCHEMA = 
    new Schema.Parser().parse(USER_SCHEMA_JSON);

  public static void main(String[] args) {
    System.out.println("User Schema Defined!");
  }
}

Producer Code: Sending Avro Data

Here's a complete Java producer application. Notice how we configure the serializers and the Schema Registry URL. We then create a GenericRecord based on our USER_SCHEMA and send it.

Try running this example!

import org.apache.kafka.clients.producer.*;
import io.confluent.kafka.serializers.KafkaAvroSerializer;
import org.apache.avro.Schema;
import org.apache.avro.generic.GenericData;
import org.apache.avro.generic.GenericRecord;

import java.util.Properties;

public class AvroProducer {

  public static final String USER_SCHEMA_JSON = 
    "{\"namespace\": \"com.coddykit\", " +
    "\"type\": \"record\", " +
    "\"name\": \"User\", " +
    "\"fields\": [" +
    "{\"name\": \"name\", \"type\": \"string\"}," +
    "{\"name\": \"age\", \"type\": \"int\"}]}";

  public static final Schema USER_SCHEMA = 
    new Schema.Parser().parse(USER_SCHEMA_JSON);

  public static void main(String[] args) {
    Properties props = new Properties();
    props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
    props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
    props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, KafkaAvroSerializer.class.getName());
    props.put("schema.registry.url", "http://localhost:8081");

    Producer<String, GenericRecord> producer = new KafkaProducer<>(props);
    String topic = "avro-users";

    GenericRecord user = new GenericData.Record(USER_SCHEMA);
    user.put("name", "Coddy");
    user.put("age", 5);

    ProducerRecord<String, GenericRecord> record = new ProducerRecord<>(topic, "user-1", user);

    try {
      producer.send(record, (metadata, exception) -> {
        if (exception == null) {
          System.out.println("Sent record to topic " + metadata.topic() + " partition " + metadata.partition() + " offset " + metadata.offset());
        } else {
          exception.printStackTrace();
        }
      });
    } finally {
      producer.flush();
      producer.close();
    }
  }
}

How it Works: Deserializers

On the consumer side, the KafkaAvroDeserializer (or equivalent) plays the opposite role.

When a consumer receives a message, the deserializer extracts the schema ID, fetches the corresponding schema from Schema Registry, and then uses that schema to correctly deserialize the message back into your application's data type (e.g., a GenericRecord or a specific Avro object).

Consumer Configuration Essentials

Similar to producers, Kafka consumers also need specific properties to work with Schema Registry:

  • key.deserializer: Often StringDeserializer or KafkaAvroDeserializer.
  • value.deserializer: Set this to io.confluent.kafka.serializers.KafkaAvroDeserializer.
  • schema.registry.url: The URL of your Schema Registry instance.
  • group.id: A unique ID for your consumer group.
  • auto.offset.reset: Defines behavior when no initial offset is found (e.g., earliest or latest).

Consumer Code: Receiving Avro Data

This consumer application is configured to read Avro messages from the 'avro-users' topic. It uses KafkaAvroDeserializer to automatically handle schema resolution.

Run this example AFTER running the producer to see the data!

import org.apache.kafka.clients.consumer.*;
import io.confluent.kafka.serializers.KafkaAvroDeserializer;
import org.apache.avro.generic.GenericRecord;

import java.time.Duration;
import java.util.Collections;
import java.util.Properties;

public class AvroConsumer {
  public static void main(String[] args) {
    Properties props = new Properties();
    props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
    props.put(ConsumerConfig.GROUP_ID_CONFIG, "avro-consumer-group");
    props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
    props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, KafkaAvroDeserializer.class.getName());
    props.put("schema.registry.url", "http://localhost:8081");
    props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");

    Consumer<String, GenericRecord> consumer = new KafkaConsumer<>(props);
    String topic = "avro-users";
    consumer.subscribe(Collections.singletonList(topic));

    System.out.println("Listening for messages on topic: " + topic);

    try {
      while (true) {
        ConsumerRecords<String, GenericRecord> records = consumer.poll(Duration.ofMillis(100));
        for (ConsumerRecord<String, GenericRecord> record : records) {
          System.out.printf("Received record (key=%s, value=%s, partition=%d, offset=%d)\n",
                            record.key(), record.value(), record.partition(), record.offset());
          GenericRecord user = record.value();
          System.out.println("  User Name: " + user.get("name") + ", Age: " + user.get("age"));
        }
      }
    } finally {
      consumer.close();
    }
  }
}

Benefits of Seamless Integration

Integrating Schema Registry with your Kafka clients offers significant advantages:

  • Data Compatibility: Ensures producers and consumers always understand each other's data formats.
  • Schema Evolution: Safely update schemas over time without breaking existing applications.
  • Data Governance: Centralized schema management provides a single source of truth for your data structures.
  • Reduced Boilerplate: Serializers/deserializers handle schema management automatically.

Quick Check: Schema Registry Setup

Which of the following properties are essential for a Kafka client (producer or consumer) to integrate with Confluent Schema Registry using Avro?

Recap: Integrating Schema Registry

In this lesson, you learned how to integrate Confluent Schema Registry with your Kafka applications.

  • We configured Kafka producers and consumers with `schema.registry.url`.
  • We used `KafkaAvroSerializer` and `KafkaAvroDeserializer` to handle Avro data automatically.
  • You saw practical examples of sending and receiving `GenericRecord`s.

This integration is key for robust, schema-driven data pipelines!

Gratis para empezar

Aprende Apache Kafka & Stream Processing Fundamentals con un tutor de IA — gratis

Escribe y ejecuta código real en tu navegador, obtén ayuda instantánea de un tutor de IA disponible 24/7 y continúa donde lo dejaste en la web o en la aplicación.

Cursos
12
Lecciones
48

Preguntas frecuentes

¿La lección «Integración de Schema Registry con Kafka» es gratis?

Sí — el texto completo de «Integración de Schema Registry con Kafka» es gratis para leer aquí en la web. Para practicarla de forma interactiva (editor de código integrado y tutor de IA 24/7) y desbloquear el resto del curso de Apache Kafka & Stream Processing Fundamentals, actualiza a CoddyKit PRO. El curso de Apache Kafka & Stream Processing Fundamentals incluye 4 lecciones en total.

¿Qué aprenderé en «Integración de Schema Registry con Kafka»?

Implemente Schema Registry en sus aplicaciones de Kafka para gestionar y aplicar esquemas de datos automáticamente. Practicas Apache Kafka & Stream Processing Fundamentals con código real que ejecutas directamente en el navegador, y un tutor de IA 24/7 responde tus preguntas mientras trabajas en la lección.

¿Necesito experiencia previa para empezar Apache Kafka & Stream Processing Fundamentals?

No se requiere experiencia previa. Apache Kafka & Stream Processing Fundamentals en CoddyKit está estructurado para principiantes hasta estudiantes avanzados, así que puedes empezar aquí o desde el inicio y avanzar a tu ritmo. Esta es la lección 3 de 4.

¿Cuánto tiempo toma la lección «Integración de Schema Registry con Kafka»?

La mayoría de las lecciones de CoddyKit toman alrededor de 5–10 minutos. Cada una es compacta e interactiva, así que avanzas constantemente y retomas exactamente por donde dejaste en la web y la app.

¿Puedo escribir y ejecutar código en esta lección de Apache Kafka & Stream Processing Fundamentals?

Sí. Cada lección de Apache Kafka & Stream Processing Fundamentals incluye un editor de código integrado, así que escribes y ejecutas código real directamente en tu navegador y obtienes retroalimentación instantánea de IA — sin configuración local necesaria.

Todas las lecciones de este curso

  1. ¿Por qué gestionar esquemas?
  2. Esquemas Avro y Protobuf
  3. Integración de Schema Registry con Kafka
  4. Evolución de esquemas y modos de compatibilidad
← Volver a Apache Kafka & Stream Processing Fundamentals