Kafka 스키마 레지스트리 운영기: Avro 스키마 진화, 호환성 전략, Consumer 배포 순서까지 깨지지 않는 파이프라인 설계
핵심 요약
이 글에서 확인할 내용
스키마 계약 없이 Kafka를 운영하면 어떤 일이 생기는가 이벤트 기반 파이프라인에서 Producer와 Consumer는 서로 다른 팀이 관리하는 경우가 대부분입니다. 단일 서버 배포로 계약을 강제할 수 있는 REST API와 달리, Kafka 메시지는 토픽에 영구 저장되기 때문에 스키마 변경 하나가 수십 개의 Consumer 서비스를 동시에 깨뜨릴 수 있습니다.
- Schema Registry의 역할과 메시지 와이어 포맷
- 스키마 호환성 모드: BACKWARD, FORWARD, FULL
- 배포 순서 전략: 언제 Consumer를 먼저 올리는가
스키마 계약 없이 Kafka를 운영하면 어떤 일이 생기는가
이벤트 기반 파이프라인에서 Producer와 Consumer는 서로 다른 팀이 관리하는 경우가 대부분입니다. 단일 서버 배포로 계약을 강제할 수 있는 REST API와 달리, Kafka 메시지는 토픽에 영구 저장되기 때문에 스키마 변경 하나가 수십 개의 Consumer 서비스를 동시에 깨뜨릴 수 있습니다.
실제 결제 이벤트 파이프라인 운영 사례를 보면, 하루 수천만 건의 메시지를 처리하는 환경에서 특정 팀이 payment_amount 필드를 amount로 이름을 바꾼 뒤 세 개의 Consumer 서비스가 동시에 역직렬화 오류를 내기 시작했습니다. 오류는 배포 직후가 아니라 한두 시간이 지나서야 DLQ(Dead Letter Queue) 급증으로 감지됐고, 밀린 메시지 재처리까지 포함하면 장애 복구에 수 시간이 소요됐습니다.
이 글은 그런 장애를 예방하기 위해 도입하는 Confluent Schema Registry 운영 체계를 다룹니다. 스키마 ID 직렬화 포맷의 내부 구조부터 BACKWARD·FORWARD·FULL 호환성 모드별 진화 가능 범위, Producer-먼저 vs Consumer-먼저 배포 순서가 파이프라인에 미치는 영향, Python과 Java 직렬화 코드 비교, 운영 중 스키마 드리프트 감지와 삭제 사고 방지까지 순서대로 설명합니다. 단순 개념 설명보다는 "언제 어느 설정을 선택하고, 배포 순서를 어떻게 결정하는가"에 집중합니다.
파이프라인 전반의 Consumer 안정성은 Kafka Consumer Lag 운영 가이드와 함께 읽으면 더 완성된 그림이 됩니다. Kafka를 이벤트 버스로 쓰는 전체 아키텍처 설계는 고가용성 이벤트 기반 아키텍처(EDA) 설계에서 다룹니다.
Schema Registry의 역할과 메시지 와이어 포맷
Schema Registry가 해결하는 문제
Kafka 메시지 자체에는 타입 정보가 없습니다. 바이트 배열을 어떻게 해석할지는 Producer와 Consumer 사이의 암묵적 합의에 달려 있습니다. 이 합의를 코드가 아니라 외부 레지스트리로 명시적으로 관리하는 것이 Schema Registry의 핵심입니다.
Schema Registry는 다음 세 가지를 제공합니다.
- 스키마 저장소: Avro, Protobuf, JSON Schema를 버전 단위로 저장합니다.
- 호환성 검증: 새 스키마를 등록할 때 기존 버전과의 호환성을 자동으로 검사합니다.
- REST API: Producer와 Consumer가 런타임에 스키마 ID를 조회하고 캐시합니다.
Confluent 와이어 포맷 (5바이트 헤더)
Confluent Schema Registry를 사용하는 메시지는 Avro 직렬화 페이로드 앞에 5바이트 헤더를 붙입니다(공식 문서 참고).
┌──────────────┬────────────────────────┬─────────────────────────┐
│ Byte 0 │ Bytes 1–4 │ Bytes 5–N │
│ Magic Byte │ Schema ID (Big-Endian)│ Avro-serialized payload│
│ 0x00 │ uint32 │ binary │
└──────────────┴────────────────────────┴─────────────────────────┘
- Magic Byte (0x00): Confluent 포맷임을 나타내는 식별자입니다. 다른 값이면 Confluent 디시리얼라이저는 처리를 거부합니다.
- Schema ID: 4바이트 빅엔디언 정수로 인코딩된 Schema Registry 내 스키마 고유 번호입니다.
- Payload: 해당 스키마 ID에 대응하는 Avro 스키마로 직렬화된 실제 데이터입니다.
Consumer는 메시지를 받으면 먼저 이 헤더를 읽어 Schema ID를 추출하고, Registry에서 해당 스키마를 가져와(또는 로컬 캐시에서 조회해) 역직렬화합니다. Schema ID를 메시지에 포함시키는 덕분에 Producer와 Consumer가 동일한 스키마 버전을 공유하지 않아도 안전하게 통신할 수 있습니다.
스키마 호환성 모드: BACKWARD, FORWARD, FULL
Confluent Schema Registry의 기본 호환성 모드는 BACKWARD입니다. 호환성 모드는 Subject 단위로 설정하거나 Registry 전체에 기본값을 지정할 수 있습니다(공식 문서).
BACKWARD (기본값)
"새 스키마로 디코딩하는 Consumer가 이전 스키마로 작성된 데이터를 읽을 수 있다."
새 스키마가 이전 스키마의 상위 호환입니다. Consumer를 먼저 배포하는 전략과 궁합이 맞습니다.
허용되는 변경:
- 기본값이 있는 필드 추가
- 선택적 필드 삭제
금지되는 변경:
- 기본값 없는 필드 추가
- 필수 필드 삭제
- 필드 타입 변경
- 필드 이름 변경(rename)
// v1 스키마
{
"type": "record",
"name": "Payment",
"fields": [
{ "name": "payment_id", "type": "string" },
{ "name": "amount", "type": "double" }
]
}
// v2 스키마 — BACKWARD 호환 (기본값 있는 필드 추가)
{
"type": "record",
"name": "Payment",
"fields": [
{ "name": "payment_id", "type": "string" },
{ "name": "amount", "type": "double" },
{ "name": "currency", "type": "string", "default": "KRW" }
]
}
v2 Consumer는 v1 메시지를 받아도 currency 필드에 기본값 "KRW"를 채워 정상 처리합니다.
FORWARD
"이전 스키마로 디코딩하는 Consumer가 새 스키마로 작성된 데이터를 읽을 수 있다."
새 스키마가 이전 스키마의 하위 호환입니다. Producer를 먼저 배포하는 전략과 맞습니다.
허용되는 변경:
- 기본값 없는 필드 추가
- 선택적 필드 삭제
FORWARD는 이전 Consumer가 새 필드를 모르더라도 무시(ignore)하면서 읽을 수 있음을 보장합니다. 단, Consumer가 실제로 모르는 필드를 무시하도록 구현되어 있어야 합니다.
FULL
BACKWARD와 FORWARD를 동시에 만족합니다. 스키마 변경이 양방향으로 호환되어야 하므로 가장 제약이 강합니다.
허용되는 변경:
- 기본값이 있는 선택적 필드 추가
- 기본값이 있는 선택적 필드 삭제
Producer와 Consumer 배포 순서에 무관하게 안전하지만, 필드 추가에 항상 기본값이 필요하기 때문에 스키마 설계 단계부터 이를 고려해야 합니다.
Transitive 변형
BACKWARD_TRANSITIVE, FORWARD_TRANSITIVE, FULL_TRANSITIVE는 직전 버전만이 아니라 모든 이전 버전에 대해 호환성을 검사합니다. 장기 보존 토픽이나 다수의 Consumer 버전이 동시에 운영되는 환경에서는 Transitive 모드가 안전합니다.
| 모드 | 검사 대상 | 배포 순서 권장 |
|---|---|---|
| BACKWARD | 직전 버전 | Consumer 먼저 |
| BACKWARD_TRANSITIVE | 모든 이전 버전 | Consumer 먼저 |
| FORWARD | 직전 버전 | Producer 먼저 |
| FORWARD_TRANSITIVE | 모든 이전 버전 | Producer 먼저 |
| FULL | 직전 버전 | 순서 무관 |
| FULL_TRANSITIVE | 모든 이전 버전 | 순서 무관 |
배포 순서 전략: 언제 Consumer를 먼저 올리는가
스키마 호환성 모드를 이해하더라도 실제 배포 순서를 잘못 잡으면 다운타임이 생깁니다.
BACKWARD 모드에서 Consumer 먼저 배포
시나리오: currency 필드(기본값 "KRW")를 추가하는 v2 스키마 도입.
1. Schema Registry에 v2 스키마 등록 (호환성 검사 통과 확인)
2. Consumer v2 배포 및 안정화 확인
3. Producer v2 배포
Consumer v2는 v1 메시지도 처리할 수 있으므로 단계 2~3 사이에 v1 메시지가 들어와도 안전합니다. Producer를 먼저 배포하면 v1 Consumer가 currency 필드를 모르는 상태에서 v2 메시지를 받아 역직렬화 오류가 날 수 있습니다.
FORWARD 모드에서 Producer 먼저 배포
시나리오: 필드를 제거하는 스키마 변경.
1. Schema Registry에 v2 스키마 등록
2. Producer v2 배포
3. Consumer v2 배포
이전 Consumer는 제거된 필드를 단순히 받지 못하는 것으로 처리하고, 새 Consumer v2가 올라오면 업그레이드가 완료됩니다.
무중단 배포 체크리스트
# 1. 호환성 사전 검증 (등록 전에 테스트)
curl -X POST -H "Content-Type: application/vnd.schemaregistry.v1+json" --data '{"schema": "{"type":"record","name":"Payment","fields":[{"name":"payment_id","type":"string"},{"name":"amount","type":"double"},{"name":"currency","type":"string","default":"KRW"}]}"}' http://localhost:8081/compatibility/subjects/payment-value/versions/latest
# 기대 응답: {"is_compatible":true}
# 2. 스키마 등록
curl -X POST -H "Content-Type: application/vnd.schemaregistry.v1+json" --data '{"schema": "..."}' http://localhost:8081/subjects/payment-value/versions
# 3. 등록된 버전 목록 확인
curl http://localhost:8081/subjects/payment-value/versions
Python: confluent-kafka로 Avro 직렬화 구현
confluent-kafka 라이브러리의 schema_registry 모듈을 사용합니다. 아래 예제는 confluent-kafka 2.x 기준입니다.
Producer
from confluent_kafka import SerializingProducer
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroSerializer
from confluent_kafka.serialization import StringSerializer
SCHEMA_STR = """
{
"type": "record",
"name": "Payment",
"namespace": "com.example",
"fields": [
{"name": "payment_id", "type": "string"},
{"name": "amount", "type": "double"},
{"name": "currency", "type": "string", "default": "KRW"}
]
}
"""
schema_registry_client = SchemaRegistryClient({"url": "http://localhost:8081"})
avro_serializer = AvroSerializer(
schema_registry_client,
SCHEMA_STR,
lambda obj, ctx: obj, # to_dict: dict 그대로 사용
)
producer_conf = {
"bootstrap.servers": "localhost:9092",
"key.serializer": StringSerializer("utf_8"),
"value.serializer": avro_serializer,
}
producer = SerializingProducer(producer_conf)
producer.produce(
topic="payments",
key="pay-001",
value={"payment_id": "pay-001", "amount": 15000.0, "currency": "KRW"},
on_delivery=lambda err, msg: print(f"Delivered: {msg.topic()} [{msg.partition()}]" if not err else f"Error: {err}"),
)
producer.flush()
AvroSerializer의 세 번째 인수는 Python 객체를 dict로 변환하는 함수입니다. 단순 dict를 사용한다면 lambda obj, ctx: obj로 처리합니다. dataclass나 Pydantic 모델을 사용할 경우 .model_dump() 등을 호출하도록 변경하면 됩니다.
Consumer
from confluent_kafka import DeserializingConsumer
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroDeserializer
from confluent_kafka.serialization import StringDeserializer
schema_registry_client = SchemaRegistryClient({"url": "http://localhost:8081"})
avro_deserializer = AvroDeserializer(
schema_registry_client,
# reader_schema를 지정하면 스키마 투영(projection)이 가능합니다.
# 지정하지 않으면 Writer 스키마(메시지에 기록된 스키마 ID)를 그대로 사용합니다.
)
consumer_conf = {
"bootstrap.servers": "localhost:9092",
"group.id": "payment-processor",
"auto.offset.reset": "earliest",
"key.deserializer": StringDeserializer("utf_8"),
"value.deserializer": avro_deserializer,
}
consumer = DeserializingConsumer(consumer_conf)
consumer.subscribe(["payments"])
try:
while True:
msg = consumer.poll(timeout=1.0)
if msg is None:
continue
if msg.error():
print(f"Consumer error: {msg.error()}")
continue
payment = msg.value()
print(f"Received: {payment}")
finally:
consumer.close()
Java: Spring Kafka + KafkaAvroSerializer 구현
Spring Boot 3.x + Spring Kafka 3.x 환경을 기준으로 합니다. io.confluent:kafka-avro-serializer 의존성을 추가해야 합니다.
build.gradle (또는 pom.xml 의존성)
// build.gradle
repositories {
maven { url "https://packages.confluent.io/maven/" }
}
dependencies {
implementation "org.springframework.kafka:spring-kafka"
implementation "io.confluent:kafka-avro-serializer:7.6.0"
}
application.yml
spring:
kafka:
bootstrap-servers: localhost:9092
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: io.confluent.kafka.serializers.KafkaAvroSerializer
properties:
schema.registry.url: http://localhost:8081
consumer:
group-id: payment-processor
auto-offset-reset: earliest
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: io.confluent.kafka.serializers.KafkaAvroDeserializer
properties:
schema.registry.url: http://localhost:8081
specific.avro.reader: true # 생성된 Avro SpecificRecord 클래스 사용 시
Producer 빈
import io.confluent.kafka.serializers.KafkaAvroSerializer;
import org.apache.kafka.common.serialization.StringSerializer;
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.core.ProducerFactory;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.util.HashMap;
import java.util.Map;
@Configuration
public class KafkaProducerConfig {
@Bean
public ProducerFactory<String, Object> producerFactory() {
Map<String, Object> props = new HashMap<>();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", StringSerializer.class);
props.put("value.serializer", KafkaAvroSerializer.class);
props.put("schema.registry.url", "http://localhost:8081");
return new DefaultKafkaProducerFactory<>(props);
}
@Bean
public KafkaTemplate<String, Object> kafkaTemplate() {
return new KafkaTemplate<>(producerFactory());
}
}
Consumer 서비스
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Service;
@Service
public class PaymentConsumerService {
@KafkaListener(topics = "payments", groupId = "payment-processor")
public void consume(Object payload) {
// specific.avro.reader=true 설정 시 payload는 생성된 SpecificRecord 타입으로 캐스팅 가능
System.out.println("Received: " + payload);
}
}
specific.avro.reader=true로 설정하면 Avro 코드 생성기(avro-tools 또는 Gradle avro 플러그인)가 생성한 SpecificRecord 구현체로 자동 캐스팅됩니다. 제네릭 GenericRecord를 사용하려면 이 옵션을 false로 유지하면 됩니다.

