스키마
스키마 (Schema) 개요
Pulsar 메시지는 구조화되지 않은 바이트 배열로 저장되고, 데이터 구조(스키마라고도 부르죠)는 데이터를 읽을 때만 적용돼요. 그래서 프로듀서와 컨슈머가 메시지의 데이터 구조(필드와 그 타입)에 대해 서로 합의해야 해요. 이번에는 Pulsar 스키마가 무엇이고, 왜 필요한지, 어떻게 동작하는지 함께 살펴볼게요.
출처: 문서
본문
정의 (Definitions)
Pulsar 메시지는 구조화되지 않은 바이트 배열로 저장되고, 데이터 구조(스키마)는 데이터를 읽을 때만 적용돼요. 따라서 프로듀서와 컨슈머 모두 메시지의 데이터 구조(필드와 관련 타입)에 대해 합의해야 해요.
Pulsar 스키마는 원시 메시지 바이트를 더 공식적인 구조 타입으로 변환하는 방법을 정의하는 메타데이터예요. 메시지를 생성하는 애플리케이션과 소비하는 애플리케이션 사이의 프로토콜 역할을 하죠. 토픽에 게시되기 전에 데이터를 원시 바이트로 직렬화(serialize)하고, 컨슈머에게 전달되기 전에 원시 바이트를 역직렬화(deserialize)해요.
Pulsar는 스키마 레지스트리(schema registry)를 중앙 저장소로 사용해 등록된 스키마 정보를 보관해요. 이를 통해 프로듀서/컨슈머가 브로커를 거쳐 토픽 메시지의 스키마를 조정할 수 있어요.
note 현재 Pulsar 스키마는 Java 클라이언트, Go 클라이언트, Python 클라이언트, Node.js 클라이언트, C++ 클라이언트, C# 클라이언트에서 사용할 수 있어요.
이점 (Benefits)
타입 안정성은 메시징·스트리밍 시스템을 중심으로 구축된 모든 애플리케이션에서 매우 중요해요. 원시 바이트는 데이터 전송 측면에서 유연하지만, 그 유연성과 중립성에는 대가가 따르죠. 시스템에 들어가는 바이트가 읽히고 성공적으로 소비될 수 있도록 데이터 타입 검사와 직렬화/역직렬화를 직접 덧붙여야 해요. 즉, 데이터가 애플리케이션에 이해 가능하고 사용 가능하도록 보장해야 한다는 뜻이에요.
Pulsar 스키마는 다음 기능으로 이런 고민을 해결해요.
- 토픽에 스키마가 정의되면 데이터 타입 안전성을 강제해요. 그 결과 프로듀서/컨슈머는 "호환되는" 스키마를 사용할 때만 연결할 수 있어요.
- 조직 내에서 사용되는 스키마 정보를 저장하는 중앙 위치를 제공해, 이 정보를 애플리케이션 팀 간 공유하는 과정을 크게 단순화해요.
- 모든 서비스와 개발 팀에서 사용되는 모든 메시지 스키마의 단일 진실 공급원(single source of truth) 역할을 해서 협업을 더 쉽게 만들어요.
- 스키마 버전 간 데이터 호환성을 유지해요. 새 스키마가 업로드되면 옛 컨슈머도 새 버전을 읽을 수 있어요.
- 추가 시스템 없이 기존 저장 계층인 BookKeeper에 저장돼요.
워크플로 (Workflow)
Pulsar 스키마는 토픽 레벨에서 적용·강제돼요. 프로듀서와 컨슈머 모두 스키마를 브로커에 업로드할 수 있으므로, Pulsar 스키마는 양쪽 모두에서 동작해요.
프로듀서 측 (Producer side)
아래 다이어그램은 프로듀서 측에서 Pulsar 스키마가 어떻게 동작하는지 보여줘요.
각 단계에 대한 설명은 다음과 같아요.
- 애플리케이션이 스키마 인스턴스를 사용해 프로듀서 인스턴스를 구성해요. 스키마 인스턴스는 프로듀서 인스턴스로 생성되는 데이터의 스키마를 정의해요. Avro를 예로 들면, Pulsar는 POJO 클래스에서 스키마 정의를 추출해
SchemaInfo를 구성해요. - 프로듀서가 전달된 스키마 인스턴스에서 추출한
SchemaInfo로 브로커에 연결을 요청해요. - 브로커가 스키마 레지스트리를 조회해 등록된 스키마인지 확인해요.
- 등록된 스키마라면 브로커가 스키마 버전을 프로듀서에게 반환해요.
- 아니면 4단계로 이동해요.
- 브로커가 스키마를 자동 업데이트할 수 있는지 확인해요.
- 자동 업데이트가 허용되지 않으면 스키마를 등록할 수 없고, 브로커는 프로듀서를 거부해요.
- 아니면 5단계로 이동해요.
- 브로커가 토픽에 정의된 스키마 호환성 검사를 수행해요.
- 스키마가 호환성 검사를 통과하면 브로커가 스키마를 스키마 레지스트리에 저장하고 스키마 버전을 프로듀서에게 반환해요. 이 프로듀서가 생성한 모든 메시지에는 스키마 버전이 태그로 붙어요.
- 아니면 브로커가 프로듀서를 거부해요.
컨슈머 측 (Consumer side)
아래 다이어그램은 컨슈머 측에서 스키마가 어떻게 동작하는지 보여줘요.
각 단계에 대한 설명은 다음과 같아요.
- 애플리케이션이 스키마 인스턴스를 사용해 컨슈머 인스턴스를 구성해요.
- 컨슈머가 전달된 스키마 인스턴스에서 추출한
SchemaInfo로 브로커에 연결해요. - 브로커가 토픽이 사용 중인지(스키마, 데이터, 활성 프로듀서, 활성 컨슈머 중 하나 이상 보유) 확인해요.
- 토픽이 위 객체 중 하나 이상을 가지면 5단계로 이동해요.
- 아니면 4단계로 이동해요.
- 브로커가 스키마를 자동 업데이트할 수 있는지 확인해요.
- 스키마를 자동 업데이트할 수 있으면 브로커가 스키마를 등록하고 컨슈머를 연결해요.
- 아니면 브로커가 컨슈머를 거부해요.
- 브로커가 스키마 호환성 검사를 수행해요.
- 스키마가 호환성 검사를 통과하면 브로커가 컨슈머를 연결해요.
- 아니면 브로커가 컨슈머를 거부해요.
사용 사례 (Use case)
간단한 데이터 타입(예: string)부터 더 복잡한 애플리케이션 전용 타입까지, 메시지를 구성·처리할 때 언어별 데이터 타입을 사용할 수 있어요.
예를 들어 User 클래스를 사용해 Pulsar 토픽으로 보낼 메시지를 정의한다고 가정해요.
public class User {
public String name;
public int age;
User() {}
User(String name, int age) {
this.name = name;
this.age = age;
}
}
스키마 없이 (Without a schema)
스키마를 지정하지 않고 프로듀서를 구성하면, 프로듀서는 byte[] 타입의 메시지만 생성할 수 있어요. POJO 클래스가 있다면 메시지를 보내기 전에 POJO를 바이트로 직접 직렬화해야 해요.
Producer<byte[]> producer = client.newProducer()
.topic(topic)
.create();
User user = new User("Tom", 28);
byte[] message = … // serialize the `user` by yourself;
producer.send(message);
스키마와 함께 (With a schema)
이 예제는 JSONSchema로 프로듀서를 구성해요. POJO를 바이트로 직렬화하는 방법을 고민하지 않고 User 클래스를 토픽으로 바로 보낼 수 있어요.
// send with json schema
Producer<User> producer = client.newProducer(JSONSchema.of(User.class))
.topic(topic)
.create();
User user = new User("Tom", 28);
producer.send(user);
// receive with json schema
Consumer<User> consumer = client.newConsumer(JSONSchema.of(User.class))
.topic(schemaTopic)
.subscriptionInitialPosition(SubscriptionInitialPosition.Earliest)
.subscriptionName("schema-sub")
.subscribe();
Message<User> message = consumer.receive();
User user = message.getValue();
assert user.age == 28 && user.name.equals("Tom");
다음 단계 (What's next?)
더 알아보기 (Learn more)
- 스키마의 다양한 타입과 호환성 검사는 schema-understand 문서에서 확인해요.
- 실제로 스키마를 구성하는 방법은 Get started with schema를 봐요.
- 관리 작업을 위한 스키마 명령은 Manage schema 문서를 참고해요.
- 클라이언트 라이브러리별 스키마 지원은 client libraries 문서를 살펴봐요.