스키마 레지스트리란?
Producer와 Consumer가 주고 받는 메세지를 서로 알게 해주고 호환을 강제한다.
- 즉, 스키마를 만들어 저장, 관리하는 confluent의 웹 어플리케이션
클라이언트 사이에는(Producer와 Consumer) 메세지 구조에 대한 강한 결합도를 가지고 있다.
- 스키마 레지스트리는 이 결합도를 낮추기 위해 고안
내부적으로는 Avro를 사용하고 REST API를 이용해 Avro 스키마를 저장/조회 가능하며, 설정에 따라 상위, 하위, 양쪽 호환성을 보장한다.
특징

- Kafka 외부에서 독립적으로 동작하며 REST API를 제공한다.
- Producer와 Consumer는 Avro 포맷의 메세지를 Kafka를 통해 송수신 할 수 있도록 한다.
- Avro
- 시스템, 프로그래밍 언어, 프로세싱 프레임워크 사이에서 데이터 교환을 도와주는 오픈소스 직렬화 시스템
- Json과 달리 데이터 필드마다 데이터 타입을 정의 가능, doc을 이용해 각 필드의 의미를 데이터 사용자들에게 정확하게 전달할 수 있다.
- namespace(이름을 식별하는 문자열), type(record, enums, arrays 등을 지원), doc(사용자들에게 보여줄 주석), name(이름 문자열 필수값), field(json 배열, 필드의 리스트)등 구성 요소를 활용해 작성할 수 있다.
- avro는 필드 명 정보는 한번만 정의하고, 그 이후에는 값만 저장한다.
- Avro 방식을 가장 추천하는 이유는 JSON과 매핑되며 JSON과 다르게 필드 네임이 포함되어 전송되지 않고 바이너리 데이터 포맷 형태로 전송되어 속도 측면에서 가장 빠르다.
- Avro
- HTTP endpoint
- Topic : Kafka broker topic
- Schema : Avro data format
- 하나의 Topic에 하나의 Schema만 produce 할 수 있지만, 여러 schema를 하나의 topic에 넣고 consumer와 분리 할 수 있다.
- Topic과 Schema를 바로 매핑 하지 않고, 중간에 subject라는 개념을 두어 Topic과 schema 간 의존성을 느슨하게 함.
- 스키마 레지스트리에서 호환성 : 스키마 레지스트리는 호환성 설정에 따라 새로운 스키마를 등록할 수 있는지 없는지 판단한다.
- Backward(기본값)
- 스키마 업데이트 순서 : Consumer → Producer
- backward_transitive : 자신과 동일한 버전을 포함한 모든 하위 버전
- 버전 업데이트 된 스키마를 적용한 Consumer가 버전 업데이트가 안된 스키마를 적용한 Producer가 보낸 메시지를 읽을 수 있도록 허용하는 호환성
- Consumer는 새로운 스키마를 사용해서 처리하지만, 가장 최근에 등록된 스키마를 사용하여 처리할 수 있다.
- Field 삭제가 가능하고, 기본 값이 지정된 필드를 추가할 수 있다.
- 이전 스키마를 사용하는 소비자가 새 스키마를 사용해 생성된 데이터를 읽을 수 있다는 보장은 없다. 새 이벤트 생성을 하기 전에 모든 소비자를 업데이트 해야 한다.
- Forward
- 스키마 업데이트 순서 : Producer → Consumer
- 버전 업데이트 된 스키마를 적용한 Producer가 버전 업데이트가 안된 스키마를 적용한 컨슈머가 보낸 메세지를 읽을 수 있도록 허용하는 호환성
- forward_transitive : 자신과 동일한 버전을 포함한 모든 상위 버전
- 컨슈머는 마지막으로 등록된 스키마를 사용해서 처리하지만, 새로운 스키마도 처리 가능하다.
- Field 추가가 가능하고, 기본값이 지정된 Field를 삭제할 수 있다.
- 새 스키마를 사용하는 소비자가 이전 스키마를 사용하여 생성된 데이터를 읽을 수 있다는 보장이 없다.
- 모든 생산자를 새 스키마 사용하도록 업데이트 한 후 업그레이드 한다.
- Full
- 새 스키마는 가장 최근에 등록된 forward, backward 호환이 가능하다(버전 기준 +1, -1)
- 기본 값이 지정된 Field 추가와 삭제가 가능하다
- 새 스키마를 사용하는 소비자는 이전 스키마를 사용하여 생성된 데이터를 읽을 수 있다는 보장이 있다. 생산자와 스키마를 독립적으로 업그레이드 할 수 있다.
- None
- 스키마 호환성을 확인하지 않는다.
- Backward(기본값)
동작 원리
Kafka에서 Avro를 사용하여 데이터 형식을 일관되게 유지할 수 있다. 스키마 레지스트리는 클라이언트가 공유 하는 중앙 스키마 저장소로 Avro 스키마를 등록하고 버전 관리할 수 있다.
- 클라이언트는 스키마 레지스트리를 통해 메세지의 스키마를 검색해 데이터를 직렬화 역직렬화 할 수 있다.
- 로컬 캐시에 없으면 스키마 전송 등록
- Confluent에서 제공하는 io.confluent.kafka.serializers.kafkaAvroserializer 라는 새로운 직렬화를 사용해 전송 하려는 메세지의 스키마 타입이 로컬 캐시에 존재하는지 확인
- 로컬 캐시에 존재 → 로컬 캐시에 등록된 정보를 사용 로컬 캐시에 존재하지 않음 → 스키마 레지스트리에 스키마가 존재하는지 여부를 확인
- 스키마 레지스트리에 스키마가 확인되지 않을 경우 Avro Producer 는 스키마를 스키마 레지스트리에 등록하고 캐시
- 스키마 레지스트리는 저장된 스키마의 정보를 카프카 내부 토픽에 저장
- 보유한 스키마와 동일한지 확인
- 스키마 레지스트리에선 프로듀서가 요청한 스키마 id의 버전과 스키마의 버전이 동일한지를 확인
- 스키마 레지스트리는 자체적으로 각 스키마에 고유 ID를 할당 → ID는 순차적으로 증가하긴 하는데, 연속적이지는 않음. (스키마가 삭제되면 그대로 빵꾸로 남기 때문)
- 스키마 레지스트리는 프로듀서에게 고유 ID를 응답
- 스키마 id 별 avro 메세지 직렬화 전송
- Producer는 스키마 레지스트리로부터 밭은 schema id를 활용해서 메세지를 Kafka로 전송 (Producer는 스키마의 전체 내용이 아닌 오로지 메세지와 schema id만 전송)
- Json은 k:v 형태로 전체 메시지를 전송하지만, avro는 producer가 schema id와 value만 메세지로 보내게 되어 kafka로 전송하는 전체 메세지 크기를 줄일 수 있다. (json보다 avro를 사용하는 편이 더 효율적인 이유)
- 스키마 id 별 avro 메세지 역직렬화 전송
- Consumer는 schema id로 io.confluent.kafka.serializer.KafkaAvroDeserializer 라는 새로운 역직렬화를 사용해서 Kafka에 저장된 메세지를 읽는다.
- 로컬캐시에 없다면 스키마 id 별도 조회
- 컨슈머가 schema id를 가지고 있지 않다면 스키마 레지스트리로부터 가져온다.

