실행 모드

실행 모드 (Execution Mode, Batch/Streaming)

DataStream API는 사용 사례의 요구사항과 작업의 특성에 따라 선택할 수 있는 여러 런타임 실행 모드를 지원합니다. 여기에는 '고전적인' DataStream 실행 동작인 STREAMING 실행 모드와, MapReduce 같은 배치 처리 프레임워크를 연상시키는 BATCH 실행 모드가 있습니다. STREAMING 모드는 연속적인 증분 처리를 요구하고 무기한 온라인 상태를 유지해야 하는 무한(unbounded) 작업에 사용하며, BATCH 모드는 입력이 고정되고 연속 실행되지 않는 유한(bounded) 작업에 사용합니다. Apache Flink의 스트림/배치 통합 방식 덕분에 유한 입력으로 실행한 DataStream 애플리케이션은 실행 모드와 무관하게 동일한 최종 결과를 만들어냅니다.

출처: 문서

본문

소개 (Introduction)

DataStream API에는 STREAMING 실행 모드와 BATCH 실행 모드가 있습니다. STREAMING 모드는 연속적인 증분 처리가 필요한 무한 작업에 사용하고, BATCH 모드는 입력이 알려진 유한 작업에 사용합니다.

Apache Flink의 스트림/배치 통합 접근 방식은 유한 입력으로 실행된 DataStream 애플리케이션이 실행 모드와 무관하게 동일한 최종 결과를 생성한다는 것을 의미합니다. 여기서 최종이 무엇을 의미하는지 주의할 필요가 있습니다. STREAMING 모드로 실행되는 작업은 증분 업데이트(데이터베이스의 upsert를 떠올려 보세요)를 생성할 수 있는 반면, BATCH 작업은 마지막에 최종 결과 하나만 생성합니다. 올바르게 해석한다면 최종 결과는 동일하지만 그 결과에 도달하는 방식은 다를 수 있습니다.

BATCH 실행을 활성화하면 Flink는 입력이 유한하다는 것을 알 때만 적용할 수 있는 추가 최적화를 적용할 수 있습니다. 예를 들어 서로 다른 join/aggregation 전략, 그리고 더 효율적인 작업 스케줄링과 실패 복구 동작을 가능하게 하는 서로 다른 셔플 구현을 사용할 수 있습니다.

BATCH 실행 모드는 언제/언제 사용해야 하나요?

BATCH 실행 모드는 *유한(bounded)*한 작업/프로그램에만 사용할 수 있습니다. 유한성(boundedness)은 데이터 소스의 속성으로, 해당 소스에서 오는 모든 입력이 실행 전에 알려져 있는지 아니면 새로운 데이터가 (잠재적으로 무기한) 계속 나타날지를 알려줍니다. 작업은 모든 소스가 유한하면 유한하고, 그렇지 않으면 무한합니다.

반면 STREAMING 실행 모드는 유한 작업과 무한 작업 모두에 사용할 수 있습니다.

일반적인 규칙으로, 프로그램이 유한하면 BATCH 실행 모드를 사용해야 합니다. 더 효율적이기 때문입니다. 프로그램이 무한하면 STREAMING 실행 모드를 사용해야 합니다. 연속 데이터 스트림을 다룰 수 있을 만큼 일반적인 모드는 이것뿐이기 때문입니다.

한 가지 명백한 예외는 무한 작업에서 사용할 작업 상태를 부트스트랩하기 위해 유한 작업을 사용하고 싶은 경우입니다. 예를 들어 STREAMING 모드로 유한 작업을 실행하고 savepoint를 만든 다음, 그 savepoint를 무한 작업에서 복원하는 경우입니다. 이는 매우 특수한 사용 사례로, BATCH 실행 작업의 추가 출력으로 savepoint를 생성할 수 있게 되면 곧 쓸모없어질 수도 있습니다.

STREAMING 모드로 유한 작업을 실행할 수 있는 또 다른 경우는 결국 무한 소스로 실행될 코드를 테스트할 때입니다. 테스트에는 유한 소스를 사용하는 것이 더 자연스러울 수 있습니다.

BATCH 실행 모드 구성

실행 모드는 execution.runtime-mode 설정으로 구성할 수 있습니다. 세 가지 값이 있습니다.

  • STREAMING: 고전적인 DataStream 실행 모드 (기본값)
  • BATCH: DataStream API에서의 배치 스타일 실행
  • AUTOMATIC: 소스의 유한성에 따라 시스템이 결정

