programming language/Kakfa

[Kafka 활용] Kafka 심화 4 - 스키마 레지스트리와 실습하기

공대키메라 2026. 7. 24. 00:28

지난 시간에는 kafka에 보안 관련해서 실습해보고 개념을 정리했다.

(지난 내용 : [Kafka 활용] Kafka 심화 3 - 보안 적용 : Kafka SSL, SASL, ACL)

 

솔직히 보안 적용하는거 굉장히 귀찮았다.

 

물론 엮여있는 기본 개념들을 찾아서 정리하고 하느라 또 늦어진 감이 있긴 한데, 이런게 있구나~ 싶은 정도만 잘 알고 있어도 나중에 쉽게 적용할 수 있지 않을까 싶어서 SSL만 자세히 정리하고 나머지는 요약하는 수준으로 마무리했다.

 

이번 글에서는 카프카에서 어떻게 스키마를 정의하고 활용할 수 있는지, 어느 유명 기업에서 어떻게 적용을 했는지 알아볼 것이다.

 

특히 이번 글의 흐름은 '실전 카프카 개발부터 운영까지' 를 참고했다

 

늘 kafka공부를 하면서 잘 정리된 흐름! 덕분에 열심히 공부하고 있다. S2 강추드립니다.

 

(알라딘 도서 링크 : 실전 카프카 개발부터 운영까지)

 

실전 카프카 개발부터 운영까지 | 고승범

국내 최초이자 유일한 컨플루언트 공인 강사 자격과 공인 관리자 자격을 보유한 『카프카,데이터 플랫폼의 최강자』 저자 고승범이 SKT, 카카오등 국내 최대 규모의 데이터 플랫폼상에서 카프카

www.aladin.co.kr


1. 스키마의 개념과 필요성

왜 스키마가 필요할까? 우선 스키마가 뭘까?

 

DB상에서는 다음과 같이 정의한다.

데이터베이스의 구조와 제약조건에 관해 전반적인 명세를 기술한 것

 

더 넓은 개념에서의 의미는 다음과 같다.

 

데이터가 어떻게 생겼는지 알려주는 데이터의 설계도

 

DB에서 미리 구조와 제약조건에 대해 명세를 기술햇기 때문에 데이터 입력시 잘못된 형태로의 입력을 방지한다. 

 

여태 Kafka로 producer, consumer, kafka cluster, topic, replication 과 docker로의 실행 및 메시지 발행 등 여러개를 진행했다.

 

 

그런데 스키마가 없이 막 데이터를 전달하게 되면 이 데이터가 정합한 데이터인지 확인하기가 힘들다.

 

즉, 스키마를 정의해서 관리하면...

 

(1) 데이터를 컴슘하는 여러 컨슈머에게 그 데이터에 대한 정확한 정의와 의미를 알려줄 수 있다.

 

또한 (2) 실수로 데이터를 잘못된 형태로 특정 토픽에 메시지를 전송하게 되면 시스템이 영향을 받게 되므로 이를 방지한다.

 

이러한 장점을 이용하기 위해서는 데이터에 대한 규약을 함께 준수하기 위해 많은 공수와 시간을 들이기 때문에

 

(3) 중앙 데이터 파이프라인 역할을 하는 카프카에서 수십, 수백의 애프리케이션들이 별다른 영향 없이 스키마를 변경할 수 있다.

 

 

실제로 업무를 해보면 결국에 DB에 어떻게 데이터를 적재하고 사용할 지가 핵심인 것 같다. 

 

아무리 도메인 주도 개발이라고 하지만 검색하고 이를 가공해서 여러 방면으로 활용해야 하는 것은 결국 데이터이다.

 

그렇기 때문에 이 스키마를 정의하는것! 매우 중요하다. 

 

2. Confluent 스키마 레지스트리 

레지스트리가 뭘까? 

 

특정한 정보를 저장하고 관리하는 데이터베이스나 시스템을 레지스트리라고한다. 

 

데이터가 어떻게 생겼는지 알려주는 데이터의 설계도를 스키마라고 했다.

 

즉, 카프카의 스키마 레지스트리는 데이터의 설계도 및 정보를 관리하는 데이터베이스나 시스템을 말한다.

 

 

여기서 시스템이라는 말이 좀 와닿았는데 하나의 시스템을 애플리케이션으로도 볼 수 있다.

 

카프카에서 스키마를 활용하기 위해서는 스키마 레지스트리라는 애플리케이션을 이용하는 것이다.

 

Confluent의 설명을 보도록 하자.

 

내용은 다음 문서에서 있는 것을 가져온것이니 알면 넘어가도 좋다

(개요 : Confluent 플랫폼용 스키마 레지스트리)

(기초 : Confluent 플랫폼의 스키마 레지스트리의 기초 개념)

 

 

스키마 레스트리의 정의를 보자.

 

스키마 레지스트리(Schema Registry)는 스트림 처리뿐만 아니라 데이터베이스, 파일, 기타 정적 데이터 저장소와 같은 저장된 데이터(data at rest)의 데이터 저장 및 교환을 지원하는 관리형 스키마 저장소입니다.

 

스키마 레지스트리는 관리형 스키마 저장소!

출처 : https://docs.confluent.io/platform/current/schema-registry/index.html



그림을 자세히 들여다보자.

 

중앙에는 스키마 레지스트리(Schema Registry)가 있다. 

 

스키마 생성자가 스키마를 등록 하거나 고도화한다.

 

Producer와 Consumer들은 스키마 레지스트리로부터 스키마를 가져와서 사용하고,

 

Kafka는 Producer의 데이터가 들어오고 Consumer가 데이터를 pulling할 때 스키마 레지스트리로부터 검증을 진행한다.

 

또한 스키마 레지스트리는 스키마 레지스트리끼리 스키마를 공유할 수 있다.

 

흐름은 그렇지만 정말 그렇게 설명이 되어있는지 보자.

 

Schema Registry는 데이터 처리 및 직렬화(바이너리 형식으로의 변환 및 복원)에 사용되는 스키마를 관리하고 검증하기 위한 중앙 집중식 저장소를 제공합니다.

 

 

스키마의 필요성을 해결해주는 부분이 많긴 한데, 스키마 레지스트리를 사용하면 다음과 같은 흔한 데이터 문제를 해결할 수 있다고 설명한다. 

 

데이터 불일치 (Data inconsistency)

