Scala for Backend Engineering & Functional Programming · 강의

파이프라인 실행하기

그래프를 구체화하고 실행해 보세요.

레슨 4/413개 단계

파이프라인 실행하기은(는) CoddyKit의 무료 Scala for Backend Engineering & Functional Programming 강의입니다. 이것은 4개 중 4번째 강의입니다. 아래에서 전체 강의를 무료로 읽을 수 있으며, 내장 코드 에디터와 24/7 AI 튜터와 함께 브라우저에서 직접 실습할 수 있습니다. 이 강의는 Scala for Backend Engineering & Functional Programming 학습 경로의 일부이며, 진행 상황이 웹과 CoddyKit 앱에 동기화됩니다. Scala for Backend Engineering & Functional Programming 강의에는 총 4개의 강의가 포함되어 있습니다.

블루프린트에서 실행으로

지금까지 파이프라인은 순수한 블루프린트였습니다. 머티리얼라이제이션은 이 블루프린트를 실제로 데이터를 이동시키는 실행 중인 액터로 바꾸는 과정입니다.

그래프를 명시적으로 실행하기 전까지는 아무 일도 일어나지 않습니다. 이 특성 덕분에 Akka Streams를 조합하고 재사용할 수 있습니다.

ActorSystem

머티리얼라이제이션에는 스트림 스테이지를 뒷받침하는 스레드와 디스패처를 제공하는 ActorSystem이 필요합니다. 최신 Akka에서는 시스템이 암시적 머티리얼라이저 역할도 합니다.

일반적으로 하나의 ActorSystem이 애플리케이션 전체와 여러 동시 스트림을 담당합니다.

import akka.actor.ActorSystem

implicit val system: ActorSystem =
  ActorSystem("data-pipeline")
import system.dispatcher // ExecutionContext

runWith

소스를 실행하는 가장 직접적인 방법은 runWith입니다. 이 방법은 싱크를 연결하고 한 단계에서 머티리얼라이즈하며, 해당 싱크의 머티리얼라이즈된 값을 반환합니다.

여기서 결과는 스트림이 끝날 때 합계와 함께 완료되는 Future[Int]입니다.

import akka.stream.scaladsl.{Source, Sink}
import scala.concurrent.Future

val total: Future[Int] =
  Source(1 to 100).runWith(Sink.fold(0)(_ + _))

RunnableGraph에서 run 실행하기

이미 to 또는 toMat으로 닫힌 RunnableGraph를 구성했다면 run()을 호출해 머티리얼라이즈하세요. 반환값은 그래프가 유지한 머티리얼라이즈된 값입니다.

이렇게 하면 파이프라인 구성과 실행을 명확하게 분리할 수 있습니다.

import akka.stream.scaladsl.Keep
import scala.concurrent.Future

val graph =
  Source(1 to 100)
    .toMat(Sink.fold(0)(_ + _))(Keep.right)

val result: Future[Int] = graph.run()

편의 실행 연산자

소스는 여러 단축 메서드를 제공합니다. runForeach, runFold, runReduce는 각각 해당하는 싱크를 연결하고 즉시 실행합니다.

소스의 일반적인 종단 연산을 간결하게 작성할 때 유용합니다.

import scala.concurrent.Future

val printed: Future[akka.Done] =
  Source(1 to 10).runForeach(println)

val sum: Future[Int] =
  Source(1 to 10).runFold(0)(_ + _)

결과 Future 다루기

종단 싱크는 스트림이 종료되거나 실패할 때 완료되는 Future를 반환합니다. onComplete로 콜백을 등록해 성공이나 오류에 대응하세요.

이 콜백의 암시적 ExecutionContext로 ActorSystem의 디스패처를 사용하세요.

import scala.util.{Success, Failure}

total.onComplete {
  case Success(value) => println(s"Sum = $value")
  case Failure(ex)    => println(s"Failed: ${ex.getMessage}")
}

현실적인 파이프라인

일반적인 데이터 처리 파이프라인은 소스에서 데이터를 읽고, 플로로 변환하며, mapAsync로 비동기 입출력을 수행하고, grouped로 일괄 처리한 다음 싱크에 기록합니다.

각 단계는 작게 구성되며, 전체 파이프라인은 단 한 번의 run으로 구체화됩니다.

val done =
  lineSource
    .map(parse)
    .mapAsync(4)(validate)
    .grouped(500)
    .runWith(bulkWriteSink)

실패한 스트림 다시 시작하기

복원력을 높이려면 소스나 플로를 RestartSource.withBackoff로 감싸십시오. 그러면 일시적인 오류(연결 끊김)가 발생했을 때 지수형 재시도 간격을 적용해 자동으로 다시 시작합니다.

이렇게 하면 수동으로 감독하지 않아도 장시간 실행되는 데이터 수집 파이프라인을 계속 유지할 수 있습니다.

import akka.stream.scaladsl.RestartSource
import akka.stream.RestartSettings
import scala.concurrent.duration._

val resilient = RestartSource.withBackoff(
  RestartSettings(1.second, 30.seconds, 0.2))(() => flakySource)

