Spark Streaming 사용자 정의 리시버

Spark Streaming 사용자 정의 리시버 (Custom Receivers)

Spark Streaming은 기본 지원되는 소스(Kafka, Kinesis, 파일, 소켓 등) 이외의 임의의 데이터 소스에서도 스트리밍 데이터를 받을 수 있어요. 다만 그러려면 해당 데이터 소스에서 데이터를 받기 위한 *리시버(receiver)*를 직접 구현해야 합니다. 이 가이드에서는 사용자 정의 리시버를 구현하고 Spark Streaming 애플리케이션에서 사용하는 과정을 함께 살펴볼게요. 참고로 사용자 정의 리시버는 Scala나 Java로 구현할 수 있어요.

출처: 문서

본문

Spark Streaming은 기본 지원되는 소스(즉, Kafka, Kinesis, 파일, 소켓 등을 제외한) 외의 임의의 데이터 소스에서도 스트리밍 데이터를 받을 수 있어요. 이를 위해서는 개발자가 해당 데이터 소스에서 데이터를 받기 위해 맞춤화된 *리시버(receiver)*를 구현해야 합니다. 이 가이드는 사용자 정의 리시버를 구현하고 Spark Streaming 애플리케이션에서 사용하는 과정을 안내해요. 참고로 사용자 정의 리시버는 Scala나 Java로 구현할 수 있습니다.

사용자 정의 리시버 구현하기

이것은 Receiver(Scala 문서, Java 문서)를 구현하는 것으로 시작합니다. 사용자 정의 리시버는 이 추상 클래스를 상속해 두 메서드를 구현해야 해요:

  • onStart(): 데이터 수신을 시작하기 위해 해야 할 일.
  • onStop(): 데이터 수신을 중지하기 위해 해야 할 일.

onStart()onStop() 모두 무한정 블록하면 안 됩니다. 일반적으로 onStart()는 데이터 수신을 담당하는 스레드를 시작하고, onStop()은 이 데이터 수신 스레드들이 중지되도록 보장합니다. 수신 스레드는 Receiver의 메서드인 isStopped()를 사용해 데이터 수신을 중지해야 하는지 확인할 수도 있어요.

데이터가 수신되면, Receiver 클래스가 제공하는 메서드인 store(data)를 호출해 그 데이터를 Spark 안에 저장할 수 있습니다. store()에는 레코드를 한 번에 하나씩 저장하거나 객체/직렬화된 바이트 컬렉션으로 저장할 수 있는 여러 변형이 있어요. 리시버를 구현할 때 사용하는 store()의 변형에 따라 신뢰성과 장애 허용(fault-tolerance) 의미론이 달라진다는 점에 주의하세요. 이에 대해서는 나중에 더 자세히 다룹니다.

수신 스레드의 예외는 리시버의 조용한 실패(silent failure)를 피하기 위해 적절히 포착하고 처리해야 합니다. restart(<exception>)onStop()을 비동기적으로 호출한 뒤 잠시 후 onStart()를 호출하여 리시버를 재시작합니다. stop(<exception>)onStop()을 호출하고 리시버를 종료해요. 또한 reportError(<error>)는 리시버를 중지/재시작하지 않고 드라이버에 오류 메시지를 보고합니다(로그와 UI에 표시됩니다).

다음은 소켓을 통해 텍스트 스트림을 수신하는 사용자 정의 리시버예요. 텍스트 스트림에서 '\n'으로 구분된 줄을 레코드로 취급해 Spark에 저장합니다. 수신 스레드에서 연결이나 수신에 오류가 있으면 리시버를 재시작해 다시 연결을 시도해요.

Scala:

class CustomReceiver(host: String, port: Int)
  extends Receiver[String](StorageLevel.MEMORY_AND_DISK_2) with Logging {

  def onStart() {
    // Start the thread that receives data over a connection
    new Thread("Socket Receiver") {
      override def run() { receive() }
    }.start()
  }

  def onStop() {
    // There is nothing much to do as the thread calling receive()
    // is designed to stop by itself if isStopped() returns false
  }

  /** Create a socket connection and receive data until receiver is stopped */
  private def receive() {
    var socket: Socket = null
    var userInput: String = null
    try {
      // Connect to host:port
      socket = new Socket(host, port)

      // Until stopped or connection broken continue reading
      val reader = new BufferedReader(
        new InputStreamReader(socket.getInputStream(), StandardCharsets.UTF_8))
      userInput = reader.readLine()
      while(!isStopped && userInput != null) {
        store(userInput)
        userInput = reader.readLine()
      }
      reader.close()
      socket.close()

      // Restart in an attempt to connect again when server is active again
      restart("Trying to connect again")
    } catch {
      case e: java.net.ConnectException =>
        // restart if could not connect to server
        restart("Error connecting to " + host + ":" + port, e)
      case t: Throwable =>
        // restart if there is any other error
        restart("Error receiving data", t)
    }
  }
}

Java:

public class JavaCustomReceiver extends Receiver<String> {

  String host = null;
  int port = -1;

  public JavaCustomReceiver(String host_ , int port_) {
    super(StorageLevel.MEMORY_AND_DISK_2());
    host = host_;
    port = port_;
  }

  @Override
  public void onStart() {
    // Start the thread that receives data over a connection
    new Thread(this::receive).start();
  }

  @Override
  public void onStop() {
    // There is nothing much to do as the thread calling receive()
    // is designed to stop by itself if isStopped() returns false
  }

