학습 목표

시나리오

회원 도메인이 users.signup 토픽에 가입 이벤트를 발행합니다. 컨슈머는 세 팀입니다 — 마케팅(웰컴 메일), 분석(가입 퍼널), CS(고객 조회).

마케팅 팀이 유입 채널(referralSource) 필드를 추가해 달라고 요청했습니다. JSON을 쓰던 시절에는 프로듀서가 필드를 추가해도 컨슈머들이 무시하기만 하면 되니 문제가 없어 보였습니다. 실제로 사고가 난 것은 필드 이름을 signupAtcreatedAt으로 바꿨을 때였습니다. 세 컨슈머 중 두 개가 NullPointerException으로 죽고, 파티션이 막혀 lag이 40분간 쌓였습니다.

이번에는 스키마를 계약으로 만들고, 호환되지 않는 변경은 배포 전에 거부되게 만듭니다.

아키텍처

Schema Registry 아키텍처 — 스키마는 레지스트리에, 메시지에는 id 만 프로듀서, Schema Registry, Kafka 토픽, 컨슈머의 관계를 6단계로 그린 그림입니다. 1단계로 프로듀서가 스키마를 subject 에 등록하면 2단계로 Schema Registry 가 schema id 를 돌려줍니다. 3단계로 프로듀서는 매직 바이트 1바이트와 schema id 4바이트를 앞에 붙인 레코드를 Kafka 토픽에 씁니다. 스키마 본문은 메시지에 들어가지 않습니다. 4단계로 컨슈머가 그 레코드를 받으면 5단계로 앞쪽 id 를 뽑아 Schema Registry 에 조회하고 6단계로 스키마를 받아 로컬에 캐시합니다. 양쪽 모두 캐시를 쓰기 때문에 메시지마다 REST 호출이 일어나지 않고 새 id 를 처음 볼 때만 조회합니다. Schema Registry 는 스키마를 자신의 내부 토픽 _schemas 에 저장하고 subject 별 버전 관리와 호환성 검사를 담당합니다. Schema Registry — 스키마는 한 번 등록하고, 메시지에는 id 만 실립니다 Schema Registry subject 별 버전 · 호환성 검사 내부 토픽 _schemas 에 저장 ① 스키마 등록 schema id 반환 (정수) ⑤ id 로 조회 ⑥ 스키마 반환 → 캐시 프로듀서 KafkaAvroSerializer 스키마 → id 로컬 캐시 같은 스키마는 다시 등록 안 함 Kafka 토픽 magic(1) + id(4) + payload 스키마 본문은 실리지 않습니다 오버헤드는 레코드당 5바이트 컨슈머 KafkaAvroDeserializer id → 스키마 로컬 캐시 처음 본 id 만 조회 메시지마다 REST 호출이 일어나지 않습니다. 처음 보는 스키마·id 일 때만 왕복하고 이후에는 캐시를 씁니다. 그래서 Schema Registry 가 잠시 죽어도 이미 캐시된 스키마로는 계속 처리됩니다 — 새 스키마 등록·조회만 막힙니다. REST: 등록은 POST /subjects/{subject}/versions · 조회는 GET /schemas/ids/{id} Schema Registry 는 Confluent 컴포넌트입니다. Apache Kafka 배포판에는 포함되지 않습니다.
Schema Registry 아키텍처 — 프로듀서가 스키마를 등록하고 id를 받아 메시지에 넣고, 컨슈머가 id로 스키마를 조회하는 흐름
wire format 바이트 레이아웃 — magic byte 1 + schema id 4 + payload Schema Registry 를 쓰는 직렬화 결과의 바이트 배치를 눈금으로 그린 그림입니다. 0번 바이트는 매직 바이트로 값이 0x00 입니다. 1번부터 4번까지 4바이트는 schema id 를 빅엔디언 정수로 담습니다. 그림의 예에서는 0x00 0x00 0x01 0x2F 이므로 schema id 는 303 입니다. 5번 바이트부터 끝까지가 payload 이며, Avro 바이너리로 인코딩된 값만 들어가고 스키마 본문은 들어가지 않습니다. 따라서 레코드당 고정 오버헤드는 정확히 5바이트입니다. 스키마 JSON 전문을 매 레코드에 넣으면 보통 수백에서 수천 바이트가 되므로, id 4바이트로 대체하는 것이 이 형식의 목적입니다. 참고로 Protobuf 는 schema id 뒤에 메시지 인덱스가 더 붙고, 매직 바이트가 0x01 인 형식은 id 대신 16바이트 GUID 를 담습니다. wire format — 레코드 앞 5바이트가 전부입니다 byte 0 1 2 3 4 5 … 레코드 value 0x00 magic 0x00 0x00 0x01 0x2F schema id = 303 payload — Avro 바이너리로 인코딩된 값 필드 이름도 스키마도 들어 있지 않습니다 1 B 4 B — big-endian 정수 가변 길이 고정 오버헤드 = 5바이트. 값이 1 KB 면 0.5% 도 되지 않습니다. 컨슈머는 앞 5바이트를 떼어 id 를 얻고, 그 id 로 읽기 스키마를 결정합니다. 레코드마다 붙는 스키마 정보의 크기 비교 id 4바이트 4 B 스키마 JSON 전문 수백 ~ 수천 B Protobuf 는 schema id 뒤에 메시지 인덱스가 더 붙습니다. Avro·JSON Schema 는 5바이트로 끝납니다. magic byte 가 0x01 인 형식은 id 4바이트 대신 16바이트 GUID 를 담습니다. 기본 형식은 0x00 입니다.
wire format 바이트 레이아웃 — magic byte 1바이트 + schema id 4바이트 + Avro payload

