Scala for Backend Engineering & Functional Programming · 강의

스트림 변환하기

흐르는 데이터를 매핑하고 필터링해 보세요.

레슨 2/413개 단계

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

변환으로서의 연산자

Akka Streams는 Source와 Flow에 다양한 연산자를 제공합니다. 이 연산자들은 Scala의 컬렉션 API와 비슷하지만 비동기적으로 실행되며 역압을 준수합니다.

각 연산자는 새로운 블루프린트를 반환하므로 스트림이 실제로 실행되기 전에 변환을 선언적으로 조합할 수 있습니다.

map과 filter

map은 모든 요소에 동기 함수를 적용하고, filter는 조건을 만족하지 않는 요소를 제거합니다. 두 연산자는 요소 단위 변환에서 가장 기본적으로 사용됩니다.

두 연산자 모두 순서를 유지하며 완료와 실패를 downstream으로 전달합니다.

val flow =
  Flow[Int]
    .filter(_ % 2 == 0)
    .map(n => n * n)

일대다 변환을 위한 mapConcat

하나의 입력에서 여러 출력이 나와야 할 때는 mapConcat을 사용합니다. 이 연산자는 iterable을 반환하는 함수를 받아 결과를 평탄화하여 스트림에 넣습니다.

빈 컬렉션을 반환하면 해당 요소가 사실상 제거됩니다.

val explode: Flow[String, String, akka.NotUsed] =
  Flow[String].mapConcat(line => line.split(",").toList)

val words = Source(List("a,b", "c,d,e"))
  .via(explode)

grouped와 sliding

grouped(n)은 연속된 요소를 최대 n개씩 묶어 Seq로 만듭니다. 대량 데이터베이스 쓰기에 유용합니다. sliding(n)은 서로 겹치는 윈도우를 내보냅니다.

배치 처리를 사용하면 입출력이 많은 파이프라인에서 요소별 처리 오버헤드를 줄일 수 있습니다.

val batches: Source[Seq[Int], akka.NotUsed] =
  Source(1 to 1000).grouped(100)

val windows =
  Source(1 to 10).sliding(3, step = 1)

scan과 fold

scan은 각 요소를 처리한 뒤 현재 누적값을 내보내므로 상태가 변화하는 스트림을 만들 수 있습니다. fold는 upstream이 완료될 때 최종 누적값만 한 번 내보냅니다.

실시간 카운터에는 scan을, 최종 집계에는 fold를 사용하세요.

val running =
  Source(1 to 5).scan(0)(_ + _) // 0,1,3,6,10,15

val total =
  Source(1 to 5).fold(0)(_ + _) // 15

비동기 작업을 위한 mapAsync

mapAsync(parallelism)은 Future를 반환하는 함수를 호출하고 결과를 순서대로 내보냅니다. 동시에 실행되는 Future의 수는 parallelism 이하로 제한됩니다.

순서가 중요한 데이터베이스 조회나 HTTP 요청 같은 비동기 호출에 사용하세요.

import scala.concurrent.Future

val enriched =
  Flow[UserId]
    .mapAsync(parallelism = 4)(id => lookup(id))

def lookup(id: UserId): Future[User] = ???

mapAsyncUnordered

mapAsyncUnordered는 mapAsync와 비슷하지만 입력 순서를 무시하고 각 결과가 완료되는 즉시 내보냅니다.

downstream에서 순서를 중요하게 여기지 않는다면 처리량을 높일 수 있습니다. 느린 Future가 더 빠른 Future를 막지 않기 때문입니다.

val fast =
  Flow[UserId]
    .mapAsyncUnordered(parallelism = 8)(id => lookup(id))

statefulMapConcat을 사용한 상태 기반 변환

요소별 변환에 변경 가능한 로컬 상태가 필요하다면 statefulMapConcat을 사용하세요. 이 연산자는 머티리얼라이제이션마다 새로운 상태를 만들고 출력 iterable을 반환합니다.

스트림 실행 간에 상태를 공유하지 않으면서 카운터나 버퍼를 유지하는 안전한 방법입니다.

val withIndex: Flow[String, (Int, String), akka.NotUsed] =
  Flow[String].statefulMapConcat { () =>
    var i = 0
    elem => { i += 1; List((i, elem)) }
  }

시간 기반 연산자

스트림은 내용뿐 아니라 시간에 따라서도 변환할 수 있습니다. throttle은 출력 속도를 제한하고, groupedWithin은 크기나 경과 시간에 따라 배치로 묶으며, takeWithin은 지속 시간을 제한합니다.

이 연산자들은 외부 API의 요청 속도를 제한할 때 필수적입니다.

import scala.concurrent.duration._

val limited =
  Source(1 to 1000)
    .throttle(10, 1.second)
    .groupedWithin(100, 500.millis)

변환 중 오류 처리

연산자 내부에서 예외가 발생하면 기본적으로 전체 스트림이 실패합니다. 대신 감독 전략을 사용해 resume으로 잘못된 요소를 버리거나 스테이지를 restart할 수 있습니다.

Flow에 withAttributes를 사용해 전략을 연결하세요.

import akka.stream.{ActorAttributes, Supervision}

val safe =
  Flow[String].map(_.toInt)
    .withAttributes(
      ActorAttributes.supervisionStrategy(_ => Supervision.Resume))

플로 조합하기

작은 플로는 via를 사용해 더 큰 플로로 조합할 수 있으며, 그 결과 하나의 재사용 가능한 플로가 만들어집니다. 이렇게 하면 각 변환의 역할이 명확해지고 독립적으로 테스트할 수 있습니다.

조합된 플로의 입력 타입은 첫 번째 플로의 입력 타입이고, 출력 타입은 마지막 플로의 출력 타입입니다.

val parse  = Flow[String].map(_.toInt)
val square = Flow[Int].map(n => n * n)

val parseAndSquare: Flow[String, Int, akka.NotUsed] =
  parse.via(square)

빠른 확인

비동기 변환과 그 순서 보장에 대해 생각해 보세요.

복습

요소별 map/filter, 일대다 변환을 위한 mapConcat, 배치 처리를 위한 grouped, scan과 fold를 사용한 누적, 그리고 mapAsync를 통한 비동기 작업 등 다양한 변환 연산자를 살펴보았습니다.

또한 상태 기반 변환, throttle과 같은 시간 기반 연산자, 오류를 위한 감독 전략, via를 사용한 플로 조합 방법도 배웠습니다. 다음에는 역압이 이러한 스테이지를 안전하게 유지하는 방법을 알아보겠습니다.

무료로 시작

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개 중 2번째 강의입니다.

“스트림 변환하기” 강의는 얼마나 걸리나요?

대부분의 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(으)로 돌아가기