스키마 시작하기

스키마 시작하기 (Get started)

이 실습형 튜토리얼에서는 스키마를 구성하는 방법을 안내와 예제로 설명해요. 관리 작업에 대한 설명은 스키마 관리 (Manage schema) 문서를 참고해요. 여기서는 언어별 클라이언트로 각 스키마 타입을 만들고 메시지를 주고받는 예제를 함께 해볼게요.

출처: 문서

본문

스키마 구성 (Construct a schema)

bytes

이 예제는 언어별 클라이언트로 bytes 스키마를 구성하고 메시지를 생성·소비하는 방법을 보여줘요.

  • Java
  • C++
  • Python
  • Go

Java:

Producer<byte[]> producer = pulsarClient.newProducer(Schema.BYTES)
       .topic("my-topic")
       .create();
Consumer<byte[]> consumer = pulsarClient.newConsumer(Schema.BYTES)
       .topic("my-topic")
       .subscriptionName("my-sub")
       .subscribe();
producer.newMessage().value("message".getBytes()).send();
Message<byte[]> message = consumer.receive(5, TimeUnit.SECONDS);

C++:

SchemaInfo schemaInfo = SchemaInfo(SchemaType::BYTES, "Bytes", "");
Producer producer;
client.createProducer("topic-bytes", ProducerConfiguration().setSchema(schemaInfo), producer);
std::array<char, 1024> buffer;
producer.send(MessageBuilder().setContent(buffer.data(), buffer.size()).build());
Consumer consumer;
res = client.subscribe("topic-bytes", "my-sub", ConsumerConfiguration().setSchema(schemaInfo), consumer);
Message msg;
consumer.receive(msg, 3000);

Python:

producer = client.create_producer(
   'bytes-schema-topic',
   schema=BytesSchema())
producer.send(b"Hello")
consumer = client.subscribe(
   'bytes-schema-topic',
	'sub',
	schema=BytesSchema())
msg = consumer.receive()
data = msg.value()

Go:

producer, err := client.CreateProducer(pulsar.ProducerOptions{
    Topic:  "my-topic",
    Schema: pulsar.NewBytesSchema(nil),
})
id, err := producer.Send(context.Background(), &pulsar.ProducerMessage{
    Value: []byte("message"),
})
consumer, err := client.Subscribe(pulsar.ConsumerOptions{
    Topic:            "my-topic",
    Schema:           pulsar.NewBytesSchema(nil),
    SubscriptionName: "my-sub",
    Type:             pulsar.Exclusive,
})

string

이 예제는 언어별 클라이언트로 string 스키마를 구성하고 메시지를 생성·소비하는 방법을 보여줘요.

  • Java
  • C++
  • Python
  • Go

Java:

Producer<String> producer = client.newProducer(Schema.STRING).create();
producer.newMessage().value("Hello Pulsar!").send();
Consumer<String> consumer = client.newConsumer(Schema.STRING).subscribe();
Message<String> message = consumer.receive();

C++:

SchemaInfo schemaInfo = SchemaInfo(SchemaType::STRING, "String", "");
Producer producer;
client.createProducer("topic-string", ProducerConfiguration().setSchema(schemaInfo), producer);
producer.send(MessageBuilder().setContent("message").build());
Consumer consumer;
client.subscribe("topic-string", "my-sub", ConsumerConfiguration().setSchema(schemaInfo), consumer);
Message msg;
consumer.receive(msg, 3000);

Python:

producer = client.create_producer(
      'string-schema-topic',
      schema=StringSchema())
producer.send("Hello")
consumer = client.subscribe(
		'string-schema-topic',
		'sub',
		schema=StringSchema())
msg = consumer.receive()
str = msg.value()

Go:

producer, err := client.CreateProducer(pulsar.ProducerOptions{
    Topic:  "my-topic",
    Schema: pulsar.NewStringSchema(nil),
})
id, err := producer.Send(context.Background(), &pulsar.ProducerMessage{
    Value: "message",
})
consumer, err := client.Subscribe(pulsar.ConsumerOptions{
    Topic:            "my-topic",
    Schema:           pulsar.NewStringSchema(nil),
    SubscriptionName: "my-sub",
    Type:             pulsar.Exclusive,
})
msg, err := consumer.Receive(context.Background())

