소스, 플로, 싱크
스트리밍을 구성하는 기본 요소를 배워 보세요.
소스, 플로, 싱크은(는) CoddyKit의 무료 Scala for Backend Engineering & Functional Programming 강의입니다. 이것은 4개 중 1번째 강의입니다. 아래에서 전체 강의를 무료로 읽을 수 있으며, 내장 코드 에디터와 24/7 AI 튜터와 함께 브라우저에서 직접 실습할 수 있습니다. 이 강의는 Scala for Backend Engineering & Functional Programming 학습 경로의 일부이며, 진행 상황이 웹과 CoddyKit 앱에 동기화됩니다. Scala for Backend Engineering & Functional Programming 강의에는 총 4개의 강의가 포함되어 있습니다.
세 가지 구성 요소
Akka Streams는 데이터 파이프라인을 처리 단계의 그래프로 모델링합니다. 세 가지 핵심 선형 단계는 요소를 생성하는 소스, 요소를 변환하는 플로, 요소를 소비하는 싱크입니다.
소스에는 출력이 하나 있고 싱크에는 입력이 하나 있으며, 플로에는 정확히 하나의 입력과 하나의 출력이 있습니다. 이들을 연결하면 언제가 아니라 무엇을 해야 하는지를 설명합니다.
Source 정의하기
Source[Out, Mat]은 Out 타입의 요소를 내보내고 Mat 타입의 머티리얼라이즈된 값을 제공합니다. 가장 단순한 소스는 메모리 내 컬렉션이나 범위에서 만들어집니다.
스트림을 실행하기 전까지 Source는 자유롭게 재사용할 수 있는 불변 청사진일 뿐입니다.
import akka.stream.scaladsl.Source
val numbers: Source[Int, akka.NotUsed] =
Source(1 to 100)
val single: Source[String, akka.NotUsed] =
Source.single("hello")Sink 정의하기
Sink[In, Mat]은 In 타입의 요소를 소비합니다. 머티리얼라이즈된 값에는 소비 결과가 담기는 경우가 많으며, 예를 들어 스트림이 끝나면 완료되는 Future가 있습니다.
Sink.foreach는 요소마다 부수 효과를 실행하고, Sink.fold는 하나의 결과를 누적합니다.
import akka.stream.scaladsl.Sink
import scala.concurrent.Future
val printSink: Sink[Int, Future[akka.Done]] =
Sink.foreach(println)
val sumSink: Sink[Int, Future[Int]] =
Sink.fold(0)(_ + _)Flow 정의하기
Flow[In, Out, Mat]은 Source와 Sink 사이에서 각 요소를 변환합니다. Flow는 단독으로 재사용할 수 있으며 어떤 종단점에 연결하기 전에 조합할 수도 있습니다.
여기서는 Flow가 정수를 두 배로 만들고 문자열로 변환합니다.
import akka.stream.scaladsl.Flow
val doubleToString: Flow[Int, String, akka.NotUsed] =
Flow[Int]
.map(_ * 2)
.map(n => s"value=$n")Source를 Sink에 연결하기
via 연산자는 Flow를 Source에 연결하고, to는 Sink를 연결합니다. Source를 to로 Sink에 직접 연결하면 닫혀 있어 실행할 수 있는 청사진인 RunnableGraph가 만들어집니다.
아직 요소가 이동하지는 않으며, 이 단계에서는 토폴로지만 설명합니다.
import akka.stream.scaladsl.{Source, Sink, RunnableGraph}
val graph: RunnableGraph[akka.NotUsed] =
Source(1 to 10).to(Sink.foreach(println))via: Flow 삽입
via를 사용하면 플로를 파이프라인에 삽입할 수 있습니다. Source.via(flow)는 출력 타입이 플로의 출력 타입과 일치하는 새로운 소스를 생성합니다.
via 호출을 연결하면 작고 테스트 가능한 플로 조각으로 긴 변환 파이프라인을 구성할 수 있습니다.
val pipeline =
Source(1 to 10)
.via(Flow[Int].filter(_ % 2 == 0))
.via(Flow[Int].map(_ * 10))
.to(Sink.foreach(println))스테이지 전체의 타입 안전성
컴파일러는 각 스테이지의 출력 타입이 다음 스테이지의 입력 타입과 일치하는지 확인합니다. 타입을 변환하는 중간 플로 없이는 Source[Int]를 Sink[String]에 연결할 수 없습니다.
이 정적 검사를 통해 실행 전에 파이프라인 연결 오류를 발견할 수 있습니다.
// Source[Int] -> Flow[Int, String] -> Sink[String]
val ok =
Source(1 to 3)
.via(Flow[Int].map(_.toString))
.to(Sink.foreach[String](println))머티리얼라이즈된 값
각 블루프린트에는 머티리얼라이즈된 값이 포함됩니다. 이는 스트림이 실행될 때 생성되는 핸들입니다. Sink.fold는 결과를 담은 Future를 머티리얼라이즈합니다. 기본적으로 스테이지를 결합하면 가장 왼쪽의 머티리얼라이즈된 값이 유지됩니다(일반 소스에서는 NotUsed).
toMat과 Keep을 사용하면 양쪽 중 어떤 값을 사용할지 선택할 수 있습니다.
import akka.stream.scaladsl.Keep
import scala.concurrent.Future
val g: RunnableGraph[Future[Int]] =
Source(1 to 100)
.toMat(Sink.fold(0)(_ + _))(Keep.right)재사용 가능한 구성 요소
소스, 플로, 싱크는 변경할 수 없는 값이므로 한 번 정의한 뒤 여러 파이프라인에서 재사용할 수 있습니다. 따라서 이름이 붙은 작은 처리 스테이지들로 라이브러리를 구성하기에 좋습니다.
구문 분석을 위해 정의한 플로를 파일 파이프라인과 HTTP 파이프라인 모두에 삽입할 수 있습니다.
val parse: Flow[String, Int, akka.NotUsed] =
Flow[String].map(_.trim.toInt)
val fromFile = lines.via(parse)
val fromHttp = requestBody.via(parse)일반적인 소스 생성자
Akka Streams에는 다양한 소스 팩토리가 포함되어 있습니다. Source.single, Source.repeat, 시간에 따라 값을 내보내는 Source.tick, Future에서 생성하는 Source.future, 그리고 Source.empty가 있습니다.
적절한 생성자를 선택하면 데이터 생산자의 의도를 명확하게 표현할 수 있습니다.
import scala.concurrent.duration._
val ticks = Source.tick(0.seconds, 1.second, "tick")
val onceF = Source.future(scala.concurrent.Future.successful(42))
val forever = Source.repeat("x")일반적인 싱크 생성자
마찬가지로 싱크에는 첫 번째 요소를 Future로 반환하는 Sink.head, 모든 요소를 Seq로 모으는 Sink.seq, 요소를 소비한 뒤 버리는 Sink.ignore, 그리고 Sink.last가 있습니다.
결과를 코드로 돌려보내는 파이프라인에서는 Sink.seq와 Sink.fold를 가장 많이 사용합니다.
import scala.concurrent.Future
val collect: Sink[Int, Future[Seq[Int]]] = Sink.seq
val firstOne: Sink[Int, Future[Int]] = Sink.head
val drain: Sink[Int, Future[akka.Done]] = Sink.ignore빠른 확인
선형 스테이지 타입에 대한 이해도를 확인해 보세요.
복습
세 가지 선형 구성 요소를 배웠습니다. 소스는 데이터를 생성하고, 플로는 변환하며, 싱크는 소비합니다. 이들은 변경할 수 없고 재사용 가능한 블루프린트이며, via와 to로 연결합니다.
소스와 싱크를 연결하면 머티리얼라이즈된 값을 담은 RunnableGraph가 생성되지만, 실행하기 전까지는 데이터가 이동하지 않습니다. 다음에는 더 다양한 연산자를 사용해 스트림을 변환해 보겠습니다.
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개 중 1번째 강의입니다.
“소스, 플로, 싱크” 강의는 얼마나 걸리나요?
대부분의 CoddyKit 강의는 약 5~10분이 소요됩니다. 각 강의는 간결하고 인터랙티브하여 꾸준한 진행이 가능하며, 웹과 앱에서 중단한 부분부터 바로 시작할 수 있습니다.
이 Scala for Backend Engineering & Functional Programming 강의에서 코드를 작성하고 실행할 수 있나요?
네. 모든 Scala for Backend Engineering & Functional Programming 강의에는 내장 코드 에디터가 포함되어 있으므로, 브라우저에서 바로 실제 코드를 작성하고 실행한 후 즉시 AI 피드백을 받을 수 있습니다 — 로컬 설정이 필요 없습니다.
이 강의의 모든 강의
- 소스, 플로, 싱크
- 스트림 변환하기
- 백프레셔
- 파이프라인 실행하기