레지스트리는 모든 시스템 데이터가 합의된 스키마를 따르도록 보장합니다. 이는 데이터 불일치 위험을 줄이고 데이터 품질을 향상시킵니다.

 

호환되지 않는 데이터 형식 (Incompatible data formats)

여러 데이터 생산자(Producer)와 소비자(Consumer)가 존재하는 환경에서는 애플리케이션마다 서로 다른 데이터 형식을 사용할 수 있습니다. 스키마 레지스트리는 중앙 집중식 스키마 관리 및 검증을 제공하여 모든 메시지 데이터의 호환성을 보장함으로써 이 문제를 해결합니다.

 

스키마 진화 (Schema evolution)

스키마는 시간이 지남에 따라 자주 변경되며, 이는 서로 다른 스키마 버전 간의 호환성 문제를 일으킬 수 있습니다. 스키마 레지스트리는 스키마 버전 관리를 지원하여, 여러 버전의 스키마를 호환성 문제 없이 동시에 사용할 수 있도록 해줍니다.

 

스키마 ID 검증 (Schema ID validation)

스키마 레지스트리는 토픽으로 발행되는 데이터가 레지스트리에 등록된 유효한 스키마 ID를 사용하는지 검증합니다. 이를 통해 데이터가 표준 형식을 준수하도록 만들어 데이터 손실이나 손상 위험을 줄입니다.

 

데이터 거버넌스 (Data governance)

스키마 레지스트리는 데이터 스키마를 관리하고 버전을 제어하는 중앙 공간을 제공합니다. 스키마 변경 사항 추적, 스키마 발전 이력 관리, 규제 요구사항에 대한 준수 여부 확인을 용이하게 하여 데이터 거버넌스를 단순화합니다.

 

얼추 맞는데? 하여간 스키마 레지스트리를 사용하면 스키마 관리를 편하게 할 수 있고, 이러쿵 저러쿵 좋다!

 

 

어떻게 스키마 레지스트리는 작동하는가?

출처 : https://docs.confluent.io/platform/current/schema-registry/fundamentals/index.html#id2

 

 

schema registry가 어떻게 producer와 consumer와 연결되는지 설명해준다. 

 

다만 순번을 차례로 나눠서 친절하게 설명해주지 않아서 보기가 불편했다.

 

 

이에 대해서는 Gemini에게 정리를 요청했다. 한 번 읽어보고 넘어가겠다!

 

2.1) Producer 측 흐름 (데이터 발행)

Producer가 객체를 Avro, Protobuf, 또는 JSON 형식으로 직렬화하여 Kafka로 전송하는 과정입니다.

  1. 스키마 캐시 확인 및 등록 (Send/Register Schema): Producer가 데이터를 전송하기 전, 해당 데이터 구조에 맞는 스키마가 자신의 로컬 캐시(Local cache for schemas) 에 존재하는지 확인합니다. 만약 캐시에 없다면, Schema Registry로 스키마를 전송하여 등록을 요청합니다.
  2. 스키마 ID 반환: Schema Registry는 해당 스키마를 저장하고(예: schema-1), 이에 매핑되는 고유한 식별자(Schema ID)를 Producer에게 반환합니다.
  3. 페이로드 조립 및 직렬화 (Serialize per schema id): Producer는 반환받은(또는 캐시에서 찾은) 스키마를 기반으로 실제 데이터(Value)를 이진(Binary) 데이터로 직렬화합니다. 이때 Kafka로 보낼 최종 페이로드의 앞부분에 4바이트의 매직 바이트와 스키마 ID를 헤더처럼 붙입니다. 즉, 메시지는 [Schema ID] + [직렬화된 Data] 형태가 됩니다.
  4. Kafka 브로커로 전송: 스키마 ID와 데이터가 결합된 경량화된 메시지를 Kafka 브로커로 전송(Publish)합니다.

 

재미있는 부분은 프로듀서에 로컬 캐시가 있다는 점이다. 스키마 정보를 사용햇다면 여기에 쌓이나 보다. 

 

2.2) Consumer 측 흐름 (데이터 소비)

Consumer가 Kafka로부터 바이너리 데이터를 읽어와 원래의 객체로 복원하는 과정입니다.

  1. Kafka에서 메시지 읽기 (Read Data): Consumer는 Kafka 브로커로부터 [Schema ID] + [직렬화된 Data] 형태로 구성된 메시지를 폴링(Polling)하여 읽어옵니다.
  2. 스키마 ID 추출 및 캐시 확인: Consumer는 페이로드의 앞부분에서 Schema ID를 추출합니다. 그리고 이 ID에 해당하는 스키마가 자신의 로컬 캐시(Local cache for schemas) 에 있는지 확인합니다.
  3. 스키마 조회 (Get schema by id): 만약 로컬 캐시에 해당 ID의 스키마가 없다면, Consumer는 Schema Registry에 해당 ID를 보내어 원본 스키마(예: schema-1)를 요청하고 응답받아 캐싱합니다.
  4. 역직렬화 (Deserialize): 조회한 스키마 구조를 바탕으로 페이로드의 나머지 부분인 [직렬화된 Data]를 읽어 들여, 애플리케이션에서 사용할 수 있는 원래의 객체(Java Object 등)로 역직렬화합니다.

 

컨슈머도 마찬가지로 로컬 캐시가 있다. 

 

3. 스키마 고도화와 호환성

 

호환성 하니 호환 호은 호환마마가 생각난다.

 

수호전 의  무송 을 습격한 호랑이

 

생각보다 쉽지 않네... 어려움이 호환급이다!

 

해당 정보의 출처는 다음과 같다.

(컨플루언스 링크 : Schema Evolution and Compatibility for Schema Registry on Confluent Platform)

 

 

스키마는 처음 작성하고 계속 두는것이 아니라 상황에 맞게 끊임없이 고도화해야한다.

 

Confluent Schema Registry에서는 producer와 consumer가 있어도 schema의 호환성을 해치지 않고 수정이 가능하다.

 

여기서 호환성이라는 말이 나오는데 호환성은 무엇인가?

여러 장치, 부품, 소프트웨어 등이 서로 충돌 없이 잘 맞아떨어져 함께 작동할 수 있는 성질

 

 

스키마 레지스트리에서는 어디가 충돌 없이 잘 맞아 떨어지게 한다는 것일까?

 

프로듀서와 컨슈머는 서로 대화하지 않는다 (대화가 필요해~)

 

특히나 프로듀서는 메시지를 발행하고 가버린다. 결국에 컨슈머가 나중에 혼자 읽을 때 이게 맞는지 아닌지 (호환성)을 확인해야한다.

 

 