key/value

이 예제는 언어별 클라이언트로 key/value 스키마를 구성하고 메시지를 생성·소비하는 방법을 보여줘요.

  • Java
  • C++
  1. INLINE 인코딩 타입으로 key/value 스키마를 구성해요.
Schema<KeyValue<Integer, String>> kvSchema = Schema.KeyValue(
    Schema.INT32,
    Schema.STRING,
    KeyValueEncodingType.INLINE);

또는 SEPARATED 인코딩 타입으로 key/value 스키마를 구성할 수도 있어요.

Schema<KeyValue<Integer, String>> kvSchema = Schema.KeyValue(
    Schema.INT32,
    Schema.STRING,
    KeyValueEncodingType.SEPARATED);
  1. key/value 스키마를 사용해 메시지를 생성해요.
Producer<KeyValue<Integer, String>> producer = client.newProducer(kvSchema)
    .topic(topicName)
    .create();
final int key = 100;
final String value = "value-100";
// send the key/value message
producer.newMessage()
    .value(new KeyValue(key, value))
    .send();
  1. key/value 스키마를 사용해 메시지를 소비해요.
Consumer<KeyValue<Integer, String>> consumer = client.newConsumer(kvSchema)
    ...
    .topic(topicName)
    .subscriptionName(subscriptionName).subscribe();
// receive key/value pair
Message<KeyValue<Integer, String>> msg = consumer.receive();
KeyValue<Integer, String> kv = msg.getValue();
  1. INLINE 인코딩 타입으로 key/value 스키마를 구성해요.
//Prepare keyValue schema
std::string jsonSchema =
R"({"type":"record","name":"cpx","fields":[{"name":"re","type":"double"},{"name":"im","type":"double"}]})";
SchemaInfo keySchema(JSON, "key-json", jsonSchema);
SchemaInfo valueSchema(JSON, "value-json", jsonSchema);
SchemaInfo keyValueSchema(keySchema, valueSchema, KeyValueEncodingType::INLINE);
  1. key/value 스키마를 사용해 메시지를 생성해요.
//Create Producer
Producer producer;
client.createProducer("my-topic", ProducerConfiguration().setSchema(keyValueSchema), producer);
//Prepare message
std::string jsonData = "{\"re\":2.1,\"im\":1.23}";
KeyValue keyValue(std::move(jsonData), std::move(jsonData));
Message msg = MessageBuilder().setContent(keyValue).setProperty("x", "1").build();
//Send message
producer.send(msg);
  1. key/value 스키마를 사용해 메시지를 소비해요.
//Create Consumer
Consumer consumer;
client.subscribe("my-topic", "my-sub", ConsumerConfiguration().setSchema(keyValueSchema), consumer);
//Receive message
Message message;
consumer.receive(message);

Avro

  • Java
  • C++
  • Python
  • Go

다음과 같은 SensorReading 클래스가 있고 이를 Pulsar 토픽으로 전송한다고 가정해요.

public class SensorReading {
    public float temperature;
    public SensorReading(float temperature) {
        this.temperature = temperature;
    }
    // A no-arg constructor is required
    public SensorReading() {
    }
    public float getTemperature() {
        return temperature;
    }
    public void setTemperature(float temperature) {
        this.temperature = temperature;
    }
}

다음과 같이 Producer<SensorReading>(또는 Consumer<SensorReading>)을 만들어요.

Java:

Producer<SensorReading> producer = client.newProducer(AvroSchema.of(SensorReading.class))
        .topic("sensor-readings")
        .create();

C++:

// Send messages
static const std::string exampleSchema =
    "{\"type\":\"record\",\"name\":\"Example\",\"namespace\":\"test\","
    "\"fields\":[{\"name\":\"a\",\"type\":\"int\"},{\"name\":\"b\",\"type\":\"int\"}]}";
