Pulsar 커넥터 개발하기

Pulsar 커넥터 개발하기 (How to develop Pulsar connectors)

이 가이드는 Pulsar와 다른 시스템 사이에서 데이터를 옮기는 Pulsar 커넥터를 개발하는 방법을 설명해요. Pulsar 커넥터는 특별한 Pulsar Functions라서, Pulsar 커넥터를 만드는 것은 Pulsar function을 만드는 것과 비슷해요.

출처: 문서

본문

Pulsar 커넥터는 두 가지 유형으로 나뉘어요.

타입 설명 예시
Source 다른 시스템에서 Pulsar로 데이터를 가져와요. RabbitMQ source 커넥터가 RabbitMQ 큐의 메시지를 Pulsar 토픽으로 가져와요.
Sink Pulsar에서 다른 시스템으로 데이터를 내보내요. Kinesis sink 커넥터가 Pulsar 토픽의 메시지를 Kinesis 스트림으로 내보내요.

개발 (Develop)

Pulsar source 커넥터와 sink 커넥터를 모두 개발할 수 있어요.

Source

source 커넥터를 개발하려면 기본적으로 Source 인터페이스를 구현하는 것이며, open 메서드와 read 메서드를 구현해야 해요.

  • open 메서드를 구현해요.
/**
 * Open connector with configuration
 * @param config initialization config
 * @param sourceContext
 * @throws Exception IO type exceptions when opening a connector
 */
void open(final Map<String, Object> config, SourceContext sourceContext) throws Exception;

이 메서드는 source 커넥터가 초기화될 때 호출돼요. 이 메서드에서 config 파라미터로 전달된 모든 커넥터 전용 설정을 가져오고 필요한 모든 리소스를 초기화할 수 있어요.

예를 들어 Kafka 커넥터는 이 open 메서드에서 Kafka 클라이언트를 만들 수 있어요.

또한 Pulsar 런타임은 커넥터가 메트릭 수집 같은 작업을 위해 런타임 리소스에 접근할 수 있도록 SourceContext를 제공해요. 구현체는 SourceContext를 저장해 두고 나중에 사용할 수 있어요.

  • read 메서드를 구현해요.
    /**
     * Reads the next message from source.
     * If source does not have any new messages, this call should block.
     * @return next message from source.  The return result should never be null
     * @throws Exception
     */
    Record<T> read() throws Exception;

반환할 것이 없다면 구현은 null을 반환하는 대신 블로킹(blocking)되어야 해요.

반환된 Record는 Pulsar IO 런타임에 필요한 다음 정보를 캡슐화해야 해요.

  • Record는 다음 변수를 제공해야 해요.
변수 필수 설명
TopicName 아니오 레코드가 비롯된 Pulsar 토픽 이름이에요.
Key 아니오 메시지는 선택적으로 키로 태그될 수 있어요. 더 자세한 내용은 Routing modes를 참고하세요.
Value 레코드의 실제 데이터예요.
EventTime 아니오 소스에서 온 레코드의 이벤트 시간이에요.
PartitionId 아니오 레코드가 파티셔닝된 소스에서 비롯됐다면 해당 PartitionId를 반환해요. PartitionId는 Pulsar IO 런타임이 메시지를 중복 제거하고 exactly-once 처리 보장을 달성하기 위해 고유 식별자의 일부로 사용해요.
RecordSequence 아니오 레코드가 순차(sequential) 소스에서 비롯됐다면 해당 RecordSequence를 반환해요. RecordSequence는 Pulsar IO 런타임이 메시지를 중복 제거하고 exactly-once 처리 보장을 달성하기 위해 고유 식별자의 일부로 사용해요.
Properties 아니오 레코드가 사용자 정의 프로퍼티를 담고 있다면 그 프로퍼티들을 반환해요.
DestinationTopic 아니오 메시지를 써야 하는 토픽이에요.
Message 아니오 사용자가 보낸 데이터를 담는 클래스예요. 더 자세한 내용은 Message.java를 참고하세요.
  • Record는 다음 메서드를 제공해야 해요.
메서드 설명
ack 레코드가 완전히 처리됐음을 acknowledge해요.
fail 레코드가 처리에 실패했음을 나타내요.
스키마 정보 처리하기 (Handle schema information)

Pulsar IO는 스키마를 자동으로 처리하고 Java 제네릭스 기반의 강타입(strongly typed) API를 제공해요. 생산하는 스키마 타입을 알고 있다면, source 선언에서 해당 타입에 대응하는 Java 클래스를 선언할 수 있어요.

public class MySource implements Source<String> {
    public Record<String> read() {}
}

어떤 스키마와도 동작하는 source를 구현하려면 byte[](또는 ByteBuffer)를 사용하고 Schema.AUTO_PRODUCE_BYTES()를 쓰면 돼요.

