Apache Beam
고친 사람 github-actions[bot]
Apache Beam 은 데이터 처리 과정을 한 번 적어 두고 여러 실행 엔진에서 돌리게 해 주는 도구입니다. 처리 로직을 적은 코드와 그 코드를 돌리는 엔진을 떼어 놓습니다. 끝이 있는 파일 데이터와 끝없이 들어오는 데이터를 같은 코드로 다룹니다.
쉽고 빠른 이해
Beam 은 데이터를 읽고 바꾸고 쓰는 과정을 코드로 한 번 적게 해 줍니다. 문장에서 단어 수를 세는 코드를 내 컴퓨터에서 돌려 본 뒤, 같은 코드를 Flink 나 Spark 클러스터에 올리는 식입니다.
이게 없으면 엔진을 바꿀 때마다 코드를 새로 짜야 합니다. 하루치 파일을 모아 처리하는 코드와 들어오는 대로 처리하는 코드도 따로 짜야 합니다.
- 읽기, 바꾸기, 쓰기 단계를 이어 붙인 파이프라인을 코드로 적습니다
- 어느 엔진에서 돌릴지 실행 옵션으로 고릅니다
- Beam 이 파이프라인을 그 엔진이 알아듣는 작업으로 바꿔 넘깁니다
끝없이 들어오는 데이터는 끝이 오지 않아 한꺼번에 처리할 수 없습니다. 그래서 시간 구간으로 잘라 구간마다 나눠 처리합니다.
대가는 엔진 위에 한 겹이 더 얹힌다는 점입니다. 엔진 하나에만 있는 기능은 쓰기 어렵습니다. 엔진마다 받아 주는 Beam 기능도 조금씩 다릅니다.
상세
Apache Beam 은 Apache 소프트웨어 재단이 관리하는 오픈 소스 프로젝트입니다. 데이터를 읽고 변환하고 쓰는 과정을 코드로 적게 해 줍니다. 적은 코드를 Beam 이 직접 돌리지는 않습니다. 따로 있는 실행 엔진에 넘겨 돌립니다.
이 절은 Beam 이 코드와 엔진을 왜 떼어 놓았는지부터 봅니다.
엔진마다 다시 짜던 코드
데이터가 컴퓨터 한 대에 다 안 들어가면 여러 대에 나눠 처리합니다. 이것을 분산 처리라고 합니다. 로그 수십 테라바이트에서 사용자별 방문 수를 세는 일이 그렇습니다.
분산 처리를 맡는 프로그램을 이 문서에서는 실행 엔진, 줄여서 엔진이라고 부릅니다. Spark 와 Flink 가 이런 엔진입니다. 엔진은 일을 여러 대에 나눠 줍니다. 한 대가 죽으면 그 몫을 다른 대에서 다시 돌립니다.
엔진마다 코드를 짜는 방법이 다릅니다. Spark 용으로 짠 코드는 Flink 에서 돌지 않습니다. 엔진을 바꾸려면 처리 로직을 새로 짜야 합니다.
데이터의 모양도 코드를 갈랐습니다. 하루치 파일을 모아 한 번에 처리하는 일을 배치 처리라고 합니다. 끝없이 들어오는 데이터를 들어오는 대로 처리하는 일은 스트림 처리라고 합니다. 같은 집계라도 이 둘을 따로 짜는 일이 흔했습니다.
Beam 은 두 갈림을 한꺼번에 풉니다. 처리 로직은 Beam 의 방식으로 한 번 적습니다. 어느 엔진에서 돌릴지는 실행할 때 고릅니다. 배치든 스트림이든 같은 코드로 적습니다.
Beam 이라는 이름도 배치(Batch)와 스트림(Stream)을 합쳐 지은 것입니다.
Beam 은 이렇게 적는 방법을 SDK(Software Development Kit, 소프트웨어 개발 키트)로 내놓습니다. SDK 는 프로그램을 짜는 데 쓰는 라이브러리와 도구의 묶음입니다. Java, Python, Go 용 SDK 가 있습니다.
파이프라인
Beam 에서는 데이터 처리 작업 하나를 파이프라인이라고 부릅니다. 입력을 읽는 단계, 데이터를 바꾸는 단계, 결과를 쓰는 단계가 한 파이프라인 안에 함께 들어갑니다. 코드에서는 Pipeline 객체 하나가 이 작업 전체를 나타냅니다.
파이프라인을 만들고 단계를 붙이는 프로그램을 드라이버 프로그램이라고 부릅니다. 드라이버 프로그램은 데이터를 직접 처리하지 않습니다. 무엇을 어떤 순서로 할지 적어서 엔진에 넘기는 일만 합니다.
PCollection
단계와 단계 사이를 흐르는 데이터 묶음이 PCollection 입니다. 여러 컴퓨터에 나뉘어 있을 수 있는 데이터 집합을 가리킵니다. 로그 한 줄 한 줄이 요소인 PCollection 을 떠올리면 됩니다.
PCollection 은 한번 만들면 고치지 않습니다. 데이터를 바꾸고 싶으면 바꾼 결과를 새 PCollection 으로 받습니다.
요소의 순서도 정해져 있지 않습니다. 엔진이 여러 대에서 나눠 처리해서 어느 요소가 먼저 올지 약속할 수 없습니다.
PCollection 은 끝이 있을 수도 있고 없을 수도 있습니다. 파일에서 읽은 PCollection 은 끝이 있습니다. 이것을 유한(bounded) PCollection 이라고 합니다.
메시지 큐에서 계속 읽어 오는 PCollection 은 끝이 없습니다. 이것은 무한(unbounded) PCollection 입니다.
PTransform
데이터를 바꾸는 단계 하나가 PTransform 입니다. PCollection 을 입력으로 받아 새 PCollection 을 내놓습니다.
읽기와 쓰기도 PTransform 입니다. 읽기는 바깥에서 데이터를 가져와 PCollection 을 내놓습니다. 쓰기는 PCollection 을 받아 바깥에 적습니다.
바깥 저장소와 데이터를 주고받는 PTransform 을 입출력 커넥터(I/O connector, Input/Output connector)라고 부릅니다. 파일, 데이터베이스, 메시지 큐마다 커넥터가 따로 있습니다.
한 PCollection 을 여러 PTransform 이 나눠 받을 수 있습니다. 그래서 파이프라인은 한 줄이 아니라 갈라지는 그래프가 됩니다. 화살표를 따라가도 처음으로 되돌아오는 길은 없습니다. 이런 그래프를 유향 비순환 그래프(DAG, Directed Acyclic Graph)라고 합니다.
아래 그림은 로그를 읽어 두 갈래로 처리하는 파이프라인입니다. 네모 칸은 PTransform 입니다. 양끝이 둥근 칸은 PCollection 입니다. 「로그 줄」 PCollection 하나를 두 PTransform 이 받습니다. 한쪽은 사용자별 방문 수를 셉니다. 다른 쪽은 오류 줄만 골라 파일에 씁니다.
flowchart TD
R["로그 읽기"] --> C1(["로그 줄"])
C1 --> U["사용자 뽑기"]
C1 --> E["오류 줄 거르기"]
U --> C2(["사용자"])
E --> C3(["오류 줄"])
C2 --> N["사용자별 방문 수 세기"]
C3 --> W["오류 파일 쓰기"]
자주 쓰는 PTransform
몇몇 변환은 요소가 키와 값의 짝이기를 바랍니다. (사용자 번호, 방문 1회)처럼 앞이 키이고 뒤가 값인 요소입니다. 이런 짝을 키-값 쌍이라고 부릅니다.
Beam 은 자주 쓰는 변환을 미리 만들어 둡니다. 대부분의 파이프라인은 아래 넷을 이어 붙여 짭니다.
| 변환 | 하는 일 | 예 |
|---|---|---|
ParDo |
요소마다 같은 함수를 돌린다. 요소 하나에서 0개 이상을 내놓는다 | 로그 줄에서 사용자 번호를 뽑는다 |
GroupByKey |
같은 키를 가진 값을 한데 모은다 | 사용자별로 방문 기록을 모은다 |
Combine |
여러 값을 하나로 줄인다 | 사용자별 방문 수를 더한다 |
Flatten |
PCollection 여럿을 하나로 합친다 | 서버 세 대의 로그를 합친다 |
ParDo 는 요소 하나만 보고 일합니다. 그래서 엔진이 요소를 여러 대에 흩어 동시에 돌릴 수 있습니다. ParDo 에 넘겨 요소마다 돌리는 함수를 DoFn 이라고 부릅니다.
GroupByKey 는 요소 하나만 보고는 일할 수 없습니다. 같은 키의 값이 여러 대에 흩어져 있어서 네트워크로 한 대에 모아야 합니다. 이렇게 옮기는 일을 셔플이라고 부릅니다. 셔플은 파이프라인에서 비용이 가장 큰 단계가 되기 쉽습니다.
짧은 코드 한 벌
아래는 Python SDK 로 짠 단어 세기 파이프라인입니다. 문장 두 개를 넣고 단어마다 몇 번 나왔는지 셉니다. 줄 오른쪽 주석은 그 줄에서 나오는 값입니다.
import apache_beam as beam
with beam.Pipeline() as p:
(p
| beam.Create(["a b", "b"])
| beam.FlatMap(str.split) # a, b, b
| beam.combiners.Count.PerElement()
| beam.Map(print)) # ('a', 1) ('b', 2)
| 는 단계를 이어 붙이는 연산자입니다. 셸의 파이프처럼 왼쪽 결과가 오른쪽으로 흘러갑니다.
첫 | 는 파이프라인 p 에 시작 단계를 붙입니다. 그 뒤의 | 는 앞 단계가 내놓은 PCollection 에 다음 PTransform 을 붙입니다.
Create 는 코드에 적은 값으로 PCollection 을 만듭니다. 이 코드에서는 문장 두 개가 요소가 됩니다.
FlatMap 은 요소 하나에서 요소 여럿을 내놓는 ParDo 입니다. 이 코드에서는 문장 하나를 단어 여럿으로 쪼갭니다.
Count.PerElement 는 같은 요소가 몇 번 나왔는지 셉니다. 안에서는 먼저 단어마다 (단어, 1) 짝을 만듭니다. 그다음 같은 단어의 1 을 Combine 으로 더합니다.
Map 은 요소 하나에서 요소 하나를 내놓는 ParDo 입니다. 여기서는 요소마다 print 를 불러 화면에 찍습니다. 찍히는 두 줄의 순서는 실행마다 바뀔 수 있습니다. PCollection 에는 순서가 없기 때문입니다.
with 블록이 끝날 때 파이프라인이 돕니다. 이 코드는 어느 엔진에서 돌지 정하지 않았습니다. 그러면 Beam 은 내 컴퓨터 한 대에서 파이프라인을 돌립니다.
러너
파이프라인을 엔진에 넘겨 돌리는 부품이 러너입니다. 러너는 PTransform 그래프를 엔진이 알아듣는 작업으로 바꿉니다. 엔진마다 러너가 하나씩 있습니다.
어느 러너를 쓸지는 코드가 아니라 실행 옵션으로 정합니다. 같은 파이프라인 코드를 옵션만 바꿔 여러 엔진에 올릴 수 있습니다. 개발할 때는 Direct 러너로 확인합니다. Direct 러너는 따로 엔진을 두지 않고 내 컴퓨터에서 파이프라인을 직접 돌립니다. 운영할 때는 클러스터에서 도는 엔진의 러너로 바꿉니다.
러너마다 받아 주는 기능이 다릅니다. Beam 의 모든 기능을 모든 러너가 똑같이 돌리지는 못합니다. Beam 프로젝트는 러너별로 어떤 기능을 받는지 표로 정리해 둡니다. 이 표를 기능 매트릭스(capability matrix)라고 부릅니다.
이벤트 시간과 윈도
무한 PCollection 에서는 전부 모아서 세는 방법이 통하지 않습니다. 데이터가 끝나지 않으니 전부가 모이는 때가 오지 않습니다. 그래서 시간을 구간으로 잘라 구간마다 셉니다. 이 구간을 윈도라고 부릅니다.
시간을 자를 때 쓸 수 있는 시각은 둘입니다. 이벤트 시간은 일이 일어난 시각입니다. 사용자가 버튼을 누른 순간이 그렇습니다. 처리 시간은 그 기록이 파이프라인에 도착해 처리되는 시각입니다.
둘은 자주 어긋납니다. 휴대폰이 지하철에서 신호를 잃으면 10시에 누른 기록이 10시 20분에 도착합니다. 처리 시간으로 자르면 이 클릭은 10시 20분 구간에 들어갑니다.
Beam 은 요소마다 시각을 하나씩 붙여 둡니다. 윈도는 이 시각으로 자릅니다. 이 시각에 이벤트 시간을 넣어 두면 늦게 도착한 클릭도 10시 구간에 들어갑니다.
윈도를 자르는 방법은 흔히 셋을 씁니다.
| 윈도 | 자르는 법 | 예 |
|---|---|---|
| 고정 윈도 (fixed window) | 같은 길이로 겹치지 않게 자른다 | 5분마다 주문 수 |
| 슬라이딩 윈도 (sliding window) | 같은 길이의 구간을 조금씩 밀며 겹치게 자른다 | 1분마다 내는 최근 10분 평균 |
| 세션 윈도 (session window) | 요소 사이 빈 시간이 정한 길이를 넘으면 끊는다 | 사용자 한 명의 방문 한 번 |
고정 윈도와 슬라이딩 윈도는 시계만 보고 자릅니다. 세션 윈도는 데이터를 보고 자릅니다.
세션 윈도는 키마다 따로 자릅니다. 사용자 번호가 키라면 사용자마다 방문이 따로 끊깁니다. 그래서 세션 윈도는 키마다 길이가 다릅니다.
워터마크와 트리거
이벤트 시간으로 자르면 새 문제가 생깁니다. 10시 구간을 언제 닫아야 할지 모릅니다. 10시에 일어난 일이 언제 더 도착할지 알 수 없기 때문입니다.
워터마크는 이 물음에 답하는 추정값입니다. 「이 시각보다 이른 이벤트는 거의 다 왔다」고 보는 시각입니다. 워터마크가 10시 5분을 넘으면 엔진은 10시부터 10시 5분까지의 윈도가 다 찼다고 보고 결과를 냅니다.
워터마크는 추정이라 틀릴 수 있습니다. 워터마크가 윈도 끝을 지난 뒤 도착한 요소를 늦은 데이터라고 부릅니다.
Beam 에서는 윈도 끝이 지난 뒤에도 얼마 동안 늦은 데이터를 받을지 정할 수 있습니다. 이 기간을 허용 지연(allowed lateness)이라고 부릅니다. 허용 지연보다 더 늦게 온 요소는 버립니다.
트리거는 윈도의 결과를 언제 내보낼지 정합니다. 워터마크가 윈도 끝을 지날 때 한 번 내는 것이 기본입니다. 허용 지연 안에 늦은 데이터가 오면 고친 결과를 다시 냅니다. 결과를 일찍 보고 싶으면 윈도가 닫히기 전에도 1분마다 중간 결과를 내게 할 수 있습니다.
아래 그림은 기본 트리거를 쓸 때 윈도 하나가 거치는 상태입니다. 처음에는 요소를 모읍니다. 워터마크가 윈도 끝을 지나면 결과를 냅니다. 허용 지연 동안은 늦은 데이터가 올 때마다 결과를 다시 냅니다.
stateDiagram-v2
state "요소를 모은다" as A
state "결과를 냈다" as B
state "닫혔다" as C
[*] --> A
A --> B: 워터마크가 윈도 끝을 지남
B --> B: 늦은 데이터가 오면 고친 결과를 다시 냄
B --> C: 허용 지연이 끝남
C --> [*]
한 겹 더 얹는 대가
Beam 은 여러 엔진이 함께 할 수 있는 일을 중심으로 모델을 정합니다. 엔진 하나에만 있는 기능은 Beam 코드에서 쓰기 어렵습니다. 그 기능이 꼭 필요하면 그 엔진의 코드로 직접 짜는 편이 맞습니다.
문제를 찾을 때도 한 겹을 더 거칩니다. 엔진의 모니터링 화면에는 러너가 바꿔 놓은 작업이 보입니다. 이 작업은 내가 적은 PTransform 과 한 줄씩 맞지 않을 때가 있습니다. 그래서 어느 단계가 늦는지 찾는 데 품이 듭니다.
엔진을 바꿀 수 있다는 약속에도 조건이 붙습니다. 옮겨 갈 러너가 내 파이프라인이 쓰는 기능을 받아 주는지 기능 매트릭스에서 먼저 봐야 합니다.
Beam 을 고르는 경우
같은 집계를 배치와 스트림 둘 다로 돌려야 할 때 맞습니다. 지난 한 달치 로그를 다시 계산하는 코드와 실시간으로 세는 코드를 하나로 둡니다. 나중에 엔진을 바꿀 여지를 남기고 싶을 때도 맞습니다.
데이터가 컴퓨터 한 대에 다 들어가면 Beam 을 세울 까닭이 적습니다. 분산 처리 자체가 필요 없기 때문입니다.
분석용 데이터를 모아 두는 데이터베이스를 데이터 웨어하우스라고 합니다. 그 안의 데이터를 SQL(Structured Query Language) 몇 줄로 바꾸면 끝나는 일에도 Beam 이 필요 없습니다. 웨어하우스 안에서 SQL 로 돌리면 데이터를 밖으로 옮기지 않아도 됩니다.
맞물림
Beam 은 혼자 데이터를 처리하지 못합니다. 파이프라인을 받아 돌릴 엔진과, 데이터를 읽어 올 원천이 있어야 합니다. 이 절은 Beam 이 가장 자주 붙는 넷이 어느 방향으로 붙는지 봅니다.
Flink, Spark, Kafka 는 Beam 과 같은 Apache 재단의 프로젝트입니다. Dataflow 는 Beam 이 비롯된 Google 의 서비스입니다.
Beam 파이프라인을 받아 돌리는 Flink
Flink 는 끝없이 들어오는 데이터를 이어서 처리하는 엔진입니다. Flink 러너는 Beam 의 PTransform 그래프를 Flink 작업으로 바꿔 Flink 클러스터에 넘깁니다. Beam 이 넘기고 Flink 가 받아 돌리는 방향입니다.
Flink 는 이벤트 시간과 워터마크를 스스로 갖추고 있습니다. 그래서 Beam 의 스트림 기능 대부분을 받아 줍니다.
Java 가 아닌 SDK 로 짜면 비용이 더 듭니다. Flink 는 JVM(Java Virtual Machine, 자바 가상 머신) 위에서 돕니다. 그래서 Python 으로 적은 DoFn 을 Flink 안에서 바로 돌리지 못합니다.
Flink 옆에 Python 프로세스를 따로 띄워 그 함수를 돌립니다. 요소가 두 프로세스 사이를 오가므로 그만큼 시간이 더 듭니다.
Beam 파이프라인을 받아 돌리는 Spark
Spark 는 데이터를 여러 대의 메모리에 나눠 올려 처리하는 엔진입니다. 배치 처리에 널리 쓰입니다. Spark 러너가 Beam 파이프라인을 Spark 작업으로 바꿔 넘깁니다.
배치 파이프라인은 Spark 러너로 옮기기 수월합니다. 스트림 파이프라인은 옮기기 전에 따져 볼 것이 있습니다. Spark 는 들어오는 데이터를 작은 배치로 끊어 처리합니다. 이 방식 때문에 Beam 의 스트림 기능 일부를 못 받는 경우가 있습니다.
그래서 스트림 파이프라인을 Spark 러너에 올릴 때는 기능 매트릭스를 먼저 봅니다.
Beam 파이프라인을 받아 돌리는 Google Cloud Dataflow
Google Cloud Dataflow 는 Google Cloud 가 운영하는 관리형 서비스입니다. 관리형이란 클러스터를 내가 띄우고 돌보지 않는다는 뜻입니다. Dataflow 러너가 파이프라인을 이 서비스에 올립니다. 서비스는 작업자 머신을 띄우고 데이터 양에 따라 늘리고 줄입니다.
Beam SDK 는 원래 이 서비스용 SDK 에서 나왔습니다. Google 이 그 SDK 를 Apache 재단에 넘긴 것이 Beam 의 출발입니다. 그래서 Dataflow 에서 도는 코드는 곧 Beam 코드입니다.
대가는 한 클라우드에 묶이는 것입니다. 작업자 머신과 과금이 Google Cloud 에 있습니다. 코드는 Beam 이라 다른 러너로 옮길 길은 남습니다.
Beam 이 읽고 쓰는 Kafka
Kafka 는 메시지를 토픽에 차례로 쌓아 두는 이벤트 스트리밍 플랫폼입니다. 이번에는 Beam 이 부르는 쪽입니다. Kafka 용 입출력 커넥터가 토픽을 읽어 무한 PCollection 을 만듭니다. 쓰기 커넥터로 결과를 다른 토픽에 다시 씁니다.
토픽 안에서 메시지의 위치를 오프셋이라고 부릅니다. 메시지마다 차례로 번호가 붙어 있어서 어디까지 읽었는지를 이 번호로 말할 수 있습니다.
엔진은 처리하던 상태를 주기적으로 저장해 둡니다. 이 저장본이 체크포인트입니다. 앞에서 본 것처럼 한 대가 죽으면 엔진은 그 몫을 다시 돌립니다. 이때 처음부터가 아니라 마지막 체크포인트부터 이어 갑니다.
Kafka 커넥터는 어디까지 읽었는지를 이 체크포인트에 오프셋으로 함께 적어 둡니다. 그래서 엔진이 다시 시작하면 저장된 오프셋부터 이어 읽습니다.
붙일 때 정할 것이 하나 있습니다. 커넥터가 요소에 어떤 시각을 붙일지입니다. 처리 시간이 붙으면 앞에서 본 이벤트 시간 윈도가 뜻대로 잘리지 않습니다. 메시지에 적힌 이벤트 시각을 쓰도록 커넥터에 정해 둡니다.
관련 항목
Apache Beam 이 속하는 상위 분류
데이터 엔지니어링 · 분산 처리 · 데이터 파이프라인 · 오픈 소스 · Apache Software Foundation
Apache Beam 파이프라인을 이루는 구성 요소
파이프라인 · PCollection · PTransform · DoFn · 드라이버 프로그램 · 입출력 커넥터 · 유향 비순환 그래프
Apache Beam 이 미리 만들어 둔 변환
변환 · ParDo · GroupByKey · Combine · Flatten · 키-값 쌍 · 셔플
Apache Beam 파이프라인을 받아 돌리는 러너와 엔진
러너 · Apache Flink · Apache Spark · Google Cloud Dataflow · Apache Samza · Direct 러너 · 클러스터 · 관리형 서비스
Apache Beam 이 한 코드로 묶는 처리 방식
배치 처리 · 스트림 처리 · 마이크로 배치 · 람다 아키텍처 · 카파 아키텍처
Apache Beam 이 무한 데이터를 자르는 시간 개념
윈도 · 세션 윈도 · 이벤트 시간 · 처리 시간 · 워터마크 · 트리거 · 늦은 데이터
Apache Beam 이 읽고 쓰는 데이터 원천
Apache Kafka · 메시지 큐 · 토픽 · 오프셋 · 체크포인트 · Amazon S3 · BigQuery · 데이터 웨어하우스
Apache Beam 파이프라인을 짜는 언어와 도구
SDK · Java · Python · Go · JVM · SQL
Apache Beam 이 비롯된 Google 의 처리 모델
맵리듀스 · FlumeJava · MillWheel · Dataflow 모델
Apache Beam 과 같은 역할을 두고 겨루는 처리 도구
Kafka Streams · Spark Structured Streaming · Apache Storm · dbt
Apache Beam 파이프라인을 일정에 맞춰 띄우는 오케스트레이터
워크플로 오케스트레이션 · Apache Airflow · Dagster · Prefect
다른 이름: Beam · 아파치 빔 · 빔