호환성 타입(compatibility type)을 사용하면 허용되는 스키마 변경 사항을 정의할 수 있다.

 

  • Backward(역방향) - 새 스키마를 사용하는 소비자는 이전 스키마로 작성된 데이터를 읽을 수 있습니다(선택적 필드 추가, 필드 제거).
  • Forward(포워드) - 기존 스키마를 사용하는 소비자는 새 스키마로 작성된 데이터를 읽을 수 있습니다(선택적 필드 제거, 필드 추가).
  • Full(완전 버전) - 이전 버전 및 이후 버전 모두 호환 가능 (선택 필드만 추가/삭제 가능)
  • Transitive(전이적 호환성) - 최신 버전뿐만 아니라 이전 모든 버전과의 호환성을 확인합니다.

 

기본 호환 모드는 Backward이다. 각 스키마 버전에는 고유 ID와 순차적으로 증가하는 버전 번호가 부여된다.

 

스키마가 업데이트되면 스키마 레지스트리는 새 버전을 수락하기 전에 호환성을 확인한다.

 

호환성 검사는 스키마 문서 두 장을 비교하는 작업이다. 실제 메시지도, 실행 중인 애플리케이션도 관여하지 않는다.

  • BACKWARD: 새 스키마로 옛 데이터를 해석 가능한가
  • FORWARD: 옛 스키마로 새 데이터를 해석 가능한가
  • FULL: 둘 다
  직전 버전만 모든 이전 버전
새->옛 데이터 BACKWARD BACKWARD_TRANSITIVE
옛->새 데이터 FORWARD FORWARD_TRANSITIVE
양방향 FULL FULL_TRANSITIVE

 

 

이걸 어떻게 프로듀서와 컨슈머 관점에서 바라 볼 수 있는지 다음 그림들로 확인하자.

 

3.1) Backward 타입

 

Backward호환성은 진화된 스키마를 적용한 컨슈머가 진화 전의 스키마가 적용된 프로듀서가 보낸 메시지를 읽을 수 잇도록 허용하는 호환성을 말한다.

 

간단히 말하면 새 스키마로 옛 데이터 읽기가 가능하다.

 

3.2) Forward 타입

 

forward는 상위 버전 스미카를 먼저 프로듀서에게 적용한 다음, 컨슈머에게 적용한다.

 

이전 버전 스키마를 쓰는 컨슈머가, 새 데이터를 읽을 수 있다. 즉, 옛 스키마로 새 데이터 읽기 가능하다.

 

3.3) Full 타입

 

 

full 은 말 그대로 둘 다! backward와 forward를 둘 다 지원한다.

 

간단히 말하면 새 스키마로 옛 데이터 읽기가 가능하며 옛 스키마로 새 데이터 읽기 가능하다.

 

4) TRANSITIVE가 하는 일

딱 하나로 등록을 거부한다. 어떻게? 전부 다 확인해서! 

 

 

시나리오 한 줄 요약

BACKWARD : v3 등록 요청 → v2와만 비교 → 통과 → v3 생성됨

BACKWARD_TRANSITIVE : v3 등록 요청 → v2,v1 다 비교 → v1과 안 맞음 → 거부 → v3 안 생김

FORWARD : v3 등록 요청 → v2와만 비교 → 통과 → v3 생성됨

FORWARD_TRANSITIVE : v3 등록 요청 → v2, v1 다 비교 → v1과 안 맞음 → 거부 → v3 안 생김

FULL : v3 등록 요청 → v2와만 비교(양방향) → 통과 → v3 생성됨

FULL_TRANSITIVE : v3 등록 요청 → v2, v1 다 비교(양방향) → v1과 안 맞음 → 거부 → v3 안 생김

 

 

 

4. 사례 분석

4.1) 요기오의 Conflunent Schema 도입

요기오에서는 R&D 센터의 데이터 파이프라인 엔지니어링 팀에서 Confluent Schem Registry의 도입기를 소개한다.

 

필자가 원했던 부분은 구체적으로 현재 회사에서 어떤 어떤 부분이 있는데 이 부분에서 적용이 필요하다 판단했기에 이것을 사용했다는 히스토리였다. 해당 글의 아쉬운 점은 개념 설명 뒤 이러한 히스토리없이 구현체 설명에 집중해서 아쉬웠다.

 

그런가 보다~ 하고 참고할 정도로는 적당하다 생각한다.

 

(참고글 : 요기오 - Confluent Schema Registry 도입기!)

 

4.2)  11번가의 Live 11 과 Schema Registry

11번가에서 스키마 레지스트리를 어떻게 도입했고 "Kafka: The Definitive Guide" 의 내용을 요약해서 설명해주고 있다.

(참고글 : Live 11 과 Schema Registry)

 

아주 친절하게도 도입기에 대해서 유튜브로도 소개를 하기에 시간이 되면 꼭! 빡집중해서 시청하려고 한다.

 

우리가 직접 구현을 할 수 있겠지만 Confluent 에서 제공하는 Avro 혹은 AWS에서 제공하는 AWS Glue Schema  Registry가 있으니 설명하는대로 잘 사용하면 될 듯 하다.

 

Live11 에서는 AWS Glue Schema Registry 를 사용하고 있으니 이를 중점적으로 설명한다.

 

 

크게는 Schema Registry에 schema 정보를 등록한다.

 

Producer가  schema 의 identifier 를 record 에 달아서 serialize한다.

 

Consumer는 그 identifier를 가지고 shema registry에서 schema를 가지고 와서 deserialize한다.

 

 

당장 보이는 사례로는 요기오와 11번가의 Live 11이다. 외국의 사례도 더 찾게 되면 필자가 차가해 놓겠다.

 

5. 실습하기(Avro 사용)

컨플루언트에서는 Avro, JSON, ProtoBuf 포맷을 지원하는데 Avro 를 추천한다. 이유는 다음과 같다.

 

  1. 에이브로는 JSON과 매핑된다.
  2. 에이브로는 매우 간결한 데이터 포맷이다.
  3. JSON은 메시지마다 필드 네임들이 포함되어 전송되므로 효율이 떨어진다.
  4. 에이브로는 바이너리 형태이므로 매우 빠르다.

 우선 docker파일을 작성하자.

 

docker-compose-schema-registry.yml

networks:
  # 카프카 브로커들과 스키마 레지스트리가 서로 통신하기 위한 도커 내부망입니다.
  kafka-network:
    name: kafka-kraft-net
    driver: bridge