REST API로 스키마 관리하기
Schema Registry는 HTTP REST API를 노출합니다. CI/CD 파이프라인이나 운영 스크립트에서 직접 호출할 때 유용합니다.
# 전체 Subject 목록 조회
curl -s http://localhost:8081/subjects | jq .
# 특정 Subject의 버전 목록
curl -s http://localhost:8081/subjects/payment-value/versions | jq .
# 특정 버전의 스키마 상세 조회
curl -s http://localhost:8081/subjects/payment-value/versions/1 | jq .
# Schema ID로 스키마 조회
curl -s http://localhost:8081/schemas/ids/1 | jq .
# 호환성 수준 조회
curl -s http://localhost:8081/config/payment-value | jq .
# Subject 단위 호환성 수준 변경
curl -X PUT -H "Content-Type: application/vnd.schemaregistry.v1+json" --data '{"compatibility": "FULL_TRANSITIVE"}' http://localhost:8081/config/payment-value
# 소프트 삭제 (버전 단위)
curl -X DELETE http://localhost:8081/subjects/payment-value/versions/1
# 하드 삭제 (소프트 삭제 후 permanency 플래그 필요)
curl -X DELETE "http://localhost:8081/subjects/payment-value/versions/1?permanent=true"
스키마 드리프트 감지와 삭제 사고 방지
스키마 드리프트란
운영 환경에서 스키마가 의도치 않게 복수의 버전으로 갈라지는 현상입니다. 개발팀이 각자 Producer를 배포하면서 Registry에 다른 버전의 스키마를 등록하거나, 로컬 개발용 스키마와 프로덕션 스키마가 달라지는 경우에 발생합니다.
드리프트를 조기에 감지하는 방법:
- CI 단계에서 호환성 사전 검사: 스키마 파일을 Git으로 관리하고, PR 시점에
/compatibilityAPI를 호출해 통과해야만 병합을 허용합니다.
# CI 스크립트 예시
SCHEMA=$(cat payment.avsc | python3 -c "import sys,json; print(json.dumps({'schema': sys.stdin.read()}))")
RESULT=$(curl -s -X POST -H "Content-Type: application/vnd.schemaregistry.v1+json" --data "$SCHEMA" http://schema-registry:8081/compatibility/subjects/payment-value/versions/latest)
echo "$RESULT" | jq -e '.is_compatible == true' || exit 1
- Schema Registry UI 또는 모니터링: Confluent Control Center, Kafdrop, Kafka UI 등 관리 도구에서 Subject별 버전 추이를 주기적으로 확인합니다.
삭제 사고 방지
Schema Registry의 삭제는 소프트 삭제와 하드 삭제로 구분됩니다. 소프트 삭제(DELETE /subjects/{subject}/versions/{version})는 해당 버전을 조회 불가로 표시하지만 Schema ID와 실제 데이터는 보존됩니다. 하드 삭제(?permanent=true)는 완전히 제거하며 복구가 불가능합니다.
운영 환경에서는 다음 정책을 권장합니다.
delete.subject.and.compatibility.level권한을 운영 계정에서 분리합니다.- 하드 삭제는 비활성화하거나 별도 승인 프로세스를 거치도록 제한합니다.
- Schema Registry 메타데이터를 주기적으로 백업합니다(Registry 자체가 내부적으로 Kafka 토픽인
_schemas에 데이터를 저장하므로, 해당 토픽의 보존 기간을 충분히 설정합니다).
Avro 스키마 진화 실전 패턴
필드 추가: 항상 기본값을 함께 정의한다
// 권장: 기본값 명시
{ "name": "is_refunded", "type": "boolean", "default": false }
// 비권장: 기본값 없음 — BACKWARD 호환성 위반
{ "name": "is_refunded", "type": "boolean" }
필드 삭제: 단계적으로 진행한다
필드를 즉시 삭제하면 해당 필드에 의존하는 Consumer가 예외를 낼 수 있습니다. 다음 순서를 권장합니다.
1단계: 삭제 예정 필드에 기본값 추가 (아직 Consumer 코드에서 사용)
2단계: Consumer 코드에서 해당 필드 참조 제거 → 배포
3단계: Schema Registry에서 필드 제거 (스키마 v+1 등록)
Union 타입으로 Optional 필드 표현
Avro에서 null 허용 필드는 union 타입으로 표현합니다.
{
"name": "refund_amount",
"type": ["null", "double"],
"default": null
}
null이 union의 첫 번째 타입이어야 default: null이 유효합니다. 순서를 바꾸면 Avro 파서가 오류를 냅니다.
필드 이름 변경: aliases 활용
Avro의 aliases를 활용하면 필드 이름 변경 시 하위 호환성을 유지할 수 있습니다.
// v2 스키마: payment_amount → amount 로 변경
{
"name": "amount",
"type": "double",
"aliases": ["payment_amount"]
}
BACKWARD 호환 읽기 시 Reader 스키마의 aliases를 기준으로 Writer 스키마의 이전 필드명과 매핑합니다. 다만 aliases 활용은 Avro 버전과 라이브러리 구현에 따라 지원 범위가 다를 수 있으므로, 실제 환경에서 반드시 사전 검증이 필요합니다.
성능 고려 사항
스키마 캐시
Producer와 Consumer 모두 내부적으로 Schema ID를 캐시합니다. 매 메시지마다 Registry를 호출하지 않으므로 네트워크 오버헤드는 초기 조회 시에만 발생합니다. 캐시 용량은 max.schemas.per.subject 등 설정으로 조정합니다.
페이로드 크기
Avro 바이너리 직렬화는 JSON 대비 페이로드 크기를 줄여주지만, 실제 절감 폭은 데이터 구조와 필드 수에 따라 달라집니다. 수치적 효과는 실제 운영 데이터를 기준으로 벤치마크를 통해 측정하는 것을 권장합니다. 단순 문자열 중심 데이터에서는 압축 효과가 제한적일 수 있습니다.
Schema Registry 가용성
Schema Registry가 다운되면 초기 스키마 조회에 실패할 수 있습니다. 이미 캐시된 스키마는 계속 동작하지만, 새로운 스키마 ID를 처음 만나는 경우에는 오류가 납니다. 고가용성 구성을 위해 Schema Registry를 다중 인스턴스로 운영하고, 로드 밸런서를 앞에 둡니다.
실무 체크리스트
-
호환성 모드는 Subject 단위로 명시적으로 설정하라. 전역 기본값(BACKWARD)에만 의존하지 말고, 토픽 목적에 따라 FULL_TRANSITIVE 등을 적극 활용하라.
-
모든 스키마 파일을 Git으로 버전 관리하고, CI에서
/compatibilityAPI 검사를 의무화하라. 스키마 변경이 코드 리뷰를 거치도록 강제하는 것이 장애 예방의 가장 효과적인 수단이다. -
배포 순서를 호환성 모드에 맞게 문서화하라. BACKWARD 모드라면 Consumer 먼저, FORWARD 모드라면 Producer 먼저 배포한다는 규칙을 팀 내 런북에 명시한다.
-
하드 삭제 권한을 운영 계정과 분리하고,
_schemas토픽의 retention을compact로 유지하라. 소프트 삭제만으로도 대부분의 정리 목적을 달성할 수 있다. -
스키마 캐시 TTL과 Registry 가용성 SLA를 함께 설계하라. 캐시 덕분에 Registry 순단 시에도 기존 스키마로 계속 처리할 수 있지만, 신규 스키마 배포 타이밍과 Registry 장애 시나리오를 미리 정의해 둬야 한다.