이는 bin/flink run ...의 명령줄 파라미터로 구성하거나, StreamExecutionEnvironment를 만들거나 구성할 때 프로그래밍 방식으로 구성할 수 있습니다.

명령줄로 실행 모드를 구성하는 방법은 다음과 같습니다.

$ bin/flink run -Dexecution.runtime-mode=BATCH <jarFile>

코드에서 실행 모드를 구성하는 예시입니다.

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setRuntimeMode(RuntimeExecutionMode.BATCH);

프로그램에서 런타임 모드를 설정하지 말고, 권장사항으로 애플리케이션을 제출할 때 명령줄에서 설정할 것을 권장합니다. 애플리케이션 코드를 구성 없이 유지하면 동일한 애플리케이션을 어떤 실행 모드로도 실행할 수 있어 더 유연합니다.

실행 동작 (Execution Behavior)

이 섹션은 BATCH 실행 모드의 실행 동작에 대한 개요를 제공하고 STREAMING 실행 모드와 대조합니다. 자세한 내용은 이 기능을 소개한 FLIP 문서인 FLIP-134FLIP-140를 참고하세요.

작업 스케줄링과 네트워크 셔플 (Task Scheduling And Network Shuffle)

Flink 작업은 데이터플로우 그래프로 연결된 서로 다른 연산들로 구성됩니다. 시스템은 이러한 연산들을 서로 다른 프로세스/머신(TaskManager)에서 어떻게 실행할지, 그리고 그 사이에서 데이터를 어떻게 셔플(전송)할지 결정합니다.

여러 연산/오퍼레이터는 체이닝(chaining)이라는 기능을 사용해 서로 연결될 수 있습니다. Flink가 스케줄링 단위로 간주하는 하나 이상의 (체인으로 연결된) 오퍼레이터 그룹을 *작업(task)*이라고 합니다. 여러 TaskManager에서 병렬로 실행되는 개별 작업 인스턴스를 지칭할 때 흔히 *서브태스크(subtask)*라는 용어를 사용하지만, 여기서는 *작업(task)*이라는 용어만 사용하겠습니다.

작업 스케줄링과 네트워크 셔플은 BATCHSTREAMING 실행 모드에서 다르게 동작합니다. 주된 이유는 BATCH 실행 모드에서는 입력 데이터가 유한하다는 것을 알 수 있어 Flink가 더 효율적인 데이터 구조와 알고리즘을 사용할 수 있기 때문입니다.

작업 스케줄링과 네트워크 전송의 차이를 설명하기 위해 다음 예시를 사용합니다.

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

DataStreamSource<String> source = env.fromData(...);

source.name("source")
	.map(...).name("map1")
	.map(...).name("map2")
	.rebalance()
	.map(...).name("map3")
	.map(...).name("map4")
	.keyBy((value) -> value)
	.map(...).name("map5")
	.map(...).name("map6")
	.sinkTo(...).name("sink");

map(), flatMap(), filter()처럼 연산 간 1-to-1 연결 패턴을 의미하는 연산은 데이터를 다음 연산으로 바로 전달할 수 있어 서로 체이닝될 수 있습니다. 즉 Flink는 보통 이들 사이에 네트워크 셔플을 삽입하지 않습니다.

반면 keyBy()rebalance() 같은 연산은 작업의 서로 다른 병렬 인스턴스 사이에서 데이터를 셔플해야 하므로 네트워크 셔플이 발생합니다.

위 예시에서 Flink는 연산들을 다음과 같은 작업으로 그룹화합니다.

  • Task1: source, map1, map2
  • Task2: map3, map4
  • Task3: map5, map6, sink

그리고 Task1과 Task2 사이, Task2와 Task3 사이에 네트워크 셔플이 있습니다. 이 작업의 시각적 표현은 다음과 같습니다.

예시 작업 그래프

STREAMING 실행 모드

STREAMING 실행 모드에서는 모든 작업이 항상 온라인/실행 중이어야 합니다. 덕분에 Flink는 새 레코드를 파이프라인 전체를 통해 즉시 처리할 수 있는데, 이는 연속적이고 낮은 지연 시간의 스트림 처리를 위해 필요합니다. 또한 작업에 할당된 TaskManager가 모든 작업을 동시에 실행할 수 있는 충분한 리소스를 가져야 한다는 것을 의미합니다.

