투기적 실행
투기적 실행 (Speculative Execution)
이 페이지는 투기적 실행(speculative execution)의 배경, 사용 방법, 그리고 그 효과를 확인하는 방법을 설명해요. 문제가 있는 노드로 인한 batch job의 느려짐을 완화하는 메커니즘이에요.
출처: 문서
본문
배경 (Background)
투기적 실행은 문제가 있는 노드(problematic node)로 인한 job의 느려짐을 완화하는 메커니즘이에요. 문제가 있는 노드는 하드웨어 문제, 우연한 I/O 병목, 높은 CPU 부하가 있을 수 있어요. 이러한 문제는 호스팅된 task를 다른 노드의 task보다 훨씬 느리게 실행하게 하고, batch job의 전체 실행 시간에 영향을 줄 수 있어요.
이러한 경우 투기적 실행은 문제가 있다고 감지되지 않은 노드에서 느린 task의 새 attempt를 시작해요. 새 attempt는 같은 입력 데이터를 처리하고 이전 것과 같은 데이터를 생성해요. 이전 attempt는 영향을 받지 않고 계속 실행돼요. 먼저 완료된 attempt가 채택되고, 그 출력이 다운스트림 task에 보이고 소비되며, 나머지 attempt는 취소돼요.
이를 달성하기 위해 Flink는 느린 task 감지기(slow task detector)를 사용해 느린 task를 감지해요. 느린 task가 위치한 노드는 문제가 있는 노드로 식별되고 blocklist 메커니즘을 통해 차단돼요. 스케줄러는 느린 task를 위한 새 attempt를 만들고 이를 차단되지 않은 노드에 배포해요.
사용법 (Usage)
이 절에서는 투기적 실행을 사용하는 방법을 설명해요. 활성화 방법, 튜닝 방법, 그리고 custom source를 투기적 실행과 함께 동작하도록 개발/개선하는 방법을 포함해요.
투기적 실행 활성화 (Enable Speculative Execution)
다음 구성 항목으로 투기적 실행을 활성화할 수 있어요.
execution.batch.speculative.enabled: true
참고: 현재 Adaptive Batch Scheduler만 투기적 실행을 지원해요. 그리고 Flink batch job은 다른 스케줄러가 명시적으로 구성되지 않는 한 기본적으로 이 스케줄러를 사용해요.
튜닝 구성 (Tuning Configuration)
투기적 실행이 서로 다른 job에서 더 잘 동작하도록, 다음 스케줄러 구성 옵션을 튜닝할 수 있어요.
execution.batch.speculative.max-concurrent-executionsexecution.batch.speculative.block-slow-node-duration
또한 느린 task 감지기의 다음 구성 옵션을 튜닝할 수 있어요.
slow-task-detector.check-intervalslow-task-detector.execution-time.baseline-lower-boundslow-task-detector.execution-time.baseline-multiplierslow-task-detector.execution-time.baseline-ratio
현재 투기적 실행은 실행 시간 기반의 느린 task 감지기를 사용해 느린 task를 감지해요. 감지기는 주기적으로 모든 완료된 실행을 집계하며, 완료된 실행 비율이 구성된 비율(slow-task-detector.execution-time.baseline-ratio)에 도달하면 기준선(baseline)은 실행 시간 중앙값에 구성된 배수(slow-task-detector.execution-time.baseline-multiplier)를 곱한 값으로 정의돼요. 그런 다음 실행 시간이 기준선을 초과하는 실행 중인 task가 느린 task로 감지돼요.
실행 시간은 실행 정점(execution vertex)의 입력 데이터 볼륨으로 가중치가 부여된다는 점을 언급할 가치가 있어요. 따라서 데이터 스큐가 발생할 때 데이터 볼륨 차이는 크지만 계산 성능은 비슷한 실행은 느린 task로 감지되지 않아요. 이는 불필요한 투기적 attempt 시작을 피하는 데 도움이 돼요.
참고: 노드가 Source이거나 Hybrid Shuffle 모드를 사용하면, 입력 데이터 볼륨을 판단할 수 없으므로 실행 시간에 입력 데이터 볼륨으로 가중치를 부여하는 최적화가 적용되지 않아요.
투기적 실행을 위한 Source 활성화 (Enable Sources for Speculative Execution)
job이 custom Source를 사용하고, 소스가 custom SourceEvent를 사용한다면, 그 소스의 SplitEnumerator를 변경해 SupportsHandleExecutionAttemptSourceEvent 인터페이스를 구현해야 해요.
public interface SupportsHandleExecutionAttemptSourceEvent {
void handleSourceEvent(int subtaskId, int attemptNumber, SourceEvent sourceEvent);
}
즉 SplitEnumerator가 이벤트를 보내는 attempt를 인식해야 해요. 그렇지 않으면 job manager가 task에서 source event를 받을 때 예외가 발생하고 job 실패로 이어져요.
다른 소스는 투기적 실행과 함께 동작하기 위해 추가 변경이 필요 없어요. 여기에는 SourceFunction 소스, InputFormat 소스, 그리고 new source가 포함돼요. Apache Flink가 제공하는 모든 소스 커넥터는 투기적 실행과 함께 동작할 수 있어요.
투기적 실행을 위한 Sink 활성화 (Enable Sinks for Speculative Execution)
투기적 실행은 SupportsConcurrentExecutionAttempts 인터페이스를 구현하지 않는 한 sink에 대해 기본적으로 비활성화돼요. 이는 호환성 고려 때문이에요.
public interface SupportsConcurrentExecutionAttempts {}
SupportsConcurrentExecutionAttempts는 Sink, SinkFunction, OutputFormat에 대해 동작해요.
task의 어떤 연산자라도 투기적 실행을 지원하지 않으면 전체 task가 투기적 실행을 지원하지 않는 것으로 표시돼요. 즉 Sink가 투기적 실행을 지원하지 않으면, Sink 연산자를 포함하는 task는 투기적으로 실행될 수 없어요.
Sink 구현의 경우, Flink는 Committer(WithPreCommitTopology와 WithPostCommitTopology로 확장된 연산자 포함)에 대해 투기적 실행을 비활성화해요. 동시 커밋은 사용자가 경험하지 않았다면 예상치 못한 문제를 일으킬 수 있기 때문이에요. 그리고 committer가 batch job의 병목이 될 가능성은 매우 낮아요.
투기적 실행의 효과 확인하기 (Checking the Effectiveness of Speculative Execution)
투기적 실행을 활성화한 뒤 느린 task가 투기적 실행을 유발하면, web UI가 job 페이지의 정점(vertices) SubTasks 탭에 투기적 attempt를 표시해요. web UI는 또한 Flink 클러스터 Overview와 Task Managers 페이지에서 차단된 taskmanager도 표시해요. 이러한 메트릭을 확인해 투기적 실행의 효과를 볼 수도 있어요.