public class MySource implements Source<byte[]> {
    public Record<byte[]> read() {
        Schema wantedSchema = ....
        Record<byte[]> myRecord = new MyRecordImplementation();
        ....
    }
    class MyRecordImplementation implements Record<byte[]> {
         public byte[] getValue() {
            return ....encoded byte[]...that represents the value
         }
         public Schema<byte[]> getSchema() {
             return Schema.AUTO_PRODUCE_BYTES(wantedSchema);
         }
    }
}

KeyValue 타입을 제대로 처리하려면 레코드 구현에 대한 지침을 따르세요.

  • Record 인터페이스를 구현하고 getKeySchema, getValueSchema, getKeyValueEncodingType을 구현해야 해요.
  • Record.getValue()KeyValue 객체를 반환해야 해요.
  • Record.getSchema()에서 null을 반환할 수 있어요.

Pulsar IO 런타임이 KVRecord를 만나면 다음 변경을 자동으로 수행해요.

  • 적절한 KeyValueSchema를 설정해요.
  • KeyValueEncoding(SEPARATED 또는 INLINE)에 따라 메시지 키와 메시지 값을 인코딩해요.

: source 커넥터를 만드는 방법에 대한 더 자세한 내용은 KafkaSource를 참고하세요.

Sink

sink 커넥터를 개발하는 것은 source 커넥터 개발과 비슷해요. Sink 인터페이스를 구현해야 하며, open 메서드와 write 메서드 모두를 구현해야 해요.

  • open 메서드를 구현해요.
    /**
     * Open connector with configuration
     *
     * @param config initialization config
     * @param sinkContext
     * @throws Exception IO type exceptions when opening a connector
     */
    void open(final Map<String, Object> config, SinkContext sinkContext) throws Exception;
  • write 메서드를 구현해요.
    /**
     * Write a message to Sink
     * @param record record to write to sink
     * @throws Exception
     */
    void write(Record<T> record) throws Exception;

구현 중에 ValueKey를 실제 sink에 어떻게 쓸지 결정하고, PartitionIdRecordSequence 같은 제공된 모든 정보를 활용해 다양한 처리 보장을 달성할 수 있어요.

또한 레코드를 ack(메시지가 성공적으로 전송된 경우)하거나 fail(메시지 전송에 실패한 경우)해야 해요.

스키마 정보 처리하기 (Handle schema information)

Pulsar IO는 스키마를 자동으로 처리하고 Java 제네릭스 기반의 강타입 API를 제공해요. 소비하는 스키마 타입을 알고 있다면, Sink 선언에서 해당 타입에 대응하는 Java 클래스를 선언할 수 있어요.

public class MySink implements Sink<String> {
    public void write(Record<String> record) {}
}

어떤 스키마와도 동작하는 sink를 구현하려면 특별한 GenericObject 인터페이스를 사용할 수 있어요.

public class MySink implements Sink<GenericObject> {
    public void write(Record<GenericObject> record) {
        Schema schema = record.getSchema();
        GenericObject genericObject = record.getValue();
        if (genericObject != null) {
            SchemaType type = genericObject.getSchemaType();
            Object nativeObject = genericObject.getNativeObject();
            ...
        }
        ....
    }
}

AVRO, JSON, Protobuf 레코드(schemaType=AVRO, JSON, PROTOBUF_NATIVE)의 경우 genericObject 변수를 GenericRecord로 캐스팅하고 getFields()getField() API를 사용할 수 있어요. genericObject.getNativeObject()로 네이티브 AVRO 레코드에 접근할 수 있어요.

KeyValue 타입의 경우 이 코드로 키 스키마와 값 스키마 양쪽에 접근할 수 있어요.

public class MySink implements Sink<GenericObject> {
    public void write(Record<GenericObject> record) {
        Schema schema = record.getSchema();
        GenericObject genericObject = record.getValue();
        SchemaType type = genericObject.getSchemaType();
        Object nativeObject = genericObject.getNativeObject();
        if (type == SchemaType.KEY_VALUE) {
            KeyValue keyValue = (KeyValue) nativeObject;
            Object key = keyValue.getKey();
            Object value = keyValue.getValue();
            KeyValueSchema keyValueSchema = (KeyValueSchema) schema;
            Schema keySchema = keyValueSchema.getKeySchema();
            Schema valueSchema = keyValueSchema.getValueSchema();
        }
        ...
    }
}

테스트 (Test)

Pulsar IO 커넥터는 두 시스템(mock하기 어려울 수 있는 Pulsar와 커넥터가 연결되는 시스템)과 상호작용하기 때문에 커넥터 테스트는 까다로울 수 있어요.

외부 서비스를 mock하면서 커넥터 기능을 아래처럼 테스트하는 특별한 테스트를 작성하는 것이 권장돼요.

단위 테스트 (Unit test)

커넥터에 대한 단위 테스트를 만들 수 있어요.

통합 테스트 (Integration test)