KillSwitch를 사용한 정상적인 종료

KillSwitch를 사용하면 외부 코드에서 실행 중인 스트림을 깔끔하게 중지할 수 있습니다. viaMat을 통해 KillSwitches.single을 삽입하고, 구체화된 값을 보관해 두었다가 나중에 shutdown()을 호출하십시오.

애플리케이션이 종료될 때 중지해야 하는 장수 스트림에는 이 기능이 필수적입니다.

import akka.stream.{KillSwitches, KillSwitch}
import akka.stream.scaladsl.Keep

val (switch, done) =
  source
    .viaMat(KillSwitches.single)(Keep.right)
    .toMat(Sink.ignore)(Keep.both)
    .run()
// later: switch.shutdown()

리소스 해제하기

애플리케이션이 종료될 때 스레드를 해제하려면 ActorSystem을 종료하십시오. 스트림의 완료를 나타내는 비동기 결과 뒤에 terminate 호출을 연결하면 종료가 질서 있게 진행됩니다.

ActorSystem을 누수시키면 JVM이 계속 실행되고 리소스가 열린 상태로 유지됩니다.

done.onComplete { _ =>
  system.terminate()
}

스트림 실행기 재사용하기

동일한 설계도를 여러 번 구체화하면 서로 독립적으로 실행되는 스트림이 생성되며, 이 스트림들은 ActorSystem의 리소스를 공유합니다. 설계도 자체는 변경되지 않고 부작용도 없는 상태로 유지됩니다.

따라서 파이프라인을 한 번 정의한 뒤 들어오는 각 작업에 필요할 때마다 실행해도 안전합니다.

val blueprint =
  Source(1 to 5).toMat(Sink.seq)(Keep.right)

val run1 = blueprint.run()
val run2 = blueprint.run() // independent execution

빠른 확인

스트림이 실제로 요소를 처리하려면 무엇이 필요한지 생각해 보십시오.

복습

파이프라인을 실행한다는 것은 ActorSystem을 사용해 설계도를 구체화하는 것을 의미하며, run, runWith 또는 편의 연산자를 사용하면 각각 비동기 결과를 반환합니다.

현실적인 다단계 파이프라인, 재시도 간격을 적용한 자동 재시작, KillSwitch를 통한 정상적인 종료, system.terminate()를 사용한 리소스 정리, 그리고 변경할 수 없는 설계도를 서로 독립적인 실행에서 안전하게 재사용하는 방법을 살펴보았습니다.

무료로 시작

AI 튜터와 함께 Scala을(를) 배우세요 — 무료

브라우저에서 실제 코드를 작성하고 실행하며, 24/7 AI 튜터로부터 즉각적인 도움을 받고, 웹이나 앱에서 중단한 부분부터 계속 학습하세요.

코스
39
레슨
143

자주 묻는 질문

“파이프라인 실행하기” 강의는 무료인가요?

네 — “파이프라인 실행하기” 전체 내용을 이 웹사이트에서 무료로 읽을 수 있습니다. 인터랙티브하게 실습하려면(내장 코드 에디터와 24/7 AI 튜터), CoddyKit PRO로 업그레이드하면 Scala for Backend Engineering & Functional Programming 강의 전체를 잠금 해제할 수 있습니다. Scala for Backend Engineering & Functional Programming 강의에는 총 4개의 강의가 포함되어 있습니다.

“파이프라인 실행하기”에서 뭘 배우나요?

그래프를 구체화하고 실행해 보세요. 브라우저에서 직접 실행하는 실습 코드로 Scala for Backend Engineering & Functional Programming을(를) 배우며, 24/7 AI 튜터가 강의를 진행하면서 질문에 답변해줍니다.

Scala for Backend Engineering & Functional Programming을(를) 시작하는 데 경험이 필요한가요?

사전 경험은 필요하지 않습니다. CoddyKit의 Scala for Backend Engineering & Functional Programming은(는) 초급자부터 고급 학습자까지를 위해 구성되어 있으므로, 여기서 시작하거나 처음부터 시작할 수 있으며 자신의 속도대로 진행할 수 있습니다. 이것은 4개 중 4번째 강의입니다.

“파이프라인 실행하기” 강의는 얼마나 걸리나요?

대부분의 CoddyKit 강의는 약 5~10분이 소요됩니다. 각 강의는 간결하고 인터랙티브하여 꾸준한 진행이 가능하며, 웹과 앱에서 중단한 부분부터 바로 시작할 수 있습니다.

이 Scala for Backend Engineering & Functional Programming 강의에서 코드를 작성하고 실행할 수 있나요?

네. 모든 Scala for Backend Engineering & Functional Programming 강의에는 내장 코드 에디터가 포함되어 있으므로, 브라우저에서 바로 실제 코드를 작성하고 실행한 후 즉시 AI 피드백을 받을 수 있습니다 — 로컬 설정이 필요 없습니다.

이 강의의 모든 강의

  1. 소스, 플로, 싱크
  2. 스트림 변환하기
  3. 백프레셔
  4. 파이프라인 실행하기
← Scala for Backend Engineering & Functional Programming(으)로 돌아가기