스키마 자체는 메시지에 들어가지 않습니다. 메시지에는 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 Kafka4.3.1예제 1의 클러스터
Confluent Schema Registry8.2.2Docker 이미지 confluentinc/cp-schema-registry:8.2.2
Apache Avro1.12.1Maven Central. avro + avro-maven-plugin
kafka-avro-serializer8.2.2Schema Registry와 같은 버전. Maven Central이 아니라 packages.confluent.io에서 받습니다
Java17 이상

전체 코드

디렉터리 구조
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 컨테이너

avro-evolution/docker-compose.schema-registry.yml
# 예제 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

src/main/avro/UserSignup-v1.avsc — 최초 버전 (참고용, 빌드 대상 아님)
{
  "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 호환으로 필드 추가

src/main/avro/UserSignup.avsc — 진화한 버전 (빌드 대상)
{
  "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 호환입니다" }
  ]
}
스키마 호환성 모드 매트릭스 — 변경 유형별 허용·거부 가로는 스키마 변경 유형 일곱 가지, 세로는 호환성 모드 일곱 가지인 표입니다. 변경 유형은 기본값 있는 필드 추가, 기본값 없는 필드 추가, 기본값 있던 필드 삭제, 기본값 없던 필드 삭제, 타입 확대 int 에서 long, 타입 축소 long 에서 int, 별칭 없는 필드 이름 변경입니다. 호환성 모드는 BACKWARD, BACKWARD_TRANSITIVE, FORWARD, FORWARD_TRANSITIVE, FULL, FULL_TRANSITIVE, NONE 입니다. BACKWARD 는 새 스키마가 이전 스키마로 쓰인 데이터를 읽을 수 있는지 검사하고, FORWARD 는 새 스키마로 쓴 데이터를 이전 스키마가 읽을 수 있는지 검사하며, FULL 은 둘 다 검사합니다. TRANSITIVE 가 붙으면 직전 버전만이 아니라 모든 이전 버전과 비교합니다. BACKWARD 에서는 기본값 있는 필드 추가, 필드 삭제, 타입 확대가 허용되고 기본값 없는 필드 추가, 타입 축소, 이름 변경은 거부됩니다. FORWARD 에서는 필드 추가와 기본값 있던 필드 삭제, 타입 축소가 허용되고 기본값 없던 필드 삭제, 타입 확대, 이름 변경은 거부됩니다. FULL 은 양쪽 모두 허용인 경우만 허용하므로 기본값 있는 필드 추가와 기본값 있던 필드 삭제만 통과합니다. NONE 은 검사를 하지 않으므로 모두 통과하지만 깨진 조합도 등록됩니다. 이 표는 이전 버전이 하나인 v1 에서 v2 로의 변경을 기준으로 하므로 TRANSITIVE 열의 판정이 같습니다. 컨트롤에서 변경 유형과 호환성 모드를 고르면 그 칸이 강조되고 판정 이유가 표시됩니다. 스키마 호환성 매트릭스 — 이 변경을 이 모드에서 등록할 수 있는가 이 표는 v1 → v2 (이전 버전이 1개) 기준입니다. TRANSITIVE 는 판정 규칙이 아니라 비교 대상 버전 범위가 다릅니다. 호환성 모드 ↓ 변경 유형 → 필드 추가 기본값 O 필드 추가 기본값 X 필드 삭제 기본값 O 필드 삭제 기본값 X 타입 확대 int→long 타입 축소 long→int 필드 이름 변경 BACKWARD 직전 버전과 비교 허용 거부 허용 허용 허용 거부 거부 BACKWARD_TRANSITIVE 모든 이전 버전과 비교 허용 거부 허용 허용 허용 거부 거부 FORWARD 직전 버전과 비교 허용 허용 허용 거부 거부 허용 거부 FORWARD_TRANSITIVE 모든 이전 버전과 비교 허용 허용 허용 거부 거부 허용 거부 FULL 직전 버전 · 양방향 허용 거부 허용 거부 거부 거부 거부 FULL_TRANSITIVE 모든 이전 · 양방향 허용 거부 허용 거부 거부 거부 거부 NONE 검사하지 않음 허용 허용 허용 허용 허용 허용 허용 허용 BACKWARD × 필드 추가 (기본값 O) — 직전 버전과 비교 새 스키마(reader)의 새 필드는 옛 데이터에 없지만 default 로 채워집니다. 업그레이드 순서: 컨슈머를 먼저 올립니다. 허용 — 등록됩니다 거부 — 등록이 실패합니다 NONE — 검사 없이 통과 (보장도 없음) 기준: Avro 해석 규칙 — 필드는 이름(또는 reader 별칭)으로 대응하고, 못 찾으면 reader 쪽 default 가 있어야 합니다. 타입 승격은 int→long→float→double 과 bytes↔string 만 허용됩니다. 이름 변경은 별칭(alias)이나 default 가 있으면 결과가 달라집니다.
호환성 모드 매트릭스 — 변경 유형(필드 추가/삭제/타입 변경/이름 변경) × 모드별 허용 여부

