실무 예제 · 7
Avro + Schema Registry 스키마 진화
JSON으로 이벤트를 주고받다가 필드 하나를 바꾸면
어느 컨슈머가 깨지는지 아무도 모릅니다.
Schema Registry는 그 계약을 브로커 밖에 명시적으로 두고,
호환되지 않는 변경을 등록 시점에 거부합니다.
이 예제는 BACKWARD 호환으로 필드를 추가하고,
왜 컨슈머를 먼저 배포해야 하는지를 실행으로 확인합니다.
학습 목표
- Avro wire format(magic byte + schema id + payload)을 바이트 단위로 확인할 수 있습니다.
- 호환성 모드 7종을 구분하고, 어떤 변경이 어느 모드에서 허용되는지 판단할 수 있습니다.
BACKWARD일 때 컨슈머를 먼저,FORWARD일 때 프로듀서를 먼저 배포해야 하는 이유를 설명할 수 있습니다.auto.register.schemas가 왜 프로덕션에서 위험한지 알 수 있습니다.
시나리오
회원 도메인이 users.signup 토픽에 가입 이벤트를 발행합니다.
컨슈머는 세 팀입니다 — 마케팅(웰컴 메일), 분석(가입 퍼널), CS(고객 조회).
마케팅 팀이 유입 채널(referralSource) 필드를 추가해 달라고 요청했습니다.
JSON을 쓰던 시절에는 프로듀서가 필드를 추가해도 컨슈머들이 무시하기만 하면 되니 문제가 없어 보였습니다.
실제로 사고가 난 것은 필드 이름을 signupAt → createdAt으로 바꿨을 때였습니다.
세 컨슈머 중 두 개가 NullPointerException으로 죽고,
파티션이 막혀 lag이 40분간 쌓였습니다.
이번에는 스키마를 계약으로 만들고, 호환되지 않는 변경은 배포 전에 거부되게 만듭니다.
아키텍처
스키마 자체는 메시지에 들어가지 않습니다. 메시지에는 4바이트 정수 schema id만 들어가고, 컨슈머는 그 id로 Registry에서 스키마를 받아 캐시합니다. 그래서 Avro 메시지는 JSON보다 훨씬 작습니다 — 필드 이름이 반복되지 않기 때문입니다.
| 형식 | 스키마 | 크기 | 스키마 진화 | 적합한 상황 |
|---|---|---|---|---|
| JSON (스키마 없음) | 없음 | 가장 큼 (필드명 반복) | 검증 불가 — 깨지는 시점이 런타임 | 프로토타입, 사람이 읽어야 하는 로그 |
| JSON Schema | Registry | 큼 | 가능 | 기존 JSON에서 점진적으로 넘어갈 때 |
| Avro 이 예제 | Registry (id 4바이트) | 작음 | 가장 성숙 — default 규칙이 명확 | Kafka 생태계의 사실상 표준. Connect·Streams 지원이 가장 두텁습니다 |
| Protobuf | Registry (id 4바이트) | 작음 | 가능 (필드 번호 기반) | gRPC를 이미 쓰는 조직. 다언어 환경 |
사전 요구사항
| 항목 | 버전 | 비고 |
|---|---|---|
| Apache Kafka | 4.3.1 | 예제 1의 클러스터 |
| Confluent Schema Registry | 8.2.2 | Docker 이미지 confluentinc/cp-schema-registry:8.2.2 |
| Apache Avro | 1.12.1 | Maven Central. avro + avro-maven-plugin |
kafka-avro-serializer | 8.2.2 | Schema Registry와 같은 버전. Maven Central이 아니라 packages.confluent.io에서 받습니다 |
| Java | 17 이상 |
전체 코드
avro-evolution/
├── docker-compose.schema-registry.yml # 예제 1 클러스터에 붙이는 Registry
├── pom.xml
└── src/main/
├── avro/
│ ├── UserSignup-v1.avsc # 최초 스키마
│ └── UserSignup.avsc # 진화한 스키마 (v2, 빌드 대상)
└── java/com/example/avro/
├── SignupProducer.java
├── SignupConsumer.java
└── Main.java
Schema Registry 컨테이너
# 예제 1의 kafka-lab 네트워크에 Schema Registry 를 붙입니다.
# 실행: docker compose -f docker-compose.schema-registry.yml up -d
---
name: kafka-lab-sr
services:
schema-registry:
image: confluentinc/cp-schema-registry:8.2.2
hostname: schema-registry
container_name: schema-registry
ports:
- '8081:8081'
environment:
SCHEMA_REGISTRY_HOST_NAME: schema-registry
# 컨테이너 네트워크 안에서 브로커를 찾습니다(예제 1의 PLAINTEXT 리스너).
SCHEMA_REGISTRY_KAFKA_BOOTSTRAP_SERVERS: 'PLAINTEXT://kafka-1:19092,PLAINTEXT://kafka-2:19092,PLAINTEXT://kafka-3:19092'
SCHEMA_REGISTRY_LISTENERS: 'http://0.0.0.0:8081'
# 스키마는 _schemas 라는 compact 토픽에 저장됩니다.
# 브로커가 3대이므로 RF 3 으로 둡니다(1이면 Registry 전체가 단일 장애점).
SCHEMA_REGISTRY_KAFKASTORE_TOPIC: '_schemas'
SCHEMA_REGISTRY_KAFKASTORE_TOPIC_REPLICATION_FACTOR: 3
# 전역 기본 호환성 모드. Schema Registry 의 기본값은 backward 입니다.
# 명시해 의도를 남깁니다(subject 별로 덮어쓸 수 있습니다).
SCHEMA_REGISTRY_SCHEMA_COMPATIBILITY_LEVEL: 'backward'
networks:
- kafka-lab_default
healthcheck:
test: ['CMD-SHELL', 'curl -sf http://localhost:8081/subjects || exit 1']
interval: 10s
timeout: 5s
retries: 12
start_period: 30s
networks:
# 예제 1의 compose 가 만든 네트워크를 재사용합니다.
# 이름은 "{프로젝트명}_default" 규칙을 따릅니다(예제 1의 name: kafka-lab).
kafka-lab_default:
external: true
스키마 v1
{
"type": "record",
"name": "UserSignup",
"namespace": "com.example.avro.model",
"doc": "회원 가입 이벤트 v1",
"fields": [
{ "name": "userId", "type": "string", "doc": "회원 식별자" },
{ "name": "email", "type": "string" },
{ "name": "signupAt", "type": "long",
"doc": "가입 시각 (epoch millis)" },
{ "name": "plan",
"type": { "type": "enum", "name": "Plan",
"symbols": ["FREE", "PRO", "ENTERPRISE"] },
"default": "FREE" }
]
}
스키마 v2 — BACKWARD 호환으로 필드 추가
{
"type": "record",
"name": "UserSignup",
"namespace": "com.example.avro.model",
"doc": "회원 가입 이벤트 v2 — referralSource / marketingOptIn 추가",
"fields": [
{ "name": "userId", "type": "string", "doc": "회원 식별자" },
{ "name": "email", "type": "string" },
{ "name": "signupAt", "type": "long",
"doc": "가입 시각 (epoch millis). v2 에서도 이름을 바꾸지 않습니다 — 이름 변경은 BACKWARD 위반입니다" },
{ "name": "plan",
"type": { "type": "enum", "name": "Plan",
"symbols": ["FREE", "PRO", "ENTERPRISE"] },
"default": "FREE" },
{ "name": "referralSource",
"type": ["null", "string"],
"default": null,
"doc": "유입 채널. BACKWARD 호환의 핵심: 새 필드에는 반드시 default 가 있어야 합니다. default 가 없으면 새 스키마로 읽는 컨슈머가 v1 데이터를 해석할 수 없습니다" },
{ "name": "marketingOptIn",
"type": "boolean",
"default": false,
"doc": "마케팅 수신 동의. null 을 허용하지 않아도 default 가 있으면 BACKWARD 호환입니다" }
]
}
pom.xml
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.example</groupId>
<artifactId>avro-evolution</artifactId>
<version>1.0.0</version>
<properties>
<maven.compiler.release>17</maven.compiler.release>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<kafka.version>4.3.1</kafka.version>
<avro.version>1.12.1</avro.version>
<!-- Schema Registry Docker 이미지 태그와 같은 버전으로 맞춥니다.
Confluent 는 Registry 와 serializer 를 같은 버전으로 릴리스합니다. -->
<confluent.version>8.2.2</confluent.version>
<slf4j.version>2.0.17</slf4j.version>
</properties>
<repositories>
<!-- io.confluent 아티팩트는 Maven Central 에 없습니다.
이 저장소를 추가하지 않으면 kafka-avro-serializer 를 받을 수 없습니다. -->
<repository>
<id>confluent</id>
<url>https://packages.confluent.io/maven/</url>
</repository>
</repositories>
<dependencies>
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>${kafka.version}</version>
</dependency>
<dependency>
<groupId>org.apache.avro</groupId>
<artifactId>avro</artifactId>
<version>${avro.version}</version>
</dependency>
<!-- KafkaAvroSerializer / KafkaAvroDeserializer -->
<dependency>
<groupId>io.confluent</groupId>
<artifactId>kafka-avro-serializer</artifactId>
<version>${confluent.version}</version>
</dependency>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-simple</artifactId>
<version>${slf4j.version}</version>
</dependency>
</dependencies>
<build>
<plugins>
<!-- .avsc → SpecificRecord Java 클래스 생성 -->
<plugin>
<groupId>org.apache.avro</groupId>
<artifactId>avro-maven-plugin</artifactId>
<version>${avro.version}</version>
<executions>
<execution>
<phase>generate-sources</phase>
<goals><goal>schema</goal></goals>
<configuration>
<sourceDirectory>${project.basedir}/src/main/avro</sourceDirectory>
<outputDirectory>${project.build.directory}/generated-sources/avro</outputDirectory>
<!-- v1 은 참고용이라 생성 대상에서 제외합니다
(같은 클래스명이 두 번 생성되면 충돌합니다) -->
<excludes>
<exclude>**/UserSignup-v1.avsc</exclude>
</excludes>
<!-- java.time 타입으로 매핑합니다(logicalType 사용 시) -->
<stringType>String</stringType>
</configuration>
</execution>
</executions>
</plugin>
<plugin>
<groupId>org.codehaus.mojo</groupId>
<artifactId>exec-maven-plugin</artifactId>
<version>3.5.0</version>
<configuration>
<mainClass>com.example.avro.Main</mainClass>
</configuration>
</plugin>
</plugins>
</build>
</project>
프로듀서
package com.example.avro;
import com.example.avro.model.Plan;
import com.example.avro.model.UserSignup;
import io.confluent.kafka.serializers.AbstractKafkaSchemaSerDeConfig;
import io.confluent.kafka.serializers.KafkaAvroSerializer;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.errors.SerializationException;
import org.apache.kafka.common.serialization.StringSerializer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.time.Duration;
import java.time.Instant;
import java.util.Properties;
public class SignupProducer implements AutoCloseable {
private static final Logger log = LoggerFactory.getLogger(SignupProducer.class);
private final Producer<String, UserSignup> producer;
private final String topic;
public SignupProducer(String bootstrapServers, String schemaRegistryUrl,
String topic, boolean autoRegister) {
this.topic = topic;
Properties p = new Properties();
p.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
p.put(ProducerConfig.CLIENT_ID_CONFIG, "signup-producer");
p.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
p.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, KafkaAvroSerializer.class.getName());
// 내구성 — 예제 3과 같은 조합입니다.
p.put(ProducerConfig.ACKS_CONFIG, "all");
p.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
p.put(ProducerConfig.LINGER_MS_CONFIG, 20);
// --- Schema Registry 설정 -------------------------------------------
// Registry 주소. 여러 개를 쉼표로 적으면 장애 시 다른 노드를 씁니다.
p.put(AbstractKafkaSchemaSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, schemaRegistryUrl);
// auto.register.schemas 의 기본값은 true 입니다.
// 개발 환경에서는 편하지만 프로덕션에서는 위험합니다 —
// 애플리케이션이 배포되는 순간 새 스키마가 조용히 등록되어,
// 리뷰 없이 계약이 바뀝니다.
// false 로 두면 CI 단계에서 등록·호환성 검사를 강제할 수 있습니다.
p.put(AbstractKafkaSchemaSerDeConfig.AUTO_REGISTER_SCHEMAS, autoRegister);
// auto.register.schemas=false 일 때, 클라이언트가 이미 등록된
// 최신 버전을 조회해 쓰도록 합니다. 기본값 false.
// 이 값이 false 이고 스키마가 등록돼 있지 않으면 직렬화가 실패합니다.
p.put(AbstractKafkaSchemaSerDeConfig.USE_LATEST_VERSION, !autoRegister);
// subject 이름 결정 방식.
// Schema Registry 8.x 의 기본값은 AssociatedNameStrategy 이며,
// Registry 에 연관 정보가 없으면 TopicNameStrategy 로 폴백합니다.
// 동작을 예측 가능하게 만들려면 명시하는 편이 안전합니다.
// TopicNameStrategy → subject 는 "{topic}-value" / "{topic}-key"
p.put(AbstractKafkaSchemaSerDeConfig.VALUE_SUBJECT_NAME_STRATEGY,
io.confluent.kafka.serializers.subject.TopicNameStrategy.class.getName());
this.producer = new KafkaProducer<>(p);
}
public void send(String userId, String email, Plan plan, String referralSource) {
// Avro 생성 클래스는 builder 를 제공합니다.
// default 가 있는 필드는 생략할 수 있습니다.
UserSignup event = UserSignup.newBuilder()
.setUserId(userId)
.setEmail(email)
.setSignupAt(Instant.now().toEpochMilli())
.setPlan(plan)
.setReferralSource(referralSource) // nullable
.setMarketingOptIn(false)
.build();
try {
producer.send(new ProducerRecord<>(topic, userId, event), (md, ex) -> {
if (ex != null) {
log.error("발행 실패 userId={}", userId, ex);
return;
}
log.info("발행 성공 userId={} → {}-{}@{}",
userId, md.topic(), md.partition(), md.offset());
});
} catch (SerializationException e) {
// 스키마 등록 실패 또는 호환성 위반이 여기로 옵니다.
// 재시도해도 결과가 같으므로 배포를 되돌려야 합니다.
log.error("직렬화 실패 — 스키마 호환성 또는 Registry 접근을 확인하세요. userId={}",
userId, e);
throw e;
}
}
@Override
public void close() {
producer.flush();
producer.close(Duration.ofSeconds(130));
}
}
컨슈머
package com.example.avro;
import com.example.avro.model.UserSignup;
import io.confluent.kafka.serializers.AbstractKafkaSchemaSerDeConfig;
import io.confluent.kafka.serializers.KafkaAvroDeserializer;
import io.confluent.kafka.serializers.KafkaAvroDeserializerConfig;
import org.apache.kafka.clients.consumer.CloseOptions;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.errors.WakeupException;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.time.Duration;
import java.util.List;
import java.util.Properties;
import java.util.concurrent.atomic.AtomicBoolean;
public class SignupConsumer implements Runnable {
private static final Logger log = LoggerFactory.getLogger(SignupConsumer.class);
private final org.apache.kafka.clients.consumer.Consumer<String, UserSignup> consumer;
private final String topic;
private final AtomicBoolean running = new AtomicBoolean(true);
public SignupConsumer(String bootstrapServers, String schemaRegistryUrl,
String topic, String groupId) {
this.topic = topic;
Properties p = new Properties();
p.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
p.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
p.put(ConsumerConfig.CLIENT_ID_CONFIG, groupId + "-1");
p.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
p.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, KafkaAvroDeserializer.class.getName());
p.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
p.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
// --- Schema Registry 설정 -------------------------------------------
p.put(AbstractKafkaSchemaSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, schemaRegistryUrl);
// specific.avro.reader 의 기본값은 false 이고, 그때는 GenericRecord 를 돌려줍니다.
// true 로 두면 avro-maven-plugin 이 생성한 SpecificRecord 클래스로 역직렬화합니다.
// 타입 안전을 얻는 대신, 컨슈머가 자기 스키마(reader schema)로 읽게 되므로
// 호환성 규칙이 실제로 작동하는지가 중요해집니다.
p.put(KafkaAvroDeserializerConfig.SPECIFIC_AVRO_READER_CONFIG, true);
this.consumer = new KafkaConsumer<>(p);
}
@Override
public void run() {
consumer.subscribe(List.of(topic));
try {
while (running.get()) {
ConsumerRecords<String, UserSignup> records =
consumer.poll(Duration.ofSeconds(1));
if (records.isEmpty()) {
continue;
}
for (ConsumerRecord<String, UserSignup> record : records) {
UserSignup event = record.value();
// v1 데이터를 v2 스키마로 읽으면 referralSource 는 default(null),
// marketingOptIn 은 default(false) 로 채워집니다.
// 이것이 BACKWARD 호환이 하는 일입니다.
log.info("수신 {}-{}@{} userId={} plan={} referralSource={} marketingOptIn={}",
record.topic(), record.partition(), record.offset(),
event.getUserId(), event.getPlan(),
event.getReferralSource(), event.getMarketingOptIn());
}
// 처리 후 커밋 → at-least-once (예제 4 참고)
consumer.commitSync(records.nextOffsets(), Duration.ofSeconds(15));
}
} catch (WakeupException e) {
log.info("종료 신호 수신");
} finally {
consumer.close(CloseOptions.timeout(Duration.ofSeconds(30)));
}
}
public void shutdown() {
running.set(false);
consumer.wakeup();
}
}
package com.example.avro;
import com.example.avro.model.Plan;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public final class Main {
private static final Logger log = LoggerFactory.getLogger(Main.class);
private static final String BOOTSTRAP =
"localhost:29092,localhost:39092,localhost:49092";
private static final String SCHEMA_REGISTRY = "http://localhost:8081";
private static final String TOPIC = "users.signup";
public static void main(String[] args) throws InterruptedException {
String mode = args.length > 0 ? args[0] : "produce";
if ("consume".equals(mode)) {
String groupId = args.length > 1 ? args[1] : "signup-analytics";
SignupConsumer consumer =
new SignupConsumer(BOOTSTRAP, SCHEMA_REGISTRY, TOPIC, groupId);
Thread t = new Thread(consumer, "signup-consumer");
Runtime.getRuntime().addShutdownHook(new Thread(consumer::shutdown));
t.start();
t.join();
return;
}
// 프로덕션에서는 autoRegister=false 로 두고 CI 에서 등록합니다.
boolean autoRegister = !(args.length > 1 && "no-auto-register".equals(args[1]));
log.info("auto.register.schemas = {}", autoRegister);
try (SignupProducer producer =
new SignupProducer(BOOTSTRAP, SCHEMA_REGISTRY, TOPIC, autoRegister)) {
producer.send("U-2001", "a@example.com", Plan.PRO, "google-ads");
producer.send("U-2002", "b@example.com", Plan.FREE, null);
producer.send("U-2003", "c@example.com", Plan.ENTERPRISE, "partner-referral");
}
log.info("발행 완료");
}
}
실행 방법
# 0. 예제 1의 클러스터가 떠 있어야 합니다.
cd kafka-lab && docker compose ps
# 1. 토픽 생성
./kcli kafka-topics.sh --create --if-not-exists \
--topic users.signup --partitions 3 --replication-factor 3 \
--config min.insync.replicas=2
# 2. Schema Registry 기동 (예제 1 네트워크에 붙습니다)
cd ../avro-evolution
docker compose -f docker-compose.schema-registry.yml up -d
curl -s http://localhost:8081/subjects # → [] 가 나오면 준비 완료
# 3. v1 스키마를 먼저 등록합니다 (진화를 실습하기 위해 수동 등록)
# jq 로 .avsc 를 문자열로 감싸 REST 로 보냅니다.
curl -s -X POST http://localhost:8081/subjects/users.signup-value/versions \
-H 'Content-Type: application/vnd.schemaregistry.v1+json' \
-d "$(jq -Rs '{schema: .}' < src/main/avro/UserSignup-v1.avsc)"
# → {"id":1}
# 4. v2 가 BACKWARD 호환인지 "등록하지 않고" 먼저 검사합니다
curl -s -X POST \
http://localhost:8081/compatibility/subjects/users.signup-value/versions/latest \
-H 'Content-Type: application/vnd.schemaregistry.v1+json' \
-d "$(jq -Rs '{schema: .}' < src/main/avro/UserSignup.avsc)" | jq
# → {"is_compatible": true}
# 5. 빌드 (avro-maven-plugin 이 UserSignup.java 를 생성합니다)
mvn -q clean package
# 6. 컨슈머를 "먼저" 띄웁니다. BACKWARD 이므로 이 순서가 중요합니다.
mvn -q exec:java -Dexec.args="consume signup-analytics" &
# 7. 프로듀서 실행 (v2 스키마로 발행)
mvn -q exec:java -Dexec.args="produce"
검증 방법
1. subject와 버전이 만들어졌는가
# subject 목록
curl -s http://localhost:8081/subjects | jq
# → ["users.signup-value"]
# 이 subject 의 버전 목록
curl -s http://localhost:8081/subjects/users.signup-value/versions | jq
# → [1,2]
# 최신 버전의 스키마
curl -s http://localhost:8081/subjects/users.signup-value/versions/latest | jq
# subject 별 호환성 모드 (설정하지 않으면 전역 기본값을 씁니다)
curl -s http://localhost:8081/config/users.signup-value | jq
curl -s http://localhost:8081/config | jq
# → {"compatibilityLevel":"BACKWARD"}
subject 이름이 users.signup-value인 것이
TopicNameStrategy가 동작한 결과입니다 —
{topic}-value 규칙입니다.
키에도 스키마를 쓴다면 users.signup-key가 함께 생깁니다.
2. wire format을 바이트로 확인
Avro 메시지의 앞 5바이트를 직접 보면 스키마가 메시지에 없다는 것을 확인할 수 있습니다.
cd ../kafka-lab
# 세그먼트 파일을 직접 덤프합니다.
./kcli kafka-dump-log.sh \
--files /var/lib/kafka/data/users.signup-0/00000000000000000000.log \
--print-data-log --value-decoder-class kafka.serializer.DefaultDecoder \
| head -20
바이트 0 : 0x00 magic byte (항상 0)
바이트 1 ~ 4 : 0x00000001 schema id (4바이트 big-endian int) → 여기서는 1
바이트 5 ~ : ... Avro binary payload (필드 이름 없음, 순서와 타입만)
magic byte가 0x00이 아니면 그 메시지는
Schema Registry 형식이 아닙니다 —
StringSerializer로 쓴 레코드를 KafkaAvroDeserializer로 읽으려 하면
여기서 실패합니다. 흔한 사고 경로입니다.
3. v1 데이터를 v2 컨슈머가 읽을 수 있는가 (BACKWARD의 실제 효과)
# Registry 의 v1 스키마 id 를 확인합니다.
curl -s http://localhost:8081/subjects/users.signup-value/versions/1 | jq '.id'
# Confluent CLI 가 없다면, v1 필드만 채운 이벤트를 발행하는
# 별도 프로듀서를 만들어 확인하는 것이 가장 확실합니다.
# 이미 v1 으로 발행된 레코드가 남아 있다면 그것을 그대로 읽으면 됩니다.
INFO c.e.avro.SignupConsumer - 수신 users.signup-1@0 userId=U-1001 plan=FREE referralSource=null marketingOptIn=false
INFO c.e.avro.SignupConsumer - 수신 users.signup-1@1 userId=U-2001 plan=PRO referralSource=google-ads marketingOptIn=false
첫 줄은 v1로 쓰인 레코드입니다.
referralSource가 데이터에 없으므로 스키마의 default: null이 채워졌습니다.
이 default가 없으면 이 지점에서 역직렬화가 실패합니다 —
그것이 Registry가 등록을 거부하는 이유입니다.
4. 호환되지 않는 변경이 거부되는가
# signupAt → createdAt 으로 이름만 바꾼 스키마를 만듭니다.
sed 's/"signupAt"/"createdAt"/' src/main/avro/UserSignup.avsc > /tmp/UserSignup-bad.avsc
# 등록 전 호환성 검사
curl -s -X POST \
http://localhost:8081/compatibility/subjects/users.signup-value/versions/latest \
-H 'Content-Type: application/vnd.schemaregistry.v1+json' \
-d "$(jq -Rs '{schema: .}' < /tmp/UserSignup-bad.avsc)" | jq
{
"is_compatible": false
}
verbose=truecurl -s -X POST \
'http://localhost:8081/compatibility/subjects/users.signup-value/versions/latest?verbose=true' \
-H 'Content-Type: application/vnd.schemaregistry.v1+json' \
-d "$(jq -Rs '{schema: .}' < /tmp/UserSignup-bad.avsc)" | jq
이 검사를 CI에 넣는 것이 이 예제의 실질적 결론입니다.
머지 전에 /compatibility 엔드포인트로 검사하면
스키마 변경이 프로덕션에 도달하기 전에 막힙니다.
시나리오의 40분 장애는 이 한 줄로 예방됩니다.
호환성 모드 7종과 배포 순서
| 모드 | 보장 | 검사 대상 | 배포 순서 |
|---|---|---|---|
BACKWARD 기본값 |
새 스키마로 직전 버전의 데이터를 읽을 수 있음 | 직전 1개 버전 | 컨슈머 먼저 |
BACKWARD_TRANSITIVE |
새 스키마로 모든 이전 버전의 데이터를 읽을 수 있음 | 전체 버전 | 컨슈머 먼저 |
FORWARD |
직전 버전 스키마로 새 데이터를 읽을 수 있음 | 직전 1개 버전 | 프로듀서 먼저 |
FORWARD_TRANSITIVE |
모든 이전 버전 스키마로 새 데이터를 읽을 수 있음 | 전체 버전 | 프로듀서 먼저 |
FULL |
양방향 (직전 버전과) | 직전 1개 버전 | 순서 무관 |
FULL_TRANSITIVE |
양방향 (모든 버전과) | 전체 버전 | 순서 무관 |
NONE |
검사하지 않음 | — | — (사고가 납니다) |
BACKWARD는 컨슈머 먼저, FORWARD는 프로듀서 먼저.
뒤바꿨을 때 무슨 일이 생기는지
# 이 subject 만 BACKWARD_TRANSITIVE 로 강화합니다.
curl -s -X PUT http://localhost:8081/config/users.signup-value \
-H 'Content-Type: application/vnd.schemaregistry.v1+json' \
-d '{"compatibility": "BACKWARD_TRANSITIVE"}' | jq
# 전역 기본값 변경 (모든 subject 에 영향)
curl -s -X PUT http://localhost:8081/config \
-H 'Content-Type: application/vnd.schemaregistry.v1+json' \
-d '{"compatibility": "FULL_TRANSITIVE"}' | jq
프로덕션 고려사항
| 항목 | 이 예제 | 프로덕션 |
|---|---|---|
auto.register.schemas |
true(기본값)로 실습 |
false. 애플리케이션 배포가 계약 변경을 일으키면 안 됩니다. CI에서 /compatibility 검사 후 등록합니다 |
| Registry 가용성 | 단일 인스턴스 | 2대 이상 + schema.registry.url에 모두 나열. _schemas 토픽 RF 3. Registry가 죽으면 새 스키마의 직렬화·역직렬화가 실패합니다(캐시된 것은 계속 동작) |
| 역직렬화 실패 | 처리하지 않음 | ErrorHandlingDeserializer로 감싸지 않으면 잘못된 레코드가 포이즌 필이 되어 파티션을 영구히 막습니다(예제 6) |
| 호환성 모드 | 전역 BACKWARD |
토픽 성격별로 subject 단위 설정. 보관이 긴 토픽은 BACKWARD_TRANSITIVE, 양방향 배포가 필요하면 FULL |
| 보안 | 인증 없음 | Registry에 인증(basic.auth.credentials.source)과 TLS를 적용하고, 쓰기 권한을 CI로 제한합니다. 인증 없는 Registry는 누구나 계약을 바꿀 수 있습니다 |
| 스키마 소유권 | 애플리케이션 안 .avsc |
스키마 전용 저장소에 두고 리뷰를 거쳐 변경합니다. 프로듀서·컨슈머 팀이 같은 파일을 참조해야 합니다 |
specific.avro.reader |
true |
타입 안전을 얻지만 스키마 변경 시 재빌드가 필요합니다. 스키마를 모르고 라우팅만 하는 서비스는 GenericRecord(기본값 false)가 유리합니다 |
자주 하는 실수
관련 케이스 스터디
이어서 볼 곳
공식 문서 출처
Kafka 설정은 Apache Kafka 4.3.1 문서에서, Schema Registry의 설정명·기본값·호환성 모드 목록은
confluentinc/schema-registry v8.2.2 소스에서 확인했습니다.
- Confluent Schema Registry — 아키텍처, wire format(magic byte + 4바이트 schema id + payload),
_schemas토픽 - Schema Evolution and Compatibility — 호환성 모드 7종(
BACKWARD,BACKWARD_TRANSITIVE,FORWARD,FORWARD_TRANSITIVE,FULL,FULL_TRANSITIVE,NONE), 모드별 배포 순서, 전역 기본값BACKWARD - Schema Registry API Reference —
/subjects,/subjects/{subject}/versions,/compatibility/subjects/{subject}/versions/{version},/config - Serializers, Deserializers —
schema.registry.url,auto.register.schemas(기본true),use.latest.version(기본false),specific.avro.reader(기본false),key.subject.name.strategy/value.subject.name.strategy - Apache Avro Specification — 스키마 해석(resolution) 규칙,
default의 역할, 타입 승격 - Producer Configs —
acks,enable.idempotence,linger.ms - Consumer Configs —
enable.auto.commit,auto.offset.reset