services:
  # --------------------------------------------------------
  # 브로커 1 (컨트롤러 겸 브로커)
  # --------------------------------------------------------
  kafka-kraft-1:
    image: apache/kafka:latest
    container_name: kafka-kraft-1
    user: root
    ports:
      # 외부(Spring Boot)에서 이 브로커로 접속할 때 사용하는 포트
      - "19092:19092"
    environment:
      # 클러스터 식별자. 3대의 브로커가 동일한 클러스터로 묶이기 위해 동일한 ID를 공유합니다.
      CLUSTER_ID: 'O1N_F_D-S2K6V8-A_G-p-A'
      KAFKA_NODE_ID: 1
      # KRaft 모드이므로 Zookeeper 없이 이 노드가 브로커와 컨트롤러 역할을 모두 수행합니다.
      KAFKA_PROCESS_ROLES: 'controller,broker'
      
      # [투표자 설정] 뗏목(Raft) 합의 알고리즘을 위해 1, 2, 3번 노드가 9093 포트로 통신하며 리더를 선출합니다.
      KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka-kraft-1:9093,2@kafka-kraft-2:9093,3@kafka-kraft-3:9093'
      
      # [바인딩할 포트] INTERNAL(내부망), CONTROLLER(리더 선출용), EXTERNAL(외부 접속용) 포트를 개방합니다.
      KAFKA_LISTENERS: 'INTERNAL://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093,EXTERNAL://0.0.0.0:19092'
      # [광고할 포트] 클라이언트(Spring Boot나 Schema Registry)가 접속을 요청할 때 "이 주소로 들어와!"라고 알려주는 주소입니다.
      KAFKA_ADVERTISED_LISTENERS: 'INTERNAL://kafka-kraft-1:9092,EXTERNAL://localhost:19092'
      
      # 이번 구성은 순수 동작 확인용이므로 모든 통신을 암호화 없는 PLAINTEXT로 설정합니다.
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT'
      KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER'
      KAFKA_INTER_BROKER_LISTENER_NAME: 'INTERNAL'
      KAFKA_LOG_DIRS: '/tmp/kraft-combined-logs'
    volumes:
      # 카프카 토픽의 세그먼트 파일(로그)들이 컨테이너가 내려가도 유지되도록 볼륨을 잡습니다.
      - kafka-kraft-1-data:/tmp/kraft-combined-logs
    networks:
      - kafka-network

  # --------------------------------------------------------
  # 브로커 2 (컨트롤러 겸 브로커)
  # - 1번 브로커와 역할은 동일하며 노드 ID와 포트(29092)만 다릅니다.
  # --------------------------------------------------------
  kafka-kraft-2:
    image: apache/kafka:latest
    container_name: kafka-kraft-2
    user: root
    ports:
      - "29092:29092"
    environment:
      CLUSTER_ID: 'O1N_F_D-S2K6V8-A_G-p-A'
      KAFKA_NODE_ID: 2
      KAFKA_PROCESS_ROLES: 'controller,broker'
      KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka-kraft-1:9093,2@kafka-kraft-2:9093,3@kafka-kraft-3:9093'
      KAFKA_LISTENERS: 'INTERNAL://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093,EXTERNAL://0.0.0.0:29092'
      KAFKA_ADVERTISED_LISTENERS: 'INTERNAL://kafka-kraft-2:9092,EXTERNAL://localhost:29092'
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT'
      KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER'
      KAFKA_INTER_BROKER_LISTENER_NAME: 'INTERNAL'
      KAFKA_LOG_DIRS: '/tmp/kraft-combined-logs'
    volumes:
      - kafka-kraft-2-data:/tmp/kraft-combined-logs
    networks:
      - kafka-network

  # --------------------------------------------------------
  # 브로커 3 (컨트롤러 겸 브로커)
  # - 노드 ID와 포트(39092)를 제외하고 위와 동일합니다.
  # --------------------------------------------------------
  kafka-kraft-3:
    image: apache/kafka:latest
    container_name: kafka-kraft-3
    user: root
    ports:
      - "39092:39092"
    environment:
      CLUSTER_ID: 'O1N_F_D-S2K6V8-A_G-p-A'
      KAFKA_NODE_ID: 3
      KAFKA_PROCESS_ROLES: 'controller,broker'
      KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka-kraft-1:9093,2@kafka-kraft-2:9093,3@kafka-kraft-3:9093'
      KAFKA_LISTENERS: 'INTERNAL://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093,EXTERNAL://0.0.0.0:39092'
      KAFKA_ADVERTISED_LISTENERS: 'INTERNAL://kafka-kraft-3:9092,EXTERNAL://localhost:39092'
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT'
      KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER'
      KAFKA_INTER_BROKER_LISTENER_NAME: 'INTERNAL'
      KAFKA_LOG_DIRS: '/tmp/kraft-combined-logs'
    volumes:
      - kafka-kraft-3-data:/tmp/kraft-combined-logs
    networks:
      - kafka-network

  # --------------------------------------------------------
  # Schema Registry (스키마 저장소)
  # - Producer가 보낸 Avro 스키마 구조(JSON)를 카프카 브로커의 `_schemas` 토픽에 저장합니다.
  # - Consumer는 역직렬화할 때 이 서버에 REST API로 스키마 구조를 질의하여 데이터를 읽습니다.
  # --------------------------------------------------------
  schema-registry:
    image: confluentinc/cp-schema-registry:7.4.0
    container_name: schema-registry
    # 카프카 브로커들이 전부 뜬 다음에 스키마 레지스트리가 떠야 하므로 의존성을 걸어줍니다.
    depends_on:
      - kafka-kraft-1
      - kafka-kraft-2
      - kafka-kraft-3
    ports:
      # 외부(Spring Boot)에서 스키마를 등록하거나 조회할 때 사용할 REST API 포트입니다.
      - "8081:8081"
    environment:
      # 도커 네트워크 안에서 자기 자신을 부를 때 사용할 호스트명입니다.
      SCHEMA_REGISTRY_HOST_NAME: schema-registry
      
      # [매우 중요] 스키마 레지스트리가 텍스트(JSON)로 된 스키마 파일을 저장할 DB로 쓸 카프카 브로커들의 주소입니다.
      # Zookeeper 방식에서는 커넥션 URL을 썼지만, KRaft 환경에서는 KAFKASTORE_BOOTSTRAP_SERVERS를 사용해 내부망(9092)으로 직접 찌릅니다.
      SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: 'kafka-kraft-1:9092,kafka-kraft-2:9092,kafka-kraft-3:9092'
      
      # 스키마 레지스트리 서버가 8081 포트로 들어오는 HTTP 요청(GET /subjects 등)을 받겠다고 선언합니다.
      SCHEMA_REGISTRY_LISTENERS: http://0.0.0.0:8081
    networks:
      - kafka-network