충분한 단위 테스트를 작성한 뒤, 엔드투엔드 기능을 검증하는 별도의 통합 테스트를 추가할 수 있어요. Pulsar는 모든 통합 테스트에 testcontainers를 사용해요.

: Pulsar 커넥터 통합 테스트를 만드는 방법에 대한 더 자세한 내용은 IntegrationTests를 참고하세요.

패키징 (Package)

커넥터를 개발하고 테스트했다면, Pulsar Functions 클러스터에 제출할 수 있도록 패키징해야 해요. 커넥터를 패키징하는 방법은 두 가지예요.

  • NAR
  • Uber JAR

참고: 커넥터를 패키징해서 다른 사람이 쓸 수 있도록 배포할 계획이라면, 자신의 코드를 적절히 라이선스하고 저작권을 표기할 책임이 있어요. 코드가 사용하는 모든 라이브러리와 배포물에 라이선스와 저작권을 추가해야 해요.

NAR 방법을 쓰면 NAR 플러그인이 생성된 NAR 패키지에 DEPENDENCIES 파일을 자동으로 만들어, 커넥터의 모든 라이브러리의 적절한 라이선스와 저작권을 포함해요.

런타임 Java 버전에 대해서는 자신의 대상 Pulsar 버전에 따른 Pulsar Runtime Java Version Recommendation을 참고하세요.

NAR

NAR는 NiFi Archive의 약자로, Apache NiFi가 사용하는 커스텀 패키징 메커니즘이며 약간의 Java ClassLoader 격리를 제공해요.

: NAR가 어떻게 동작하는지에 대한 더 자세한 내용은 여기를 참고하세요.

Pulsar는 모든 내장 커넥터를 패키징하는 데 같은 메커니즘을 사용해요. Pulsar 커넥터를 패키징하는 가장 쉬운 방법은 nifi-nar-maven-plugin으로 NAR 패키지를 만드는 거예요.

커넥터의 maven 프로젝트에 아래처럼 nifi-nar-maven-plugin을 포함하세요.

<plugins>
  <plugin>
    <groupId>org.apache.nifi</groupId>
    <artifactId>nifi-nar-maven-plugin</artifactId>
    <version>1.5.0</version>
  </plugin>
</plugins>

또한 다음 내용으로 resources/META-INF/services/pulsar-io.yaml 파일을 만들어야 해요.

name: connector name
description: connector description
sourceClass: fully qualified class name (only if source connector)
sinkClass: fully qualified class name (only if sink connector)

Gradle 사용자라면 Gradle Plugin Portal에서 사용 가능한 Gradle Nar 플러그인이 있어요.

: Pulsar 커넥터에 NAR을 사용하는 방법에 대한 더 자세한 내용은 DataGen을 참고하세요.

Uber JAR

대안적인 방법은 커넥터의 모든 JAR 파일과 다른 리소스 파일을 포함하는 uber JAR을 만드는 거예요. 내부 디렉터리 구조는 필요 없어요.

아래처럼 maven-shade-plugin으로 uber JAR을 만들 수 있어요.

<plugin>
  <groupId>org.apache.maven.plugins</groupId>
  <artifactId>maven-shade-plugin</artifactId>
  <version>3.1.1</version>
  <executions>
    <execution>
      <phase>package</phase>
      <goals>
        <goal>shade</goal>
      </goals>
      <configuration>
        <filters>
          <filter>
            <artifact>*:*</artifact>
          </filter>
        </filters>
      </configuration>
    </execution>
  </executions>
</plugin>

모니터링 (Monitor)

Pulsar 커넥터는 Pulsar 안팎으로 데이터를 쉽게 옮길 수 있게 해줘요. 실행 중인 커넥터가 항상 건강한지 확인하는 것이 중요해요. 배포된 Pulsar 커넥터는 다음 방법으로 모니터링할 수 있어요.

  • Pulsar가 제공하는 메트릭을 확인해요.

Pulsar 커넥터는 Java 커넥터의 건강을 모니터링하는 데 사용할 수 있는 메트릭을 노출해요. 모니터링 가이드에 따라 메트릭을 확인할 수 있어요.

  • 커스터마이즈한 메트릭을 설정하고 확인해요.

Pulsar가 제공하는 메트릭 외에도, Pulsar는 Java 커넥터에 대해 메트릭을 커스터마이즈할 수 있게 해줘요. Functions 워커는 사용자 정의 메트릭을 Prometheus에 자동으로 수집하고, Grafana에서 확인할 수 있어요.

Java 커넥터에 메트릭을 커스터마이즈하는 예시는 다음과 같아요.

  • Java
public class TestMetricSink implements Sink<String> {
        @Override
        public void open(Map<String, Object> config, SinkContext sinkContext) throws Exception {
            sinkContext.recordMetric("foo", 1);
        }
        @Override
        public void write(Record<String> record) throws Exception {
        }
        @Override
        public void close() throws Exception {
        }
    }

더 알아보기 (Learn more)