Apache Spark용 Pulsar 어댑터

Apache Spark용 Pulsar 어댑터

Pulsar는 Apache Spark Streaming과 연동할 수 있는 수신기(receiver)를 제공해요. Spark Streaming용 Pulsar 수신기는 Pulsar에서 원시 데이터를 받아 Apache Spark Streaming에 전달하는 커스텀 리시버예요. 수신한 데이터는 Resilient Distributed Dataset(RDD) 형식으로 받아 다양한 방식으로 처리할 수 있어요.

출처: 문서

본문

Spark Streaming 수신기 (Spark Streaming receiver)

Spark Streaming용 Pulsar 수신기는 Pulsar에서 원시 데이터를 받아 Apache Spark Streaming에 전달하는 커스텀 리시버예요. 애플리케이션은 이 수신기를 통해 Resilient Distributed Dataset(RDD) 형식으로 데이터를 받아 다양한 방식으로 처리할 수 있어요.

사전 준비 (Prerequisites)

수신기를 사용하려면 Java 구성에 pulsar-spark 라이브러리 의존성을 포함해야 해요.

Maven

Maven을 사용한다면 pom.xml에 다음을 추가해요.

<!-- in your <properties> block -->
<pulsar.version>5.0.0-M2</pulsar.version>

<!-- in your <dependencies> block -->
<dependency>
  <groupId>org.apache.pulsar</groupId>
  <artifactId>pulsar-spark</artifactId>
  <version>${pulsar.version}</version>
</dependency>

Gradle

Gradle을 사용한다면 build.gradle 파일에 다음을 추가해요.

def pulsarVersion = "5.0.0-M2"

dependencies {
    compile group: 'org.apache.pulsar', name: 'pulsar-spark', version: pulsarVersion
}

사용법 (Usage)

SparkStreamingPulsarReceiver 인스턴스를 JavaStreamingContextreceiverStream 메서드에 전달해요.

    String serviceUrl = "pulsar://localhost:6650/";
    String topic = "persistent://public/default/test_src";
    String subs = "test_sub";

    SparkConf sparkConf = new SparkConf().setMaster("local[*]").setAppName("Pulsar Spark Example");

    JavaStreamingContext jsc = new JavaStreamingContext(sparkConf, Durations.seconds(60));

    ConsumerConfigurationData<byte[]> pulsarConf = new ConsumerConfigurationData();

    Set<String> set = new HashSet();
    set.add(topic);
    pulsarConf.setTopicNames(set);
    pulsarConf.setSubscriptionName(subs);

    SparkStreamingPulsarReceiver pulsarReceiver = new SparkStreamingPulsarReceiver(
        serviceUrl,
        pulsarConf,
        new AuthenticationDisabled());

    JavaReceiverInputDStream<byte[]> lineDStream = jsc.receiverStream(pulsarReceiver);

전체 예제는 여기를 클릭해서 볼 수 있어요. 이 예제에서는 수신한 메시지 중 "Pulsar" 문자열을 포함하는 메시지 수를 세요.

필요하다면 다른 Pulsar 인증 클래스도 사용할 수 있어요. 예를 들어 인증 시 토큰을 사용하려면 SparkStreamingPulsarReceiver 생성자의 다음 파라미터를 설정하면 돼요.

SparkStreamingPulsarReceiver pulsarReceiver = new SparkStreamingPulsarReceiver(
        serviceUrl,
        pulsarConf,
        new AuthenticationToken("token:<secret-JWT-token>"));

더 알아보기 (Learn more)

  • Spark Streaming에 대한 기본 개념은 Apache Spark 공식 문서를 참고해요.
  • Pulsar 인증 방식에 대해 더 알고 싶다면 인증 관련 문서를 살펴보세요.
  • JWT 토큰 인증 설정 방법은 보안 문서에서 확인할 수 있어요.