Integrera Schema Registry med Kafka
Implementera Schema Registry i era Kafka-applikationer för att automatiskt hantera och upprätthålla datascheman.
Integrera Schema Registry med Kafka är en gratis lektion i Grunderna i Apache Kafka och strömbearbetning på CoddyKit. Detta är lektion 3 av 4. Du kan läsa vilka 3 lektioner som helst i den här lärvägen kostnadsfritt i sin helhet – därefter låser CoddyKit PRO upp alla lektioner, plus praktisk övning med en inbyggd kodredigerare och en AI-lärare dygnet runt. Den ingår i lärvägen för Grunderna i Apache Kafka och strömbearbetning, och Era framsteg synkroniseras mellan webben och CoddyKit-appen. Kursen i Grunderna i Apache Kafka och strömbearbetning innehåller totalt 4 lektioner.
Varför integrera Schema Registry?
Du har lärt dig om Kafka och Schema Registry. Nu kopplar vi ihop dem! Att integrera Schema Registry i dina Kafka-applikationer är avgörande för att säkerställa datakvalitet och kompatibilitet.
Det fungerar som ett centralt arkiv för scheman och gör det möjligt för producenter och konsumenter att validera och utveckla dataformat på ett säkert sätt.
Så fungerar det: serialiserare
För integrationen använder Kafka-klienter särskilda serialiserare och deserialiserare som kommunicerar med Schema Registry.
När en producent skickar data tar KafkaAvroSerializer (eller motsvarande för Protobuf/JSON Schema) emot dina data, registrerar schemat (om det är nytt) och lägger sedan till ett schema-ID före datan innan den skickas till Kafka.
Det viktigaste i producentkonfigurationen
För att få din Kafka-producent att fungera med Schema Registry måste du ange vissa egenskaper. De talar om för producenten var Schema Registry finns och vilken serialiserare som ska användas.
key.serializer: OftaStringSerializerellerKafkaAvroSerializer.value.serializer: Angeio.confluent.kafka.serializers.KafkaAvroSerializer.schema.registry.url: URL:en till din Schema Registry-instans (till exempelhttp://localhost:8081).
Producentkod: definiera ett Avro-schema
Innan vi skickar data måste vi definiera dess struktur med ett Avro-schema. För enkelhetens skull skapar vi ett grundläggande schema för 'User' med fälten name och age.
Det här schemat används för att skapa en 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!");
}
}Producentkod: skicka Avro-data
Här är en komplett Java-producentapplikation. Observera hur vi konfigurerar serialiserarna och URL:en till Schema Registry. Sedan skapar vi en GenericRecord utifrån vårt USER_SCHEMA och skickar den.
Prova att köra exemplet!
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();
}
}
}Så fungerar det: deserialiserare
På konsumentsidan har KafkaAvroDeserializer (eller motsvarande) motsatt funktion.
När en konsument tar emot ett meddelande läser deserialiseraren ut schema-ID:t, hämtar motsvarande schema från Schema Registry och använder sedan schemat för att korrekt deserialisera meddelandet tillbaka till applikationens datatyp (till exempel en GenericRecord eller ett specifikt Avro-objekt).
Det viktigaste i konsumentkonfigurationen
Precis som producenter behöver Kafka-konsumenter också vissa egenskaper för att fungera med Schema Registry:
key.deserializer: OftaStringDeserializerellerKafkaAvroDeserializer.value.deserializer: Angeio.confluent.kafka.serializers.KafkaAvroDeserializer.schema.registry.url: URL:en till din Schema Registry-instans.group.id: Ett unikt ID för din konsumentgrupp.auto.offset.reset: Definierar beteendet när ingen initial offset hittas (till exempelearliestellerlatest).
Konsumentkod: ta emot Avro-data
Den här konsumentapplikationen är konfigurerad för att läsa Avro-meddelanden från ämnet 'avro-users'. Den använder KafkaAvroDeserializer för att automatiskt hantera schemaupplösning.
Kör det här exemplet EFTER att du har kört producenten för att se datan!
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();
}
}
}Fördelar med sömlös integration
Att integrera Schema Registry med dina Kafka-klienter ger betydande fördelar:
- Datakompatibilitet: Säkerställer att producenter och konsumenter alltid förstår varandras dataformat.
- Schemauveckling: Uppdatera scheman säkert över tid utan att befintliga applikationer slutar fungera.
- Datastyrning: Centraliserad schemahantering ger en gemensam källa till sanning för dina datastrukturer.
- Mindre standardkod: Serialiserare/deserialiserare hanterar schemahanteringen automatiskt.
Snabbtest: konfigurera Schema Registry
Vilka av följande egenskaper är nödvändiga för att en Kafka-klient (producent eller konsument) ska kunna integrera med Confluent Schema Registry med hjälp av Avro?
Sammanfattning: integrera Schema Registry
I den här lektionen lärde du dig att integrera Confluent Schema Registry med dina Kafka-applikationer.
- Vi konfigurerade Kafka-producenter och konsumenter med `schema.registry.url`.
- Vi använde `KafkaAvroSerializer` och `KafkaAvroDeserializer` för att automatiskt hantera Avro-data.
- Du fick se praktiska exempel på hur man skickar och tar emot `GenericRecord`-objekt.
Den här integrationen är viktig för robusta, schemadrivna datapipelines!
Lär dig Grunderna i Apache Kafka och strömbearbetning med en AI-lärare – gratis
Skriv och kör riktig kod i webbläsaren, få omedelbar hjälp av en AI-lärare dygnet runt och fortsätt där du slutade – på webben eller i appen.
- Kurser
- 12
- Lektioner
- 48
Vanliga frågor
Är lektionen ”Integrera Schema Registry med Kafka” gratis?
Ja – du kan läsa vilka 3 lektioner som helst i lärvägen Grunderna i Apache Kafka och strömbearbetning, inklusive ”Integrera Schema Registry med Kafka”, kostnadsfritt i sin helhet här på webben. Därefter låser CoddyKit PRO upp alla lektioner, plus interaktiv övning med en inbyggd kodredigerare och en AI-lärare dygnet runt. Kursen i Grunderna i Apache Kafka och strömbearbetning innehåller totalt 4 lektioner.
Vad lär jag mig i ”Integrera Schema Registry med Kafka”?
Implementera Schema Registry i era Kafka-applikationer för att automatiskt hantera och upprätthålla datascheman. Ni övar på Grunderna i Apache Kafka och strömbearbetning med praktisk kod som körs direkt i webbläsaren, medan en AI-handledare som är tillgänglig dygnet runt svarar på Era frågor under lektionen.
Behöver jag någon erfarenhet för att börja lära mig Grunderna i Apache Kafka och strömbearbetning?
Du behöver inga förkunskaper. Utbildningen i Grunderna i Apache Kafka och strömbearbetning på CoddyKit är upplagd för allt från nybörjare till avancerade elever, så att du kan börja här eller från början och gå fram i din egen takt. Detta är lektion 3 av 4.
Hur lång tid tar lektionen ”Integrera Schema Registry med Kafka”?
De flesta CoddyKit-lektioner tar cirka 5–10 minuter. Varje lektion är kort och interaktiv, så att du gör stadiga framsteg och kan fortsätta precis där du slutade – på webben eller i appen.
Kan jag skriva och köra kod i den här Grunderna i Apache Kafka och strömbearbetning-lektionen?
Ja. Varje Grunderna i Apache Kafka och strömbearbetning-lektion innehåller en inbyggd kodredigerare, så att du kan skriva och köra riktig kod direkt i webbläsaren och få omedelbar AI-feedback – utan lokal installation.
Alla lektioner i den här kursen
- Varför behövs schemanhantering?
- Avro- och Protobuf-scheman
- Integrera Schema Registry med Kafka
- Schemautveckling och kompatibilitetslägen