  /** Create a socket connection and receive data until receiver is stopped */
  private void receive() {
    Socket socket = null;
    String userInput = null;

    try {
      // connect to the server
      socket = new Socket(host, port);

      BufferedReader reader = new BufferedReader(
        new InputStreamReader(socket.getInputStream(), StandardCharsets.UTF_8));

      // Until stopped or connection broken continue reading
      while (!isStopped() && (userInput = reader.readLine()) != null) {
        System.out.println("Received data '" + userInput + "'");
        store(userInput);
      }
      reader.close();
      socket.close();

      // Restart in an attempt to connect again when server is active again
      restart("Trying to connect again");
    } catch(ConnectException ce) {
      // restart if could not connect to server
      restart("Could not connect", ce);
    } catch(Throwable t) {
      // restart if there is any other error
      restart("Error receiving data", t);
    }
  }
}

Spark Streaming 애플리케이션에서 사용자 정의 리시버 사용하기

사용자 정의 리시버는 streamingContext.receiverStream(<사용자 정의 리시버 인스턴스>)를 사용해 Spark Streaming 애플리케이션에서 사용할 수 있어요. 이렇게 하면 사용자 정의 리시버 인스턴스가 수신한 데이터로 입력 DStream이 생성됩니다:

Scala:

// Assuming ssc is the StreamingContext
val customReceiverStream = ssc.receiverStream(new CustomReceiver(host, port))
val words = customReceiverStream.flatMap(_.split(" "))
...
전체 소스 코드는 예제 [CustomReceiver.scala](https://github.com/apache/spark/blob/v4.2.0/examples/src/main/scala/org/apache/spark/examples/streaming/CustomReceiver.scala)에 있습니다.

Java:

// Assuming ssc is the JavaStreamingContext
JavaDStream<String> customReceiverStream = ssc.receiverStream(new JavaCustomReceiver(host, port));
JavaDStream<String> words = customReceiverStream.flatMap(s -> ...);
...
전체 소스 코드는 예제 [JavaCustomReceiver.java](https://github.com/apache/spark/blob/v4.2.0/examples/src/main/java/org/apache/spark/examples/streaming/JavaCustomReceiver.java)에 있습니다.

리시버 신뢰성 (Receiver Reliability)

Spark Streaming 프로그래밍 가이드에서 간략히 논의했듯이, 신뢰성과 장애 허용 의미론에 따라 두 종류의 리시버가 있어요.

  • 신뢰할 수 있는 리시버 (Reliable Receiver) - 전송된 데이터를 확인(acknowledge)해 주는 신뢰할 수 있는 소스의 경우, 신뢰할 수 있는 리시버는 데이터가 Spark에 안정적으로(즉, 성공적으로 복제되어) 수신되고 저장되었음을 소스에 올바르게 확인해 줍니다. 보통 이 리시버를 구현할 때는 소스의 확인 의미론을 신중하게 고려해야 해요.

  • 신뢰할 수 없는 리시버 (Unreliable Receiver) - 신뢰할 수 없는 리시버는 소스에 확인을 보내지 않습니다. 이는 확인을 지원하지 않는 소스에 사용할 수 있고, 확인의 복잡성까지 감수하고 싶지 않거나 필요하지 않은 경우에는 신뢰할 수 있는 소스에도 사용할 수 있어요.

신뢰할 수 있는 리시버를 구현하려면 store(multiple-records)를 사용해 데이터를 저장해야 합니다. 이 store 변형은 블로킹 호출로, 주어진 모든 레코드가 Spark에 저장된 후에만 반환해요. 리시버의 구성된 저장 수준이 복제(기본적으로 활성화)를 사용한다면, 이 호출은 복제가 완료된 후에 반환됩니다. 따라서 데이터가 안정적으로 저장되었음을 보장하며, 리시버는 이제 소스에 적절히 확인을 보낼 수 있어요. 이는 리시버가 데이터를 복제하는 중간에 실패해도 데이터가 손실되지 않도록 보장합니다. 버퍼링된 데이터는 확인되지 않으므로 나중에 소스가 다시 보내게 되기 때문이에요.

신뢰할 수 없는 리시버는 이런 로직을 구현할 필요가 없어요. 단순히 소스에서 레코드를 받아 store(single-record)를 사용해 한 번에 하나씩 삽입하면 됩니다. store(multiple-records)의 신뢰성 보장을 얻을 수는 없지만 다음과 같은 장점이 있어요:

  • 시스템이 해당 데이터를 적절한 크기의 블록으로 청킹해 줘요 (블록 간격에 대해서는 Spark Streaming 프로그래밍 가이드 참고).
  • 비율 제한(rate limits)이 지정되어 있으면 시스템이 수신 비율을 제어해 줘요.
  • 이 두 가지 덕분에 신뢰할 수 없는 리시버는 신뢰할 수 있는 리시버보다 구현이 더 간단합니다.

다음 표는 두 종류의 리시버 특성을 요약한 것이에요:

리시버 유형 특성
신뢰할 수 없는 리시버 구현이 간단함. 시스템이 블록 생성과 비율 제어를 처리함. 장애 허용 보장이 없어 리시버 실패 시 데이터를 잃을 수 있음.
신뢰할 수 있는 리시버 강력한 장애 허용 보장으로 데이터 손실을 0으로 보장할 수 있음. 블록 생성과 비율 제어는 리시버 구현이 처리해야 함. 구현 복잡도는 소스의 확인 메커니즘에 따라 달라짐.

더 알아보기 (Learn more)