네트워크 셔플은 *파이프라인 방식(pipelined)*입니다. 즉 레코드가 네트워크 계층에서 일부 버퍼링되며 다운스트림 작업으로 즉시 전송됩니다. 연속 데이터 스트림을 처리할 때는 작업(또는 작업 파이프라인)들 사이에서 데이터를 실체화(materialize)할 수 있는 자연스러운 (시간상의) 지점이 없기 때문에 이것이 필요합니다. 이는 아래에서 설명하듯 중간 결과를 실체화할 수 있는 BATCH 실행 모드와 대조됩니다.

BATCH 실행 모드

BATCH 실행 모드에서는 작업의 태스크들이 하나씩 실행할 수 있는 스테이지로 분리될 수 있습니다. 입력이 유한하므로 Flink가 파이프라인의 한 스테이지를 완전히 처리한 다음 다음으로 넘어갈 수 있기 때문입니다. 위 예시에서 작업은 셔플 경계로 분리된 세 개의 작업에 해당하는 세 개의 스테이지를 갖습니다.

레코드를 즉시 다운스트림 작업으로 보내는 대신(STREAMING 모드에서 설명한 대로), 스테이지 처리는 작업의 중간 결과를 어떤 비휘발성 저장소에 실체화해야 합니다. 그래야 업스트림 작업이 오프라인이 된 후에도 다운스트림 작업이 그 결과를 읽을 수 있습니다. 이는 처리 지연을 증가시키지만 다른 흥미로운 속성을 가져옵니다. 첫째, 실패가 발생했을 때 전체 작업을 다시 시작하는 대신 가장 최근에 사용 가능한 결과로 되돌아갈 수 있습니다(backtrack). 둘째, BATCH 작업은 (TaskManager의 가용 슬롯 측면에서) 더 적은 리소스로 실행될 수 있습니다. 시스템이 작업을 순차적으로 하나씩 실행할 수 있기 때문입니다.

TaskManager는 중간 결과를 최소한 다운스트림 작업이 소비할 때까지 유지합니다. (엄밀히 말하면, 소비하는 파이프라인 영역이 그 출력을 생성할 때까지 유지됩니다.) 그 후에는 실패 시 앞서 언급한 이전 결과로의 backtracking을 허용하기 위해 공간이 허용하는 한 유지됩니다.

상태 백엔드 / 상태 (State Backends / State)

STREAMING 모드에서는 Flink가 StateBackend를 사용해 상태가 어떻게 저장되고 체크포인팅이 어떻게 동작하는지 제어합니다.

BATCH 모드에서는 구성된 상태 백엔드가 무시됩니다. 대신 키가 있는 연산의 입력이 키별로 그룹화되고(정렬 사용), 한 키의 모든 레코드를 차례로 처리합니다. 이렇게 하면 한 번에 하나의 키에 대해서만 상태를 유지할 수 있습니다. 특정 키의 상태는 다음 키로 이동할 때 버려집니다.

자세한 배경 정보는 FLIP-140을 참고하세요.

처리 순서 (Order of Processing)

레코드가 오퍼레이터나 사용자 정의 함수(UDF)에서 처리되는 순서는 BATCHSTREAMING 실행 간에 다를 수 있습니다.

STREAMING 모드에서는 사용자 정의 함수가 들어오는 레코드의 순서에 대해 어떤 가정도 해서는 안 됩니다. 데이터는 도착하는 즉시 처리됩니다.

BATCH 실행 모드에서는 Flink가 순서를 보장하는 몇 가지 연산이 있습니다. 순서는 특정 작업 스케줄링, 네트워크 셔플, 상태 백엔드(위 참조)의 부수 효과이거나 시스템의 의식적인 선택일 수 있습니다.

구분할 수 있는 세 가지 일반적인 입력 유형이 있습니다.

  • broadcast input: broadcast 스트림에서 오는 입력 (Broadcast State 참조)
  • regular input: broadcast도 keyed도 아닌 입력
  • keyed input: KeyedStream에서 오는 입력