Producer producer;
ProducerConfiguration producerConf;
producerConf.setSchema(SchemaInfo(AVRO, "Avro", exampleSchema));
client.createProducer("topic-avro", producerConf, producer);
// Receive messages
static const std::string exampleSchema =
    "{\"type\":\"record\",\"name\":\"Example\",\"namespace\":\"test\","
    "\"fields\":[{\"name\":\"a\",\"type\":\"int\"},{\"name\":\"b\",\"type\":\"int\"}]}";
ConsumerConfiguration consumerConf;
Consumer consumer;
consumerConf.setSchema(SchemaInfo(AVRO, "Avro", exampleSchema));
client.subscribe("topic-avro", "sub-2", consumerConf, consumer)

Python에서는 다음 방법 중 하나로 AvroSchema를 선언할 수 있어요.

방법 1: Record pulsar.schema.Record를 상속하고 필드를 클래스 변수로 정의하는 클래스를 전달해 AvroSchema를 선언해요.

class Example(Record):
    a = Integer()
    b = Integer()
producer = client.create_producer(
   'avro-schema-topic',
   schema=AvroSchema(Example))
r = Example(a=1, b=2)
producer.send(r)
consumer = client.subscribe(
   'avro-schema-topic',
	'sub',
	schema=AvroSchema(Example))
msg = consumer.receive()
e = msg.value()

방법 2: JSON 정의

  1. JSON을 사용해 AvroSchema를 선언해요. 이 경우 Avro 스키마는 JSON으로 정의돼요. 아래는 JSON 파일(company.avsc)로 정의한 AvroSchema 예제예요.
{
    "doc": "this is doc",
    "namespace": "example.avro",
    "type": "record",
    "name": "Company",
    "fields": [
        {"name": "name", "type": ["null", "string"]},
        {"name": "address", "type": ["null", "string"]},
        {"name": "employees", "type": ["null", {"type": "array", "items": {
            "type": "record",
            "name": "Employee",
            "fields": [
                {"name": "name", "type": ["null", "string"]},
                {"name": "age", "type": ["null", "int"]}
            ]
        }}]},
        {"name": "labels", "type": ["null", {"type": "map", "values": "string"}]}
    ]
}
  1. avro.schema 또는 fastavro.schema를 사용해 파일에서 스키마 정의를 불러와요.

JSON 정의 방법으로 AvroSchema를 선언하려면 다음을 해야 해요.

  • Record 방법과 달리 Python dict를 사용해 메시지를 생성·소비해요.
  • AvroSchema 객체를 생성할 때 _record_cls 파라미터 값을 None으로 설정해요.

예제

from fastavro.schema import load_schema
from pulsar.schema import *
schema_definition = load_schema("examples/company.avsc")
avro_schema = AvroSchema(None, schema_definition=schema_definition)
producer = client.create_producer(
    topic=topic,
    schema=avro_schema)
consumer = client.subscribe(topic, 'test', schema=avro_schema)
company = {
    "name": "company-name" + str(i),
    "address": 'xxx road xxx street ' + str(i),
    "employees": [
        {"name": "user" + str(i), "age": 20 + i},
        {"name": "user" + str(i), "age": 30 + i},
        {"name": "user" + str(i), "age": 35 + i},
    ],
    "labels": {
        "industry": "software" + str(i),
        "scale": ">100",
        "funds": "1000000.0"
    }
}
producer.send(company)
msg = consumer.receive()
# Users could get a dict object by `value()` method.
msg.value()

다음과 같은 avroExampleStruct 클래스가 있고 이를 Pulsar 토픽으로 전송한다고 가정해요.

type avroExampleStruct struct {
    ID   int
    Name string
}
  1. 다음과 같이 avroSchemaDef를 추가해요.
var (
    exampleSchemaDef = "{\"type\":\"record\",\"name\":\"Example\",\"namespace\":\"test\"," +
  "\"fields\":[{\"name\":\"ID\",\"type\":\"int\"},{\"name\":\"Name\",\"type\":\"string\"}]}"
)
  1. 프로듀서와 컨슈머를 만들어 메시지를 보내고 받아요.
