데이터 파이프라인과 ETL
데이터 파이프라인과 ETL (Data Pipelines & ETL)
Apache Flink의 매우 흔한 사용 사례 중 하나는 하나 이상의 소스에서 데이터를 가져와 일부 변환 및/또는 강화(enrichment)를 수행한 다음 결과를 어딘가에 저장하는 ETL(extract, transform, load) 파이프라인을 구현하는 것입니다. 이 섹션에서는 Flink의 DataStream API를 사용해 이런 종류의 애플리케이션을 구현하는 방법을 살펴봅니다.
출처: 문서
본문
Apache Flink의 매우 흔한 사용 사례 중 하나는 ETL 파이프라인을 구현하는 것입니다. 이 섹션에서는 Flink의 DataStream API를 사용해 이런 종류의 애플리케이션을 구현하는 방법을 살펴봅니다.
Flink의 Table and SQL APIs는 많은 ETL 사용 사례에 잘 맞습니다. 하지만 궁극적으로 DataStream API를 직접 사용하는지 여부와 무관하게 여기에 제시된 기본 개념을 확실히 이해하는 것은 가치가 있습니다.
상태 없는 변환 (Stateless Transformations)
이 섹션은 상태 없는 변환을 구현하는 데 사용되는 기본 연산인 map()과 flatmap()을 다룹니다. 이 섹션의 예시는 flink-training-repo의 실습에서 사용되는 Taxi Ride 데이터에 익숙하다고 가정합니다.
map()
첫 번째 연습에서 taxi ride 이벤트 스트림을 필터링했습니다. 같은 코드 베이스에 위치(경도, 위도)를 약 100x100미터 크기의 영역을 가리키는 그리드 셀에 매핑하는 정적 메서드 GeoUtils.mapToGridCell(float lon, float lat)를 제공하는 GeoUtils 클래스가 있습니다.
이제 각 이벤트에 startCell과 endCell 필드를 추가해 taxi ride 객체 스트림을 강화해 봅시다. TaxiRide를 확장하고 이 필드들을 추가하는 EnrichedRide 객체를 만들 수 있습니다.
public static class EnrichedRide extends TaxiRide {
public int startCell;
public int endCell;
public EnrichedRide() {}
public EnrichedRide(TaxiRide ride) {
this.rideId = ride.rideId;
this.isStart = ride.isStart;
...
this.startCell = GeoUtils.mapToGridCell(ride.startLon, ride.startLat);
this.endCell = GeoUtils.mapToGridCell(ride.endLon, ride.endLat);
}
public String toString() {
return super.toString() + "," +
Integer.toString(this.startCell) + "," +
Integer.toString(this.endCell);
}
}
그 다음 스트림을 변환하는 애플리케이션을 만들 수 있습니다.
DataStream<TaxiRide> rides = env.addSource(new TaxiRideSource(...));
DataStream<EnrichedRide> enrichedNYCRides = rides
.filter(new RideCleansingSolution.NYCFilter())
.map(new Enrichment());
enrichedNYCRides.print();
이렇게 MapFunction으로 말입니다.
public static class Enrichment implements MapFunction<TaxiRide, EnrichedRide> {
@Override
public EnrichedRide map(TaxiRide taxiRide) throws Exception {
return new EnrichedRide(taxiRide);
}
}
flatmap()
MapFunction은 일대일 변환을 수행할 때만 적합합니다. 들어오는 모든 스트림 요소에 대해 map()은 변환된 요소 하나를 내보냅니다. 그 외의 경우에는 flatmap()을 사용하고 싶을 것입니다.
DataStream<TaxiRide> rides = env.addSource(new TaxiRideSource(...));
DataStream<EnrichedRide> enrichedNYCRides = rides
.flatMap(new NYCEnrichment());
enrichedNYCRides.print();
FlatMapFunction과 함께 말입니다.
public static class NYCEnrichment implements FlatMapFunction<TaxiRide, EnrichedRide> {
@Override
public void flatMap(TaxiRide taxiRide, Collector<EnrichedRide> out) throws Exception {
FilterFunction<TaxiRide> valid = new RideCleansing.NYCFilter();
if (valid.filter(taxiRide)) {
out.collect(new EnrichedRide(taxiRide));
}
}
}
이 인터페이스에서 제공하는 Collector 덕분에 flatmap() 메서드는 원하는 만큼 많은 스트림 요소를 내보낼 수 있으며, 전혀 내보내지 않을 수도 있습니다.
키가 있는 스트림 (Keyed Streams)
keyBy()
하나의 속성 주위로 스트림을 파티셔닝하여 해당 속성의 같은 값을 가진 모든 이벤트가 함께 그룹화되도록 하는 것은 종종 매우 유용합니다. 예를 들어 각 그리드 셀에서 시작하는 가장 긴 택시 ride를 찾고 싶다고 가정해 봅시다. SQL 쿼리 관점에서 생각하면 이는 startCell로 일종의 GROUP BY를 수행하는 것을 의미하며, Flink에서는 keyBy(KeySelector)로 수행합니다.
rides
.flatMap(new NYCEnrichment())
.keyBy(enrichedRide -> enrichedRide.startCell);
모든 keyBy는 스트림을 재파티셔닝하는 네트워크 셔플을 발생시킵니다. 일반적으로 이것은 네트워크 통신과 직렬화/역직렬화를 수반하므로 상당히 비쌉니다.
키가 계산되는 방식 (Keys are computed)
KeySelector는 이벤트에서 키를 추출하는 것에 국한되지 않습니다. 결과 키가 결정적(deterministic)이고 hashCode()와 equals()의 유효한 구현을 가진다면 원하는 어떤 방식으로든 키를 계산할 수 있습니다. 이 제약은 난수를 생성하거나 Arrays나 Enums를 반환하는 KeySelector를 배제하지만, 요소들이 이 규칙을 따른다면 예를 들어 Tuples나 POJOs를 사용한 복합 키(composite keys)를 가질 수 있습니다.
키는 필요할 때마다 다시 계산되며 스트림 레코드에 첨부되지 않으므로 결정적인 방식으로 생성되어야 합니다.
예를 들어 키로 사용할 startCell 필드를 가진 새 EnrichedRide 클래스를 만들어 다음처럼 사용하는 대신,
keyBy(enrichedRide -> enrichedRide.startCell);
이렇게 할 수 있습니다.
keyBy(ride -> GeoUtils.mapToGridCell(ride.startLon, ride.startLat));
키가 있는 스트림에서의 집계 (Aggregations on Keyed Streams)
이 코드 조각은 각 end-of-ride 이벤트에 대해 startCell과 지속 시간(분)을 포함하는 새 튜플 스트림을 만듭니다.
import org.joda.time.Interval;
DataStream<Tuple2<Integer, Minutes>> minutesByStartCell = enrichedNYCRides
.flatMap(new FlatMapFunction<EnrichedRide, Tuple2<Integer, Minutes>>() {
@Override
public void flatMap(EnrichedRide ride,
Collector<Tuple2<Integer, Minutes>> out) throws Exception {
if (!ride.isStart) {
Interval rideInterval = new Interval(ride.startTime, ride.endTime);
Minutes duration = rideInterval.toDuration().toStandardMinutes();
out.collect(new Tuple2<>(ride.startCell, duration));
}
}
});
이제 각 startCell에 대해 (그 시점까지) 지금까지 본 것 중 가장 긴 ride만 포함하는 스트림을 만들 수 있습니다.
키로 사용할 필드를 표현하는 방법은 다양합니다. 앞서 키로 사용할 필드를 이름으로 지정한 EnrichedRide POJO의 예를 보았습니다. 이 경우는 Tuple2 객체가 관련되며, 튜플 내 인덱스(0부터 시작)가 키를 지정하는 데 사용됩니다.
minutesByStartCell
.keyBy(value -> value.f0) // .keyBy(value -> value.startCell)
.maxBy(1) // duration
.print();
이제 출력 스트림은 지속 시간이 새 최대값에 도달할 때마다 각 키에 대한 레코드를 포함합니다 — 여기서는 셀 50797로 표시됩니다.
...
4> (64549,5M)
4> (46298,18M)
1> (51549,14M)
1> (53043,13M)
1> (56031,22M)
1> (50797,6M)
...
1> (50797,8M)
...
1> (50797,11M)
...
1> (50797,12M)
(암시적) 상태 (Implicit State)
이것은 이 교육에서 상태를 수반하는 첫 번째 예시입니다. 상태가 투명하게 처리되고 있지만, Flink는 각각의 고유 키에 대해 최대 지속 시간을 추적해야 합니다.
애플리케이션에 상태가 개입될 때마다 상태가 얼마나 커질 수 있는지 생각해야 합니다. 키 공간이 무한할 때마다 Flink가 필요로 하는 상태의 양도 무한해집니다.
스트림을 다룰 때는 전체 스트림에 대한 집계보다는 유한한 창(window)에 대한 집계로 생각하는 것이 일반적으로 더 합리적입니다.
reduce()와 다른 집계자 (reduce() and other aggregators)
위에서 사용한 maxBy()는 Flink의 KeyedStream에서 사용 가능한 여러 집계 함수 중 한 예시일 뿐입니다. 또한 자신만의 사용자 정의 집계를 구현하는 데 사용할 수 있는 더 일반적인 목적의 reduce() 함수도 있습니다.
상태 있는 변환 (Stateful Transformations)
왜 Flink가 상태 관리를 담당하는가?
애플리케이션은 Flink가 관리를 담당하지 않고도 상태를 사용할 수 있습니다. 하지만 Flink는 관리하는 상태에 대해 몇 가지 매력적인 기능을 제공합니다.
- local: Flink 상태는 그것을 처리하는 머신에 로컬로 유지되며 메모리 속도로 접근할 수 있습니다.
- durable: Flink 상태는 장애 허용적입니다. 즉 정기적으로 자동 체크포인트되며 실패 시 복원됩니다.
- vertically scalable: Flink 상태는 더 많은 로컬 디스크를 추가해 확장되는 임베디드 RocksDB 인스턴스에 보관될 수 있습니다.
- horizontally scalable: Flink 상태는 클러스터가 커지고 줄어들면서 재분배됩니다.
이 섹션에서는 Flink가 관리하는 keyed state를 다루는 Flink API를 사용하는 방법을 배웁니다.
Rich 함수 (Rich Functions)
이 시점에서 FilterFunction, MapFunction, FlatMapFunction을 포함한 몇 가지 Flink 함수 인터페이스를 이미 보았습니다. 이들은 모두 Single Abstract Method 패턴의 예시입니다.
이 각각의 인터페이스에 대해 Flink는 소위 "rich" 변종도 제공합니다. 예를 들어 RichFlatMapFunction은 다음과 같은 몇 가지 추가 메서드가 있습니다.
open(OpenContext context)close()getRuntimeContext()
open()은 오퍼레이터 초기화 중에 한 번 호출됩니다. 이는 정적 데이터를 로드하거나 외부 서비스에 대한 연결을 여는 기회입니다.
getRuntimeContext()는 잠재적으로 흥미로운 것들의 전체 모음에 대한 접근을 제공하지만, 가장 두드러지게 Flink가 관리하는 상태를 만들고 접근하는 방법을 제공합니다.
Keyed State를 사용한 예시 (An Example with Keyed State)
이 예시에서는 중복을 제거하고 싶은 이벤트 스트림이 있어 각 키의 첫 번째 이벤트만 유지하고 싶다고 가정합니다. Deduplicator라는 RichFlatMapFunction을 사용해 이를 수행하는 애플리케이션은 다음과 같습니다.
private static class Event {
public final String key;
public final long timestamp;
...
}
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.addSource(new EventSource())
.keyBy(e -> e.key)
.flatMap(new Deduplicator())
.print();
env.execute();
}
이를 위해 Deduplicator는 각 키에 대해 해당 키에 대한 이벤트가 이미 있었는지 여부를 어떻게든 기억해야 합니다. Flink의 keyed state 인터페이스를 사용해 이를 수행합니다.
이런 키가 있는 스트림으로 작업할 때 Flink는 관리되는 각 상태 항목에 대해 키/값 저장소를 유지합니다.
Flink는 여러 유형의 keyed state를 지원하며, 이 예시는 가장 단순한 것인 ValueState를 사용합니다. 이는 각 키에 대해 Flink가 단일 객체를 저장한다는 뜻입니다 — 이 경우에는 Boolean 타입의 객체입니다.
우리의 Deduplicator 클래스에는 open()과 flatMap() 두 가지 메서드가 있습니다. open 메서드는 ValueStateDescriptor<Boolean>을 정의해 관리 상태의 사용을 설정합니다. 생성자에 대한 인자들은 이 keyed state 항목의 이름("keyHasBeenSeen")을 지정하고, 이 객체들을 직렬화하는 데 사용할 수 있는 정보(이 경우 Types.BOOLEAN)를 제공합니다.
public static class Deduplicator extends RichFlatMapFunction<Event, Event> {
ValueState<Boolean> keyHasBeenSeen;
@Override
public void open(OpenContext ctx) {
ValueStateDescriptor<Boolean> desc = new ValueStateDescriptor<>("keyHasBeenSeen", Types.BOOLEAN);
keyHasBeenSeen = getRuntimeContext().getState(desc);
}
@Override
public void flatMap(Event event, Collector<Event> out) throws Exception {
if (keyHasBeenSeen.value() == null) {
out.collect(event);
keyHasBeenSeen.update(true);
}
}
}
flatMap 메서드가 keyHasBeenSeen.value()를 호출하면 Flink 런타임은 문맥상의 키에 대한 이 상태 조각의 값을 조회하고, 그것이 null일 때만 이벤트를 출력으로 수집합니다. 이 경우 keyHasBeenSeen도 true로 갱신합니다.
키가 우리의 Deduplicator 구현에서 명시적으로 보이지 않기 때문에 키 파티셔닝된 상태에 접근하고 갱신하는 이 메커니즘은 다소 마술처럼 보일 수 있습니다. Flink 런타임이 RichFlatMapFunction의 open 메서드를 호출할 때는 이벤트가 없으므로 그 순간 문맥상의 키도 없습니다. 하지만 flatMap 메서드를 호출할 때는 처리 중인 이벤트의 키가 런타임에 제공되며, Flink의 상태 백엔드에서 어떤 항목을 대상으로 하는지 결정하는 데 백그라운드에서 사용됩니다.
분산 클러스터에 배포되면 이 Deduplicator의 인스턴스가 여러 개 있을 것이며, 각각은 전체 키 공간의 서로소 부분집합(disjoint subset)을 담당합니다. 따라서 다음처럼 단일 ValueState 항목을 볼 때,
ValueState<Boolean> keyHasBeenSeen;
이것이 단일 Boolean을 나타내는 것이 아니라 분산되고 샤딩된 키/값 저장소를 나타낸다는 것을 이해하세요.
상태 정리 (Clearing State)
위 예시에는 잠재적인 문제가 있습니다. 키 공간이 무한하면 어떻게 될까요? Flink는 사용되는 각각의 고유 키에 대해 어딘가에 Boolean 인스턴스를 저장합니다. 키 집합이 유한하면 괜찮지만, 키 집합이 무한한 방식으로 성장하는 애플리케이션에서는 더 이상 필요하지 않은 키의 상태를 정리하는 것이 필요합니다. 이는 다음처럼 상태 객체에 대해 clear()를 호출해 수행합니다.
keyHasBeenSeen.clear();
예를 들어 특정 키에 대한 비활성 기간 후에 이 작업을 수행하고 싶을 수 있습니다. 이벤트 기반 애플리케이션 섹션에서 ProcessFunction을 배울 때 Timers를 사용해 이를 수행하는 방법을 볼 것입니다.
또한 상태 디스크립터로 구성할 수 있는 State Time-to-Live (TTL) 옵션도 있으며, 이는 낡은 키의 상태를 자동으로 정리하고 싶은 시점을 지정합니다.
키가 없는 상태 (Non-keyed State)
키가 없는(non-keyed) 문맥에서도 관리 상태로 작업하는 것이 가능합니다. 이를 operator state라고도 합니다. 관련된 인터페이스는 다소 다르며, 사용자 정의 함수가 키가 없는 상태를 필요로 하는 것은 드물기 때문에 여기서는 다루지 않습니다. 이 기능은 주로 소스와 싱크의 구현에 가장 자주 사용됩니다.
연결된 스트림 (Connected Streams)
때로는 미리 정의된 변환을 적용하는 대신,
임계값, 규칙 또는 다른 파라미터를 스트리밍함으로써 변환의 일부 측면을 동적으로 변경할 수 있기를 원할 수 있습니다. Flink에서 이를 지원하는 패턴은 *연결된 스트림(connected streams)*이라고 하며, 단일 오퍼레이터가 두 개의 입력 스트림을 갖습니다.
연결된 스트림은 스트리밍 조인을 구현하는 데에도 사용할 수 있습니다.
예시 (Example)
이 예시에서는 streamOfWords에서 걸러내야 할 단어를 지정하기 위해 제어 스트림을 사용합니다. 연결된 스트림에 적용되어 이를 수행하는 ControlFunction이라는 RichCoFlatMapFunction이 사용됩니다.
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<String> control = env
.fromData("DROP", "IGNORE")
.keyBy(x -> x);
DataStream<String> streamOfWords = env
.fromData("Apache", "DROP", "Flink", "IGNORE")
.keyBy(x -> x);
control
.connect(streamOfWords)
.flatMap(new ControlFunction())
.print();
env.execute();
}
연결되는 두 스트림은 호환 가능한 방식으로 키가 지정되어야 합니다. keyBy의 역할은 스트림의 데이터를 파티셔닝하는 것이며, 키가 있는 스트림이 연결될 때는 동일한 방식으로 파티셔닝되어야 합니다. 이는 두 스트림 모두에서 같은 키를 가진 모든 이벤트가 같은 인스턴스로 전송되도록 보장합니다. 이를 통해 예를 들어 두 스트림을 그 키로 조인하는 것이 가능해집니다.
이 경우 두 스트림 모두 DataStream<String> 타입이고, 두 스트림 모두 문자열로 키가 지정됩니다. 아래에서 볼 수 있듯이 이 RichCoFlatMapFunction은 keyed state에 Boolean 값을 저장하며, 이 Boolean은 두 스트림이 공유합니다.
public static class ControlFunction extends RichCoFlatMapFunction<String, String, String> {
private ValueState<Boolean> blocked;
@Override
public void open(OpenContext ctx) {
blocked = getRuntimeContext()
.getState(new ValueStateDescriptor<>("blocked", Boolean.class));
}
@Override
public void flatMap1(String control_value, Collector<String> out) throws Exception {
blocked.update(Boolean.TRUE);
}
@Override
public void flatMap2(String data_value, Collector<String> out) throws Exception {
if (blocked.value() == null) {
out.collect(data_value);
}
}
}
RichCoFlatMapFunction은 연결된 한 쌍의 스트림에 적용할 수 있고 rich function 인터페이스에 접근할 수 있는 일종의 FlatMapFunction입니다. 이는 상태를 가질 수 있게 해줍니다.
blocked Boolean은 control 스트림에서 언급된 키(이 경우 단어)를 기억하는 데 사용되며, 해당 단어는 streamOfWords 스트림에서 걸러집니다. 이것은 keyed state이며 두 스트림 사이에서 공유되므로 두 스트림이 같은 키 공간을 공유해야 하는 이유입니다.
flatMap1과 flatMap2는 두 연결된 스트림 각각의 요소로 Flink 런타임이 호출합니다 — 우리의 경우 control 스트림의 요소는 flatMap1로, streamOfWords의 요소는 flatMap2로 전달됩니다. 이는 두 스트림이 control.connect(streamOfWords)로 연결된 순서에 의해 결정되었습니다.
flatMap1과 flatMap2 콜백이 호출되는 순서를 제어할 수 없다는 것을 인식하는 것이 중요합니다. 이 두 입력 스트림은 서로 경쟁하고 있으며, Flink 런타임은 한 스트림 또는 다른 스트림의 이벤트를 소비하는 것에 대해 원하는 대로 할 것입니다. 타이밍 및/또는 순서가 중요하다면, 애플리케이션이 처리할 준비가 될 때까지 관리되는 Flink 상태에 이벤트를 버퍼링하는 것이 필요할 수 있습니다. (참고: 정말 절박하다면 InputSelectable을 구현하는 사용자 정의 Operator를 사용해 두 입력 오퍼레이터가 입력을 소비하는 순서를 제한적으로 제어할 수 있습니다.)
직접 해보기 (Hands-on)
이 섹션과 함께 진행되는 실습은 Rides and Fares입니다.