여러 입력 유형을 소비하는 함수나 오퍼레이터는 다음 순서로 처리합니다.

  • broadcast 입력이 먼저 처리됩니다
  • regular 입력이 두 번째로 처리됩니다
  • keyed 입력이 마지막으로 처리됩니다

여러 regular 또는 broadcast 입력을 소비하는 함수(예: CoProcessFunction)의 경우 Flink는 해당 유형의 어떤 입력에서든 임의의 순서로 데이터를 처리할 권리가 있습니다.

여러 keyed 입력을 소비하는 함수(예: KeyedCoProcessFunction)의 경우 Flink는 다음 키로 넘어가기 전에 모든 keyed 입력에서 단일 키의 모든 레코드를 처리합니다.

이벤트 시간 / 워터마크 (Event Time / Watermarks)

이벤트 시간을 지원할 때, Flink의 스트리밍 런타임은 이벤트가 순서 없이 올 수 있다는 비관적 가정에 기반합니다. 즉 타임스탬프 t를 가진 이벤트가 타임스탬프 t+1을 가진 이벤트 이후에 올 수 있습니다. 이 때문에 시스템은 특정 타임스탬프 T에 대해 타임스탬프가 t < T인 요소가 미래에 더 이상 오지 않는다고 확신할 수 없습니다. STREAMING 모드에서 Flink는 시스템을 실용적으로 만들면서 이 순서 뒤섞임이 최종 결과에 미치는 영향을 완화하기 위해 Watermark라는 휴리스틱을 사용합니다. 타임스탬프 T를 가진 워터마크는 타임스탬프가 t < T인 요소가 더 이상 오지 않음을 알립니다.

BATCH 모드에서는 입력 데이터셋이 미리 알려져 있으므로 이런 휴리스틱이 필요 없습니다. 최소한 요소를 타임스탬프로 정렬해 시간 순서대로 처리할 수 있기 때문입니다. 스트리밍에 익숙한 독자라면, BATCH에서는 "완벽한 워터마크(perfect watermarks)"를 가정할 수 있습니다.

위 사실을 바탕으로 BATCH 모드에서는 각 키와 연관된 입력의 끝, 또는 입력 스트림이 keyed가 아닌 경우 입력의 끝에서 MAX_WATERMARK만 필요합니다. 이 방식에 따라 등록된 모든 타이머는 *시간의 끝(end of time)*에 발화하며, 사용자 정의 WatermarkAssignerWatermarkGenerator는 무시됩니다. 하지만 WatermarkStrategy를 지정하는 것은 여전히 중요합니다. 그 이유는 WatermarkStrategyTimestampAssigner가 레코드에 타임스탬프를 할당하는 데 여전히 사용되기 때문입니다.

처리 시간 (Processing Time)

Processing Time은 레코드가 처리되는 특정 순간에 레코드가 처리되는 머신의 벽시계 시간(wall-clock time)입니다. 이 정의에 따르면, 처리 시간에 기반한 계산 결과는 재현 가능하지 않습니다. 같은 레코드를 두 번 처리하면 두 개의 서로 다른 타임스탬프를 갖게 되기 때문입니다.

그럼에도 불구하고 STREAMING 모드에서 처리 시간을 사용하는 것은 유용할 수 있습니다. 그 이유는 스트리밍 파이프라인이 종종 무한 입력을 실시간으로 수집하므로 이벤트 시간과 처리 시간 사이에 상관관계가 있기 때문입니다. 또한 위 사실 때문에 STREAMING 모드에서 이벤트 시간의 1h는 종종 처리 시간(벽시계 시간)의 거의 1h가 될 수 있습니다. 그래서 처리 시간은 예상 결과에 대한 힌트를 주는 조기(불완전한) 발화에 사용할 수 있습니다.

입력 데이터셋이 정적이고 미리 알려져 있는 배치 세계에서는 이런 상관관계가 존재하지 않습니다. 따라서 BATCH 모드에서는 사용자가 현재 처리 시간을 요청하고 처리 시간 타이머를 등록하는 것을 허용하지만, 이벤트 시간의 경우와 마찬가지로 모든 타이머는 입력의 끝에서 발화합니다.

개념적으로, 처리 시간은 작업 실행 중에는 앞으로 나아가지 않고 전체 입력이 처리될 때 시간의 끝으로 빨리 감는다고 상상할 수 있습니다.

실패 복구 (Failure Recovery)