volumes:
  kafka-kraft-1-data:
  kafka-kraft-2-data:
  kafka-kraft-3-data:

 

 

 

docker rm -f kafka-kraft-1 kafka-kraft-2 kafka-kraft-3 schema-registry

 

docker network rm kafka-kraft-net

 

실행 전에 필자는 위의 명령어들을 사용해서 깔끔하게 정리했다.

 

docker compose -f docker-compose.yml up -d

 

그러면 다음처럼 성공하는 것을 볼 수 있다.

 

Gemini에게 스키마 레지스트리에서 사용할 수 있는 API 를 정리해달라고 했다.

메서드 엔드포인트 설명
POST /subjects/{subject}/versions 새로운 스키마를 등록하고, 전역 고유 ID를 반환받습니다
GET /subjects 현재 등록된 모든 Subject 목록을 조회합니다.
GET /subjects/{subject}/versions/{version} 특정 Subject의 특정 버전 스키마 구조를 조회합니다.
GET /schemas/ids/{id} 스키마의 전역 고유 ID로 스키마 구조를 직접 조회합니다.
(Consumer가 주로 사용)
DELETE /subjects/{subject} 특정 Subject를 삭제합니다. (개발 환경에서 초기화 용도)

 

5.1) 스키마 등록 (POST)

users라는 토픽에 데이터를 보낸다고 가정하고, users-value라는 Subject로 간단한 Avro 스키마를 등록한다.

 

Content-Type은 application/vnd.schemaregistry.v1+json을 사용하는 것이 권장된다고 한다.

curl -X POST -H "Content-Type: application/vnd.schemaregistry.v1+json" \
--data '{
  "schema": "{\"type\": \"record\", \"name\": \"User\", \"fields\": [{\"name\": \"id\", \"type\": \"long\"}, {\"name\": \"username\", \"type\": \"string\"}]}"
}' \
http://localhost:8081/subjects/users-value/versions

 

정상 응답 예시

{"id": 2} (이 ID는 Kafka 클러스터 전체에서 해당 스키마 구조를 식별하는 고유 번호)

 

성공하면 id값을 반환한다.

 

application/vnd.schemaregistry.v1+json은 HTTP 통신에서 데이터의 타입을 정의하는 커스텀 MIME 타입(Media Type)으로

 

REST API 설계 원칙 중 '콘텐츠 협상(Content Negotiation)'을 위해 명시적으로 지정된 규격이라고 한다.

(처음보는데 이걸 쓸 일이 있나 내가...?)

 

그리고 --data 하는 부분은 아파치 에이브로 규격이라고 한다.

 

{
  "type": "record",
  "name": "User",
  "fields": [
    {"name": "id", "type": "long"},
    {"name": "username", "type": "string"}
  ]
}

 

이게 자바로 표현하면

public class User { // type: record, name: User
    private long id;
    private String username;
}

 

 

 

Apache Avro의 Specification에서는 데이터 형식을 선언하는 방식에 대해 설명한다.

 

크게는 원시 타입(Primitive Types) 과 복합 타입(Complex Types) 라고 소개한다.

 

이에 대해서 Gemini에게 정리를 요청했다.

 

원시 타입(Primitive Types)

가장 기본이 되는 단순 데이터 타입으로 속싱이 따로 필요없다.

 

  • null: 데이터가 없음(No value)을 나타냅니다.
  • boolean: 참(true) 또는 거짓(false)의 이진 값입니다.
  • int: 32비트 부호 있는 정수입니다. (일반적인 숫자)
  • long: 64비트 부호 있는 정수입니다. (매우 큰 숫자)
  • float: 32비트 부동 소수점 숫자입니다. (소수점)
  • double: 64비트 부동 소수점 숫자입니다. (더 정밀한 소수점)
  • bytes: 8비트 부호 없는 바이트들의 배열입니다. (이미지, 파일 등 이진 데이터)
  • string: 유니코드(UTF-8) 문자열입니다. (텍스트)

 

복합 타입(Complex Types)

원시 타입들을 조합하여 만드는 복잡한 데이터 구조

 

  • record: 여러 개의 변수(필드)를 하나로 묶어놓은 형태입니다. (Java의 Class, C의 Struct와 동일한 핵심 타입)
  • enum: 미리 정해둔 특정 문자열 값들 중에서만 선택할 수 있게 강제하는 열거형 타입입니다. (예: ["RED", "GREEN", "BLUE"])
  • array: 동일한 타입의 데이터들을 순서대로 담아두는 배열(리스트)입니다.
  • map: 문자열(String)을 키(Key)로 하고, 특정 타입의 데이터를 값(Value)으로 가지는 사전(Dictionary) 형태입니다.
  • union: 하나의 필드가 '여러 가지 타입 중 하나'를 가질 수 있게 해줍니다. (예: ["null", "string"]으로 설정하면, 문자열이 올 수도 있고 값이 없을 수도 있다는 뜻으로, 주로 기본값/Null 처리에 쓰입니다.)
  • fixed: 크기가 정확히 고정된 바이트 배열입니다. (일반 bytes와 달리 무조건 정해진 바이트 크기만 들어올 수 있습니다.)

 

일반적으로 primitive타입으로 통신할 일은 사실 많이 없을것 같다.

 

어플리케이션으로 각각 관리해서 통신을 하지 단순하게 응답을 한다면 복잡환 상황에서 활용도가 떨어지기 때문이다.

 

기본적으로 백엔드 서버 써본대야 Java, Typescript, Python인데 과연... 단순하게 보내서 쓸 경우는 극단적인 상황을 제외하고는 딱히 와닿지가 않는다. 

 

애초에 schema registry의 목적 자체가 데이터 형식을 하나에서 관리하고 지키기 위해서 그런 것 아닌가? (지극히 개인적인 생각)

 

 

그리고 이러한 타입을 솔직히 내가 어떻게 정하고 쓸지를 AI에게 작성해 달라고 하면 금방이다.

 