- avro의 header에 schema.id를 같이 작성해서 관리하는 형식이다.
- 스키마 레지스트레에는 스키마 정보를, 토픽에는 스키마 id와 avro 포맷의 데이터를 보내게 된다.
적용 방법
예시 : docker-compose 기준으로 작성
version: '3.6'
services:
kafka-schema-registry:
image: confluentinc/cp-schema-registry:6.2.1
container_name: full-kafka-schema-registry
hostname: kafka-schema-registry
ports:
- "8081:8081"
environment:
SCHEMA_REGISTRY_KAFKASTORE_CONNECTION_URL: zookeeper:2181
SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: PLAINTEXT://kafka:9092
SCHEMA_REGISTRY_HOST_NAME: kafka-schema-registry
SCHEMA_REGISTRY_LISTENERS: <http://0.0.0.0:8081>
depends_on:
- zookeeper
- kafka
kafka-rest-proxy:
image: confluentinc/cp-kafka-rest:6.2.1
container_name: full-kafka-rest-proxy
hostname: kafka-rest-proxy
ports:
- "8082:8082"
environment:
# KAFKA_REST_ZOOKEEPER_CONNECT: zoo1:2181
KAFKA_REST_LISTENERS: <http://0.0.0.0:8082/>
KAFKA_REST_SCHEMA_REGISTRY_URL: <http://kafka-schema-registry:8081/>
KAFKA_REST_HOST_NAME: kafka-rest-proxy
KAFKA_REST_BOOTSTRAP_SERVERS: PLAINTEXT://kafka1:19092
depends_on:
- zookeeper
- kafka
- kafka-schema-registry
- 스키마 레지스트리 api 리스트
- GET /schemas : 전체 스키마 리스트 조회
- GET /schemas/ids/id : 스키마 아이디로 조회
- GET /schemas/lds/id//versions : 스키마 id 버전
- GET /subjects : subject리스트
- GET /subjects/서브젝트 이름/versions : 특정 서브 젝트 버전 리스트 조회
- GET /config : 전역으로 설정된 호환성 레벨 조회
- GET /config/서브젝트 이름 : 서브젝트에 설정된 호환성 조회
- DELETE/subjects/서브젝트 이름 :특정 서브젝트 전체 삭제
- DELETE/subjects/서브젝트 이름/versions/버전 : 특정 서브젝트에서 특정 버전만 삭제
참고한 글들
늘 좋은 글 감사합니다..!
https://velog.io/@fj2008/카프카-스키마-레지스트리https://velog.io/@fj2008/카프카-스키마-레지스트리
https://medium.com/@gaemi/kafka-와-confluent-schema-registry-를-사용한-스키마-관리-1-cdf8c99d2c5c
https://always-kimkim.tistory.com/entry/kafka101-schema-registry
https://blog.voidmainvoid.net/463
https://data-engineer-tech.tistory.com/39
https://devidea.tistory.com/113
https://docs.confluent.io/current/schema-registry/index.html
'Hadoop eco' 카테고리의 다른 글
| [Kafka] Kafka DR(Disaster Recovery) (0) | 2023.07.14 |
|---|---|
| [Kafka] 성능 측정 지표 (0) | 2023.05.23 |
| [Kafka] Message Delivery Semantics (0) | 2023.05.23 |
| [Kafka] 기업 도입 사례들 (0) | 2023.05.23 |
| [Kafka] Kafka 브로커의 동작 (0) | 2023.02.11 |