STREAMING 실행 모드에서 Flink는 실패 복구에 체크포인트를 사용합니다. 이에 대한 실전 문서와 구성 방법은 체크포인팅 문서를 참고하세요. 개념을 더 높은 수준에서 설명하는 상태 스냅샷을 통한 장애 허용 소개 섹션도 있습니다.

실패 복구를 위한 체크포인팅의 특징 중 하나는 실패가 발생하면 Flink가 모든 실행 중인 작업을 체크포인트에서 다시 시작한다는 것입니다. 이는 BATCH 모드에서 해야 하는 것(아래 설명)보다 더 비쌀 수 있으며, 이것이 작업이 허용한다면 BATCH 실행 모드를 사용해야 하는 이유 중 하나입니다.

BATCH 실행 모드에서 Flink는 중간 결과가 아직 사용 가능한 이전 처리 스테이지로 되돌아가려고 시도합니다. 잠재적으로 실패한 작업(또는 그래프에서 그 작업의 선행자)만 다시 시작하면 되므로, 체크포인트에서 모든 작업을 다시 시작하는 것보다 처리 효율성과 작업의 전체 처리 시간을 개선할 수 있습니다.

중요한 고려사항 (Important Considerations)

고전적인 STREAMING 실행 모드와 비교할 때 BATCH 모드에서는 일부 기능이 기대한 대로 동작하지 않을 수 있습니다. 일부 기능은 약간 다르게 동작하고, 다른 기능은 지원되지 않습니다.

BATCH 모드에서 동작이 변경되는 것:

  • reduce()sum() 같은 "롤링(rolling)" 연산은 STREAMING 모드에서 도착하는 모든 새 레코드에 대해 증분 업데이트를 내보냅니다. BATCH 모드에서는 이러한 연산이 "롤링"이 아닙니다. 최종 결과만 내보냅니다.

BATCH 모드에서 지원되지 않는 것:

사용자 정의 오퍼레이터는 주의해서 구현해야 합니다. 그렇지 않으면 부적절하게 동작할 수 있습니다. 자세한 내용은 아래의 추가 설명을 참고하세요.

체크포인팅 (Checkpointing)

위에서 설명한 대로, 배치 프로그램의 실패 복구는 체크포인팅을 사용하지 않습니다.

체크포인트가 없기 때문에 CheckpointListener 같은 특정 기능과, 그 결과로 Kafka의 EXACTLY_ONCE 모드나 File SinkOnCheckpointRollingPolicy는 동작하지 않는다는 것을 기억하는 것이 중요합니다. BATCH 모드에서 동작하는 트랜잭션 싱크가 필요하다면 FLIP-143에서 제안한 Unified Sink API를 사용하는지 확인하세요.

모든 상태 프리미티브는 여전히 사용할 수 있습니다. 단지 실패 복구에 사용되는 메커니즘이 다를 뿐입니다.

사용자 정의 오퍼레이터 작성 (Writing Custom Operators)

참고: 사용자 정의 오퍼레이터는 Apache Flink의 고급 사용 패턴입니다. 대부분의 사용 사례에서는 (키가 있는) process function을 사용하는 것을 고려하세요.

사용자 정의 오퍼레이터를 작성할 때는 BATCH 실행 모드에서 가정하는 사항들을 기억하는 것이 중요합니다. 그렇지 않으면 STREAMING 모드에서는 잘 동작하는 오퍼레이터가 BATCH 모드에서는 잘못된 결과를 생성할 수 있습니다. 오퍼레이터는 특정 키에 범위가 한정되지 않으므로 Flink가 활용하려는 BATCH 처리의 어떤 속성을 볼 수 있습니다.

먼저 오퍼레이터 내부에 마지막으로 본 워터마크를 캐시하면 안 됩니다. BATCH 모드에서는 레코드를 키별로 처리합니다. 그 결과 워터마크는 각 키 사이에서 MAX_VALUE에서 MIN_VALUE로 전환됩니다. 오퍼레이터에서 워터마크가 항상 오름차순이라고 가정하면 안 됩니다. 같은 이유로 타이머는 먼저 키 순서로, 그 다음 각 키 내에서 타임스탬프 순서로 발화합니다. 또한 키를 수동으로 변경하는 연산은 지원되지 않습니다.

더 알아보기 (Learn more)