//Create producer and send message
producer, err := client.CreateProducer(pulsar.ProducerOptions{
    Topic:  "my-topic",
    Schema: pulsar.NewAvroSchema(exampleSchemaDef, nil),
})
msgId, err := producer.Send(context.Background(), &pulsar.ProducerMessage{
    Value: avroExampleStruct{
       ID:   10,
       Name: "avroExampleStruct",
   },
})
//Create Consumer and receive message
consumer, err := client.Subscribe(pulsar.ConsumerOptions{
    Topic:            "my-topic",
    Schema:           pulsar.NewAvroSchema(exampleSchemaDef, nil),
    SubscriptionName: "my-sub",
    Type:             pulsar.Shared,
})
message, err := consumer.Receive(context.Background())

JSON

  • Java
  • C++
  • Python
  • Go

AvroSchema 사용법과 유사하게, 클래스를 전달해 JsonSchema를 선언할 수 있어요. 유일한 차이는 스키마 타입을 정의할 때 AvroSchema 대신 JsonSchema를 사용한다는 점이에요. 아래에서 확인할 수 있어요. record로 AvroSchema를 사용하는 방법은 방법 1 - Record를 참고해요.

Java:

static class SchemaDemo {
   public String name;
   public int age;
}
Producer<SchemaDemo> producer = pulsarClient.newProducer(Schema.JSON(SchemaDemo.class))
       .topic("my-topic")
       .create();
Consumer<SchemaDemo> consumer = pulsarClient.newConsumer(Schema.JSON(SchemaDemo.class))
       .topic("my-topic")
       .subscriptionName("my-sub")
       .subscribe();
SchemaDemo schemaDemo = new SchemaDemo();
schemaDemo.name = "puslar";
schemaDemo.age = 20;
producer.newMessage().value(schemaDemo).send();
Message<SchemaDemo> message = consumer.receive(5, TimeUnit.SECONDS);

C++에서 JSON 스키마를 선언하려면 다음을 해요.

  1. 다음과 같이 JSON 문자열을 전달해요.
Std::string jsonSchema = R"({"type":"record","name":"cpx","fields":[{"name":"re","type":"double"},{"name":"im","type":"double"}]})";
SchemaInfo schemaInfo = SchemaInfo(JSON, "JSON", jsonSchema);
  1. 프로듀서를 만들고 메시지를 보내요.
client.createProducer("my-topic", ProducerConfiguration().setSchema(schemaInfo), producer);
std::string jsonData = "{\"re\":2.1,\"im\":1.23}";
producer.send(MessageBuilder().setContent(std::move(jsonData)).build());
  1. 컨슈머를 만들고 메시지를 받아요.
Consumer consumer;
client.subscribe("my-topic", "my-sub", ConsumerConfiguration().setSchema(schemaInfo), consumer);
Message msg;
consumer.receive(msg);