pom.xml

avro-evolution/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>

프로듀서

src/main/java/com/example/avro/SignupProducer.java
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));
    }
}

컨슈머

src/main/java/com/example/avro/SignupConsumer.java
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();
    }
}
src/main/java/com/example/avro/Main.java
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와 버전이 만들어졌는가

Registry REST API 조회
# 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가 함께 생깁니다.

subject naming strategy — Topic · Record · TopicRecord Schema Registry 가 스키마를 등록할 subject 이름을 정하는 세 가지 전략을 비교한 표입니다. 예시는 토픽 orders 와 레코드 타입 com.ex.OrderCreated, com.ex.OrderCancelled 입니다. TopicNameStrategy 는 기본 동작으로 subject 가 토픽 이름에 -key 또는 -value 를 붙인 형태입니다. 예를 들어 orders-value 가 되며, 한 토픽의 value 에는 사실상 스키마 하나만 둘 수 있습니다. RecordNameStrategy 는 subject 가 레코드의 전체 이름입니다. 예를 들어 com.ex.OrderCreated 가 되며, 여러 토픽이 같은 스키마를 공유하고 토픽과 무관하게 전역으로 진화합니다. TopicRecordNameStrategy 는 subject 가 토픽 이름과 레코드 전체 이름을 이은 형태입니다. 예를 들어 orders-com.ex.OrderCreated 가 되며, 한 토픽에 여러 이벤트 타입을 담으면서 토픽별로 독립적으로 진화시킬 수 있습니다. 설정 키는 key.subject.name.strategy 와 value.subject.name.strategy 입니다. 호환성 검사는 subject 단위로 이루어지므로 전략을 바꾸면 subject 이름이 바뀌어 기존 진화 이력과 끊깁니다. subject naming strategy — 호환성 검사는 subject 단위로 이루어집니다 예시 — 토픽 orders · 레코드 타입 com.ex.OrderCreated, com.ex.OrderCancelled 전략 subject 이름 언제 쓰는가 TopicName Strategy 기본 동작 {topic}-key {topic}-value → orders-value 한 토픽 = 한 스키마. 대부분의 경우 이걸로 충분합니다. 한 토픽에 여러 타입을 섞기 어렵습니다. RecordName Strategy {레코드 전체 이름} → com.ex.OrderCreated 토픽 이름이 들어가지 않습니다 여러 토픽이 같은 스키마를 공유할 때. 타입 단위로 전역 진화합니다. 한 토픽에 여러 타입을 담을 수 있습니다. TopicRecord NameStrategy {topic}-{레코드 이름} → orders-com.ex. OrderCreated 한 토픽에 여러 이벤트 타입을 담고 토픽별로 독립 진화시킬 때. 가장 유연하지만 subject 수가 늘어납니다. key.subject.name.strategy value.subject.name.strategy 기본: {topic}-value 전략을 바꾸면 subject 이름이 바뀌어 기존 진화 이력과 끊깁니다 — 운영 중 변경은 신중히. 클래스 이름은 Schema Registry 버전에 따라 다를 수 있습니다. 최신 배포판의 기본 클래스도 조회 결과가 없으면 TopicNameStrategy 로 동작합니다.
subject naming strategy 3종 — Topic / Record / TopicRecord가 만드는 subject 이름의 차이

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
wire format 구조
바이트 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의 실제 효과)

v1 스키마로 레코드를 하나 발행한 뒤, v2 컨슈머로 읽습니다
# Registry 의 v1 스키마 id 를 확인합니다.
curl -s http://localhost:8081/subjects/users.signup-value/versions/1 | jq '.id'