내가 쓰는 언어에 맞춰서 어떻게 써야하는지 알면 되는 정도이니 그런가보구나~ 하는 정도로 이해하고 넘어갈 것이다.

 

 

하여간 현재 폴더는 users-value로 schema를 생성한 것이다.

 

그 폴더 안에 우리가 이름이 "User"라는 스키마를 id, username이라는 필드를 가지로 long, string타입으로 선언한 것이다.

 

그런데 왜 하필 폴더는 users-value일까?

 

물론 내가 그렇게 등록했다. (하라는 대로 했을 뿐입니다... ㅠㅠ)

 

그렇다면 이 이름은 어디서 온 것일까?

http://localhost:8081/subjects/users-value/versions
                               ^^^^^^^^^^^ 직접 타이핑한 값

 

Schema Registry는 Kafka와 별개의 서버다.

 

포트부터 다르다(9092 / 8081). 따라서 토픽이라는 개념을 알지 못하며, URL에 적힌 문자열을 그대로 폴더 이름처럼 사용할 뿐이다.

 

실제로 이 시점에 users라는 토픽은 존재하지도 않았지만 등록은 정상적으로 완료됐다.

그렇다면 왜 굳이 -value를 붙였는가. 프로듀서 실행 시 출력되는 로그에 답이 있다.

key.subject.name.strategy   = class io.confluent...TopicNameStrategy
value.subject.name.strategy = class io.confluent...TopicNameStrategy

 

클라이언트는 Subject 이름을 다음 규칙으로 계산한다.

subject = {토픽명}-key
subject = {토픽명}-value

 

Kafka 메시지는 key와 value로 구성되고 둘은 서로 다른 스키마를 가질 수 있으므로, 하나의 토픽에 최대 두 개의 Subject가 대응된다. 즉 -value강제된 규칙이 아니라 클라이언트와의 약속이다. hello라고 지어도 등록은 되지만, 나중에 프로듀서가 찾지 못한다.

 

뭔가 헷갈린다. 다시 생각해보자. 

Kafka 입장에서 메시지는 바이트 두 덩어리

다.

메시지 = [ key 바이트 ] + [ value 바이트 ]

이 둘은 직렬화 방식조차 서로 다를 수 있다.

key.serializer   = ByteArraySerializer      ← 그냥 바이트
value.serializer = ByteArraySerializer      ← Avro 래퍼가 감싸서 처리

 

그러니까 Avro는 key랑 value를 별도의 schema를 쓰도록 설정할 수 있는데, 현재 필자는 value만 한 것이다.

 

굳이 key는 key일 뿐이지 얼마나 복잡하면 또 그걸 뭐 schema까지 설정해야하나 싶다.

 
 

5.2) Subject 목록 조회 (GET)

현재 레지스트리에 어떤 Subject들이 등록되어 있는지 확인한다.

curl -X GET http://localhost:8081/subjects

 

 

현재 필자가 테스트하면서 2개가 있다. ("schema-test-topic-value", "users-value")

 

5.3) 특정 버전의 스키마 상세 조회 (GET)

등록된 users-value 스키마의 가장 최신 버전(latest) 구조를 확인한다.

curl -X GET http://localhost:8081/subjects/users-value/versions/latest

 

 

5.4) 스키마 ID로 조회 (GET)

Consumer 애플리케이션은 메시지를 읽을 때 데이터 앞에 붙어있는 4바이트의 스키마 ID(예: 1)만 읽어낸다.

 

이후 아래 API를 통해 스키마 구조를 가져와 역직렬화를 수행한다.

 

curl -X GET http://localhost:8081/schemas/ids/1

 

난리난 나의 실습상황 ㅠㅠ...

 

여기까지는 맛보기였다. schema-registry로 요청을 보내서 작동하는것을 봤다면 

 

실제로 kafka에서 터미널 환경에서 메시지를 발행하고 소비해보도록 하겠다.

 

5.5) 컨슈머 + 프로듀서 실행을 통한 메시지 전송

 

필자는 mini pc 환경에서 이를 테스트 하기 위해 window 에서 원격 접속 후 cmd 창을 두개로 사용했다.

 

 

5.5.1) 컨슈머 실행

 

먼저 데이터를 수신할 컨슈머를 백그라운드에 대기시킨다.

docker exec -it schema-registry kafka-avro-console-consumer \
  --bootstrap-server kafka-kraft-1:9092 \
  --topic schema-test-topic \
  --property schema.registry.url=http://schema-registry:8081 \
  --from-beginning

 

각 줄 하나씩 설명하겠다.

 

  1. docker exec ... kafka-avro-console-consumer: Avro 데이터 전용 읽기 도구 실행
  2. --bootstrap-server: 접속할 카프카 서버 주소
  3. --topic: 데이터를 읽어올 카프카 토픽
  4. --property schema.registry.url=...: 바이너리 데이터를 해석할 '데이터 구조(스키마)'를 가져올 주소
  5. --from-beginning: 해당 토픽에 저장된 맨 처음 과거 데이터부터 모두 읽기

 

5.5.2) V1 스키마 정의 및 정상 메시지 발행

 

name 필드만 존재하는 V1 스키마로 프로듀서를 실행하고 메시지를 보낸다.

docker exec -i schema-registry kafka-avro-console-producer \
  --bootstrap-server kafka-kraft-1:9092 \
  --topic schema-test-topic \
  --property schema.registry.url=http://schema-registry:8081 \
  --property value.schema='{"type":"record","name":"User","fields":[{"name":"name","type":"string"}]}' <<EOF
{"name":"Developer"}
EOF

 

 

그러면 이제 5.5.1) 화면에서 {"name": "Developer"} 가 나온다.

 

여러 뭐 많이 뜨지만 가장 하단을 보라!

 

 

 

스키마 등록하고 등록한 스키마의 형식에 맞게 메시지를 보낸 것이다. 

 

그리고 다음 명령어를 실행하면 버전을 알 수 있다.

 

docker exec -it schema-registry curl -s http://localhost:8081/subjects/schema-test-topic-value/versions

 

 

아직 스키마 변경을 안해서 version이 1이다.

 

5.5.3) 스키마 진화!

 

이전의 프로듀서를 종료한다.

 

이번에는 필수 필드인 age 를 기본값없이 추가한 v2 스키마로 다시 프로듀서를 실행한다.

 