pulsar.schema.Record를 상속하고 필드를 클래스 변수로 정의하는 클래스를 전달해 JsonSchema를 선언할 수 있어요. 이는 AvroSchema 사용법과 유사해요. 유일한 차이는 스키마 타입을 정의할 때 AvroSchema 대신 JsonSchema를 사용한다는 점이에요. 아래에서 확인할 수 있어요. record로 AvroSchema를 사용하는 방법은 (#method-1-record)을 참고해요.

producer = client.create_producer(
   'avro-schema-topic',
   schema=JsonSchema(Example))
consumer = client.subscribe(
	'avro-schema-topic',
	'sub',
	schema=JsonSchema(Example))

다음과 같은 avroExampleStruct 클래스가 있고 이를 JSON 형태로 Pulsar 토픽에 전송한다고 가정해요.

type jsonExampleStruct struct {
    ID   int    `json:"id"`
    Name string `json:"name"`
}
  1. 다음과 같이 jsonSchemaDef를 추가해요.
   jsonSchemaDef = "{\"type\":\"record\",\"name\":\"Example\",\"namespace\":\"test\"," +
  "\"fields\":[{\"name\":\"ID\",\"type\":\"int\"},{\"name\":\"Name\",\"type\":\"string\"}]}"
  1. 프로듀서/컨슈머를 만들어 메시지를 보내고 받아요.
   //Create producer and send message
   producer, err := client.CreateProducer(pulsar.ProducerOptions{
       Topic:  "my-topic",
       Schema: pulsar.NewJSONSchema(jsonSchemaDef, nil),
   })
   msgId, err := producer.Send(context.Background(), &pulsar.ProducerMessage{
       Value: jsonExampleStruct{
           ID:   10,
           Name: "jsonExampleStruct",
     },
   })
   //Create Consumer and receive message
   consumer, err := client.Subscribe(pulsar.ConsumerOptions{
       Topic:            "my-topic",
       Schema:           pulsar.NewJSONSchema(jsonSchemaDef, nil),
       SubscriptionName: "my-sub",
       Type:             pulsar.Exclusive,
   })
   message, err := consumer.Receive(context.Background())

ProtobufNative

  • Java
  • C++ 다음 예제는 Java로 ProtobufNative 스키마를 사용해 프로듀서/컨슈머를 만드는 방법을 보여줘요.
  1. Protobuf3 이상 버전으로 DemoMessage 클래스를 생성해요.
syntax = "proto3";
message DemoMessage {
   string stringField = 1;
   double doubleField = 2;
   int32 intField = 6;
   TestEnum testEnum = 4;
   SubMessage nestedField = 5;
   repeated string repeatedField = 10;
   proto.external.ExternalMessage externalMessage = 11;
}
  1. 프로듀서/컨슈머를 만들어 메시지를 보내고 받아요.
Producer<DemoMessage> producer = pulsarClient.newProducer(Schema.PROTOBUF_NATIVE(DemoMessage.class))
    .topic("my-topic")
    .create();
Consumer<DemoMessage> consumer = pulsarClient.newConsumer(Schema.PROTOBUF_NATIVE(DemoMessage.class))
    .topic("my-topic")
    .subscriptionName("my-sub")
    .subscribe();
SchemaDemo schemaDemo = new SchemaDemo();
schemaDemo.name = "puslar";
schemaDemo.age = 20;
producer.newMessage().value(DemoMessage.newBuilder().setStringField("string-field-value")
    .setIntField(1).build()).send();
Message<DemoMessage> message = consumer.receive(5, TimeUnit.SECONDS);

다음 예제는 ProtobufNative 스키마로 프로듀서/컨슈머를 만드는 방법을 보여줘요.

  1. Protobuf3 이상 버전으로 User 클래스를 생성해요.
syntax = "proto3";
message User {
    string name = 1;
    int32 age = 2;
}
  1. 소스 코드에 ProtobufNativeSchema.h를 포함해요. 프로젝트에 Protobuf 의존성이 추가되어 있는지 확인해요.
#include <pulsar/ProtobufNativeSchema.h>
  1. 프로듀서를 만들어 User 인스턴스를 보내요.
ProducerConfiguration producerConf;
producerConf.setSchema(createProtobufNativeSchema(User::GetDescriptor()));
Producer producer;
client.createProducer("topic-protobuf", producerConf, producer);
User user;
user.set_name("my-name");
user.set_age(10);
std::string content;
user.SerializeToString(&content);
producer.send(MessageBuilder().setContent(content).build());
  1. 컨슈머를 만들어 User 인스턴스를 받아요.
ConsumerConfiguration consumerConf;
consumerConf.setSchema(createProtobufNativeSchema(User::GetDescriptor()));
consumerConf.setSubscriptionInitialPosition(InitialPositionEarliest);
Consumer consumer;
client.subscribe("topic-protobuf", "my-sub", consumerConf, consumer);
Message msg;
consumer.receive(msg);
User user2;
user2.ParseFromArray(msg.getData(), msg.getLength());

Protobuf

  • Java
  • C++
  • Go

Java로 protobuf 스키마를 구성하는 것은 ProtobufNative 스키마 구성과 유사해요. 유일한 차이는 스키마 타입을 정의할 때 PROTOBUF_NATIVE 대신 PROTOBUF를 사용한다는 점이에요. 아래에서 확인할 수 있어요.

  1. Protobuf3 이상 버전으로 DemoMessage 클래스를 생성해요.
syntax = "proto3";
message DemoMessage {
   string stringField = 1;
   double doubleField = 2;
   int32 intField = 6;
   TestEnum testEnum = 4;
   SubMessage nestedField = 5;
   repeated string repeatedField = 10;
   proto.external.ExternalMessage externalMessage = 11;
}
  1. 프로듀서/컨슈머를 만들어 메시지를 보내고 받아요.
Producer<DemoMessage> producer = pulsarClient.newProducer(Schema.PROTOBUF(DemoMessage.class))
       .topic("my-topic")
       .create();
Consumer<DemoMessage> consumer = pulsarClient.newConsumer(Schema.PROTOBUF(DemoMessage.class))
       .topic("my-topic")
       .subscriptionName("my-sub")
       .subscribe();
SchemaDemo schemaDemo = new SchemaDemo();
schemaDemo.name = "puslar";
schemaDemo.age = 20;
producer.newMessage().value(DemoMessage.newBuilder().setStringField("string-field-value")
    .setIntField(1).build()).send();
Message<DemoMessage> message = consumer.receive(5, TimeUnit.SECONDS);

C++로 protobuf 스키마를 구성하는 것은 JSON 사용법과 유사해요. 유일한 차이는 스키마 타입을 정의할 때 JSON 대신 PROTOBUF를 사용한다는 점이에요. 아래에서 확인할 수 있어요.

std::string jsonSchema =
  R"({"type":"record","name":"cpx","fields":[{"name":"re","type":"double"},{"name":"im","type":"double"}]})";
SchemaInfo schemaInfo = SchemaInfo(pulsar::PROTOBUF, "PROTOBUF", jsonSchema);
  1. 프로듀서를 만들어 메시지를 보내요.
Producer producer;
client.createProducer("my-topic", ProducerConfiguration().setSchema(schemaInfo), producer);
std::string jsonData = "{\"re\":2.1,\"im\":1.23}";
producer.send(MessageBuilder().setContent(std::move(jsonData)).build());
  1. 컨슈머를 만들어 메시지를 받아요.
Consumer consumer;
client.subscribe("my-topic", "my-sub", ConsumerConfiguration().setSchema(schemaInfo),
  consumer);
Message msg;
consumer.receive(msg);

다음과 같은 protobufDemo 클래스가 있고 이를 JSON 형태로 Pulsar 토픽에 전송한다고 가정해요.

type protobufDemo struct {
    Num                  int32    `protobuf:"varint,1,opt,name=num,proto3" json:"num,omitempty"`
    Msf                  string   `protobuf:"bytes,2,opt,name=msf,proto3" json:"msf,omitempty"`
    XXX_NoUnkeyedLiteral struct{} `json:"-"`
    XXX_unrecognized     []byte   `json:"-"`
    XXX_sizecache        int32    `json:"-"`
}
  1. 다음과 같이 protoSchemaDef를 추가해요.
var (
    protoSchemaDef = "{\"type\":\"record\",\"name\":\"Example\",\"namespace\":\"test\"," +
        "\"fields\":[{\"name\":\"num\",\"type\":\"int\"},{\"name\":\"msf\",\"type\":\"string\"}]}"
)
  1. 프로듀서/컨슈머를 만들어 메시지를 보내고 받아요.
psProducer := pulsar.NewProtoSchema(protoSchemaDef, nil)
producer, err := client.CreateProducer(pulsar.ProducerOptions{
    Topic:  "proto",
    Schema: psProducer,
})
msgId, err := producer.Send(context.Background(), &pulsar.ProducerMessage{
    Value: &protobufDemo{
        Num: 100,
        Msf: "pulsar",
  },
})
psConsumer := pulsar.NewProtoSchema(protoSchemaDef, nil)
consumer, err := client.Subscribe(pulsar.ConsumerOptions{
    Topic:                       "proto",
    SubscriptionName:            "sub-1",
    Schema:                      psConsumer,
    SubscriptionInitialPosition: pulsar.SubscriptionPositionEarliest,
})
msg, err := consumer.Receive(context.Background())

Native Avro

이 예제는 native Avro 스키마를 구성하는 방법을 보여줘요.

org.apache.avro.Schema nativeAvroSchema = … ;
Producer<byte[]> producer = pulsarClient.newProducer().topic("ingress").create();
byte[] content = … ;
producer.newMessage(Schema.NATIVE_AVRO(nativeAvroSchema)).value(content).send();

AUTO_PRODUCE

Pulsar 토픽 P가 있고, Kafka 토픽 K의 메시지를 처리하는 프로듀서가 있으며, 애플리케이션이 K의 메시지를 읽어 P에 쓰는 상황을 가정해요.

이 예제는 AUTO_PRODUCE 스키마를 구성해 K가 생성한 바이트를 P로 보낼 수 있는지 확인하는 방법을 보여줘요.

Produce<byte[]> pulsarProducer = client.newProducer(Schema.AUTO_PRODUCE_BYTES())
    …
    .create();
byte[] kafkaMessageBytes = … ;
pulsarProducer.produce(kafkaMessageBytes);

AUTO_CONSUME

Pulsar 토픽 PP에서 메시지를 받는 컨슈머 MySQL이 있고, 이 메시지에 애플리케이션이 집계해야 하는 정보가 있는지 확인하려는 상황을 가정해요.

이 예제는 AUTO_CONSUME 스키마를 구성해 P가 생성한 바이트를 MySQL로 보낼 수 있는지 확인하는 방법을 보여줘요.

Consumer<GenericRecord> pulsarConsumer = client.newConsumer(Schema.AUTO_CONSUME())
    …
    .subscribe();
Message<GenericRecord> msg = consumer.receive() ;
GenericRecord record = msg.getValue();
record.getFields().forEach((field -> {
   if (field.getName().equals("theNeedFieldName")) {
       Object recordField = record.getField(field);
       //Do some things
   }
}));

스키마 저장소 커스터마이징 (Customize schema storage)

기본적으로 Pulsar는 다양한 데이터 타입의 스키마를 Pulsar와 함께 배포되는 Apache BookKeeper에 저장해요. 필요하다면 다른 저장소 시스템을 사용할 수도 있어요.

Pulsar 스키마에 기본값이 아닌(BookKeeper 외) 저장소 시스템을 사용하려면, 커스텀 스키마 저장소 배포 전에 다음 Java 인터페이스를 구현해야 해요.

SchemaStorage 인터페이스 구현 (Implement SchemaStorage interface)

SchemaStorage 인터페이스에는 다음 메서드가 있어요.

public interface SchemaStorage {
    // How schemas are updated
    CompletableFuture<SchemaVersion> put(String key, byte[] value, byte[] hash);
    // How schemas are fetched from storage
    CompletableFuture<StoredSchema> get(String key, SchemaVersion version);
    // How schemas are deleted
    CompletableFuture<SchemaVersion> delete(String key);
    // Utility method for converting a schema version byte array to a SchemaVersion object
    SchemaVersion versionFromBytes(byte[] version);
    // Startup behavior for the schema storage client
    void start() throws Exception;
    // Shutdown behavior for the schema storage client
    void close() throws Exception;
}

tip schema storage 구현의 완전한 예제는 BookKeeperSchemaStorage 클래스를 참고해요.

SchemaStorageFactory 인터페이스 구현 (Implement SchemaStorageFactory interface)

SchemaStorageFactory 인터페이스에는 다음 메서드가 있어요.

public interface SchemaStorageFactory {
    @NotNull
    SchemaStorage create(PulsarService pulsar) throws Exception;
}

tip schema storage factory 구현의 완전한 예제는 BookKeeperSchemaStorageFactory 클래스를 참고해요.

커스텀 스키마 저장소 배포 (Deploy custom schema storage)

커스텀 스키마 저장소 구현을 사용하려면 다음 단계를 수행해요.

  1. 구현을 JAR 파일로 패키징해요.
  2. Pulsar 바이너리 또는 소스 배포판의 lib 폴더에 JAR 파일을 추가해요.
  3. conf/broker.conf 파일의 schemaRegistryStorageClassName 구성을 커스텀 팩토리 클래스로 변경해요.
  4. Pulsar를 시작해요.

더 알아보기 (Learn more)