# Confluent CLI 가 없다면, v1 필드만 채운 이벤트를 발행하는
# 별도 프로듀서를 만들어 확인하는 것이 가장 확실합니다.
# 이미 v1 으로 발행된 레코드가 남아 있다면 그것을 그대로 읽으면 됩니다.
기대 컨슈머 로그 — v1 데이터가 default로 채워집니다
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=true
curl -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종과 배포 순서

Schema Registry 호환성 모드 (REST /configcompatibilityLevel에 설정)
모드 보장 검사 대상 배포 순서
BACKWARD 기본값 새 스키마로 직전 버전의 데이터를 읽을 수 있음 직전 1개 버전 컨슈머 먼저
BACKWARD_TRANSITIVE 새 스키마로 모든 이전 버전의 데이터를 읽을 수 있음 전체 버전 컨슈머 먼저
FORWARD 직전 버전 스키마로 새 데이터를 읽을 수 있음 직전 1개 버전 프로듀서 먼저
FORWARD_TRANSITIVE 모든 이전 버전 스키마로 새 데이터를 읽을 수 있음 전체 버전 프로듀서 먼저
FULL 양방향 (직전 버전과) 직전 1개 버전 순서 무관
FULL_TRANSITIVE 양방향 (모든 버전과) 전체 버전 순서 무관
NONE 검사하지 않음 — (사고가 납니다)
스키마 배포 순서 — BACKWARD 는 컨슈머 먼저, FORWARD 는 프로듀서 먼저 호환성 모드에 따른 배포 순서를 좌우로 비교한 그림입니다. 왼쪽 BACKWARD 는 새 스키마가 이전 스키마로 쓰인 데이터를 읽을 수 있음을 보장합니다. 허용되는 변경은 필드 삭제와 기본값 있는 필드 추가입니다. 따라서 새 스키마를 등록한 뒤 컨슈머를 먼저 새 스키마로 배포하고, 그 다음에 프로듀서를 배포합니다. 순서를 뒤바꿔 프로듀서를 먼저 올리면 새 데이터에서 필드가 사라졌는데 아직 옛 스키마를 쓰는 컨슈머가 그 필드를 요구하므로 역직렬화가 실패합니다. 오른쪽 FORWARD 는 새 스키마로 쓴 데이터를 이전 스키마가 읽을 수 있음을 보장합니다. 허용되는 변경은 필드 추가와 기본값 있던 필드 삭제입니다. 따라서 프로듀서를 먼저 새 스키마로 배포하고, 그 다음에 컨슈머를 배포합니다. 순서를 뒤바꿔 컨슈머를 먼저 올리면 새 컨슈머가 기본값 없는 새 필드를 요구하는데 아직 옛 스키마를 쓰는 프로듀서는 그 필드를 쓰지 않으므로 새 컨슈머가 실패합니다. FULL 과 FULL_TRANSITIVE 는 양방향이 보장되므로 순서가 자유롭고, NONE 은 어느 순서로 올려도 깨질 수 있습니다. 배포 순서 — 무엇을 먼저 올리는지가 장애를 결정합니다 BACKWARD 새 스키마가 옛 데이터를 읽음 허용: 필드 삭제 · 기본값 있는 필드 추가 ① 새 스키마를 레지스트리에 등록 ② 컨슈머를 먼저 새 스키마로 배포 ③ 그 다음 프로듀서를 배포 뒤바꾸면 — 프로듀서를 먼저 올리면 새 데이터에서 필드가 사라졌는데 옛 스키마 컨슈머는 그 필드를 요구합니다. → 구 컨슈머 역직렬화 실패 FORWARD 새 데이터를 옛 스키마가 읽음 허용: 필드 추가 · 기본값 있던 필드 삭제 ① 새 스키마를 레지스트리에 등록 ② 프로듀서를 먼저 새 스키마로 배포 ③ 그 다음 컨슈머를 배포 뒤바꾸면 — 컨슈머를 먼저 올리면 새 컨슈머가 기본값 없는 새 필드를 요구하는데 옛 프로듀서는 아직 그 필드를 쓰지 않습니다. → 새 컨슈머 역직렬화 실패 외우는 방법: BACKWARD = 컨슈머(뒤에서 읽는 쪽) 먼저 · FORWARD = 프로듀서(앞에서 쓰는 쪽) 먼저 FULL · FULL_TRANSITIVE 는 양방향이므로 순서가 자유롭습니다. NONE 은 어느 순서로 올려도 깨질 수 있습니다.
배포 순서 — BACKWARD는 컨슈머 먼저, FORWARD는 프로듀서 먼저. 뒤바꿨을 때 무슨 일이 생기는지
subject별 호환성 모드 변경
# 이 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 소스에서 확인했습니다.