docker exec -i schema-registry kafka-avro-console-producer \
  --bootstrap-server kafka-kraft-1:9092 \
  --topic schema-test-topic \
  --property schema.registry.url=http://schema-registry:8081 \
  --property value.schema='{"type":"record","name":"User","fields":[{"name":"name","type":"string"},{"name":"age","type":"int"}]}' <<EOF
{"name":"Kim","age":30}
EOF

 

실패하면 다음과 같은 메시지가 나온다.

 

$ docker exec -it schema-registry curl -s http://localhost:8081/subjects/schema-test-topic-value/versions
[1docker exec -i schema-registry kafka-avro-console-producer \r \
  --bootstrap-server kafka-kraft-1:9092 \
  --topic schema-test-topic \
  --property schema.registry.url=http://schema-registry:8081 \
  --property value.schema='{"type":"record","name":"User","fields":[{"name":"name","type":"string"},{"name":"age","type":"int"}]}' <<EOF
{"name":"Kim","age":30}
EOF
[2026-07-23 15:19:39,336] INFO Registered kafka:type=kafka.Log4jController MBean (kafka.utils.Log4jControllerRegistration$)
[2026-07-23 15:19:39,442] INFO KafkaAvroSerializerConfig values:
        auto.register.schemas = true
        avro.reflection.allow.null = false
        avro.remove.java.properties = false
        avro.use.logical.type.converters = false
        basic.auth.credentials.source = URL
        basic.auth.user.info = [hidden]
        bearer.auth.cache.expiry.buffer.seconds = 300
        bearer.auth.client.id = null
        bearer.auth.client.secret = null
        bearer.auth.credentials.source = STATIC_TOKEN
        bearer.auth.identity.pool.id = null
        bearer.auth.issuer.endpoint.url = null
        bearer.auth.logical.cluster = null
        bearer.auth.scope = null
        bearer.auth.scope.claim.name = scope
        bearer.auth.sub.claim.name = sub
        bearer.auth.token = [hidden]
        context.name.strategy = class io.confluent.kafka.serializers.context.NullContextNameStrategy
        id.compatibility.strict = true
        key.subject.name.strategy = class io.confluent.kafka.serializers.subject.TopicNameStrategy
        latest.cache.size = 1000
        latest.cache.ttl.sec = -1
        latest.compatibility.strict = true
        max.schemas.per.subject = 1000
        normalize.schemas = false
        proxy.host =
        proxy.port = -1
        rule.actions = []
        rule.executors = []
        schema.format = null
        schema.reflection = false
        schema.registry.basic.auth.user.info = [hidden]
        schema.registry.ssl.cipher.suites = null
        schema.registry.ssl.enabled.protocols = [TLSv1.2, TLSv1.3]
        schema.registry.ssl.endpoint.identification.algorithm = https
        schema.registry.ssl.engine.factory.class = null
        schema.registry.ssl.key.password = null
        schema.registry.ssl.keymanager.algorithm = SunX509
        schema.registry.ssl.keystore.certificate.chain = null
        schema.registry.ssl.keystore.key = null
        schema.registry.ssl.keystore.location = null
        schema.registry.ssl.keystore.password = null
        schema.registry.ssl.keystore.type = JKS
        schema.registry.ssl.protocol = TLSv1.3
        schema.registry.ssl.provider = null
        schema.registry.ssl.secure.random.implementation = null
        schema.registry.ssl.trustmanager.algorithm = PKIX
        schema.registry.ssl.truststore.certificates = null
        schema.registry.ssl.truststore.location = null
        schema.registry.ssl.truststore.password = null
        schema.registry.ssl.truststore.type = JKS
        schema.registry.url = [http://schema-registry:8081]
        use.latest.version = false
        use.latest.with.metadata = null
        use.schema.id = -1
        value.subject.name.strategy = class io.confluent.kafka.serializers.subject.TopicNameStrategy
 (io.confluent.kafka.serializers.KafkaAvroSerializerConfig)
[2026-07-23 15:19:39,788] INFO ProducerConfig values:
        acks = -1
        auto.include.jmx.reporter = true
        batch.size = 16384
        bootstrap.servers = [kafka-kraft-1:9092]
        buffer.memory = 33554432
        client.dns.lookup = use_all_dns_ips
        client.id = console-producer
        compression.type = none
        confluent.proxy.protocol.client.address = null
        confluent.proxy.protocol.client.port = null
        confluent.proxy.protocol.client.version = NONE
        connections.max.idle.ms = 540000
        delivery.timeout.ms = 120000
        enable.idempotence = true
        interceptor.classes = []
        key.serializer = class org.apache.kafka.common.serialization.ByteArraySerializer
        linger.ms = 1000
        max.block.ms = 60000
        max.in.flight.requests.per.connection = 5
        max.request.size = 1048576
        metadata.max.age.ms = 300000
        metadata.max.idle.ms = 300000
        metric.reporters = []
        metrics.num.samples = 2
        metrics.recording.level = INFO
        metrics.sample.window.ms = 30000
        partitioner.adaptive.partitioning.enable = true
        partitioner.availability.timeout.ms = 0
        partitioner.class = null
        partitioner.ignore.keys = false
        receive.buffer.bytes = 32768
        reconnect.backoff.max.ms = 1000
        reconnect.backoff.ms = 50
        request.timeout.ms = 1500
        retries = 3
        retry.backoff.ms = 100
        sasl.client.callback.handler.class = null
        sasl.jaas.config = null
        sasl.kerberos.kinit.cmd = /usr/bin/kinit
        sasl.kerberos.min.time.before.relogin = 60000
        sasl.kerberos.service.name = null
        sasl.kerberos.ticket.renew.jitter = 0.05
        sasl.kerberos.ticket.renew.window.factor = 0.8
        sasl.login.callback.handler.class = null
        sasl.login.class = null
        sasl.login.connect.timeout.ms = null
        sasl.login.read.timeout.ms = null
        sasl.login.refresh.buffer.seconds = 300
        sasl.login.refresh.min.period.seconds = 60
        sasl.login.refresh.window.factor = 0.8
        sasl.login.refresh.window.jitter = 0.05
        sasl.login.retry.backoff.max.ms = 10000
        sasl.login.retry.backoff.ms = 100
        sasl.mechanism = GSSAPI
        sasl.oauthbearer.clock.skew.seconds = 30
        sasl.oauthbearer.expected.audience = null
        sasl.oauthbearer.expected.issuer = null
        sasl.oauthbearer.jwks.endpoint.refresh.ms = 3600000
        sasl.oauthbearer.jwks.endpoint.retry.backoff.max.ms = 10000
        sasl.oauthbearer.jwks.endpoint.retry.backoff.ms = 100
        sasl.oauthbearer.jwks.endpoint.url = null
        sasl.oauthbearer.scope.claim.name = scope
        sasl.oauthbearer.sub.claim.name = sub
        sasl.oauthbearer.token.endpoint.url = null
        security.protocol = PLAINTEXT
        security.providers = null
        send.buffer.bytes = 102400
        socket.connection.setup.timeout.max.ms = 30000
        socket.connection.setup.timeout.ms = 10000
        ssl.cipher.suites = null
        ssl.enabled.protocols = [TLSv1.2, TLSv1.3]
        ssl.endpoint.identification.algorithm = https
        ssl.engine.factory.class = null
        ssl.key.password = null
        ssl.keymanager.algorithm = SunX509
        ssl.keystore.certificate.chain = null
        ssl.keystore.key = null
        ssl.keystore.location = null
        ssl.keystore.password = null
        ssl.keystore.type = JKS
        ssl.protocol = TLSv1.3
        ssl.provider = null
        ssl.secure.random.implementation = null
        ssl.trustmanager.algorithm = PKIX
        ssl.truststore.certificates = null
        ssl.truststore.location = null
        ssl.truststore.password = null
        ssl.truststore.type = JKS
        transaction.timeout.ms = 60000
        transactional.id = null
        value.serializer = class org.apache.kafka.common.serialization.ByteArraySerializer
 (org.apache.kafka.clients.producer.ProducerConfig)
[2026-07-23 15:19:39,825] INFO [Producer clientId=console-producer] Instantiated an idempotent producer. (org.apache.kafka.clients.producer.KafkaProducer)
[2026-07-23 15:19:39,908] INFO Kafka version: 7.4.0-ce (org.apache.kafka.common.utils.AppInfoParser)
[2026-07-23 15:19:39,909] INFO Kafka commitId: aee97a585bd06866 (org.apache.kafka.common.utils.AppInfoParser)
[2026-07-23 15:19:39,909] INFO Kafka startTimeMs: 1784819979897 (org.apache.kafka.common.utils.AppInfoParser)
org.apache.kafka.common.errors.InvalidConfigurationException: Schema being registered is incompatible with an earlier schema for subject "schema-test-topic-value", details: [{errorType:'READER_FIELD_MISSING_DEFAULT_VALUE', description:'The field 'age' at path '/fields/1' in the new schema has no default value and is missing in the old schema', additionalInfo:'age'}, {oldSchemaVersion: 1}, {oldSchema: '{"type":"record","name":"User","fields":[{"name":"name","type":"string"}]}'}, {compatibility: 'BACKWARD'}]; error code: 409
[2026-07-23 15:19:40,260] INFO [Producer clientId=console-producer] Closing the Kafka producer with timeoutMillis = 9223372036854775807 ms. (org.apache.kafka.clients.producer.KafkaProducer)
[2026-07-23 15:19:40,595] INFO [Producer clientId=console-producer] Cluster ID: O1N_F_D-S2K6V8-A_G-p-A (org.apache.kafka.clients.Metadata)
[2026-07-23 15:19:40,596] INFO [Producer clientId=console-producer] ProducerId set to 6 with epoch 0 (org.apache.kafka.clients.producer.internals.TransactionManager)
[2026-07-23 15:19:40,604] INFO Metrics scheduler closed (org.apache.kafka.common.metrics.Metrics)
[2026-07-23 15:19:40,604] INFO Closing reporter org.apache.kafka.common.metrics.JmxReporter (org.apache.kafka.common.metrics.Metrics)
[2026-07-23 15:19:40,604] INFO Metrics reporters closed (org.apache.kafka.common.metrics.Metrics)
[2026-07-23 15:19:40,604] INFO App info kafka.producer for console-producer unregistered (org.apache.kafka.common.utils.AppInfoParser)

 

 

맨 하단 부분에 잘 보면

 

org.apache.kafka.common.errors.InvalidConfigurationException: Schema being registered is incompatible with an earlier schema for subject "schema-test-topic-value"

 

이런게 있다. 호환이 안된다는 말이다. 왜? 

 

errorType: 'READER_FIELD_MISSING_DEFAULT_VALUE' description: The field 'age' at path '/fields/1' in the new schema has no default value and is missing in the old schema oldSchemaVersion: 1 oldSchema: {"fields":[{"name":"name","type":"string"}]} compatibility: 'BACKWARD' error code: 409

 

새 스키마의 age필드는 default가 없는데, 옛 스키마(v1)에는 그 필드가 아예 없다.
그러니 v2로 v1 데이터를 읽을 때 age를 채울 방법이 없다. 고로 거부!

 

 

5.5.3) 호환성에 맞게 스키마 수정 후 정상 처리

 

docker exec -i schema-registry kafka-avro-console-producer \
  --bootstrap-server kafka-kraft-1:9092 \
  --topic schema-test-topic \
  --property schema.registry.url=http://schema-registry:8081 \
  --property value.schema='{"type":"record","name":"User","fields":[{"name":"name","type":"string"},{"name":"age","type":"int","default":0}]}' <<EOF
{"name":"Kim","age":30}
EOF

 

 

컨슈머에서 보면 기존 컨슈머는 중단 없이 새 데이터를 읽는다.

 

 

왼쪽에서 {"name":"Kim","age":30} 가 보이는 것을 확인했다.


 

 

어휴... 찾아보고 읽느라 진땀 뺏다.

 

해당 글을 정리하는데 너무 오랜 기간이 걸렸다. 무엇이 이렇게 어려운건지... 

 

무엇보다 이게 정확하게 어떻게 돌아가는지 머릿속으로 이해가 안됏다. 

 

호환성, 그것이 정확하게 어디서 이루어 지는건데? 하는 부분을 이해하기 위해서 하루종일 고생을 했다.

 

무엇보다도 kafka... 그냥 분산 메시지 스트리밍 서비스인줄 알앗는데... 기능이 엄청 많다.

 

역시 데이터 파이프라인 플랫폼의 최강자 답다. 내가 정말 내 두 손 다 든다 증말!

 

참고 글 정리:

요기오 - Confluent Schema Registry 도입기!

Kafka 와 Confluent Schema Registry 를 사용한 스키마 관리 #1

Kafka 와 Confluent Schema Registry 를 사용한 스키마 관리 #2

Apache Avro의 Specification

Live 11 과 Schema Registry