병렬 실행
병렬 실행 (Parallel Execution)
이 섹션은 Flink에서 프로그램의 병렬 실행을 구성하는 방법을 설명해요. Flink 프로그램은 여러 작업(변환/연산자, 데이터 소스, 싱크)으로 구성돼요. 작업은 실행을 위해 여러 병렬 인스턴스로 나뉘며, 각 병렬 인스턴스는 작업 입력 데이터의 일부를 처리해요. 작업의 병렬 인스턴스 수를 parallelism이라고 해요.
출처: 문서
본문
savepoints를 사용하려면 최대 병렬도(max parallelism)를 설정하는 것도 고려해야 해요. savepoint에서 복원할 때 특정 연산자나 전체 프로그램의 병렬도를 바꿀 수 있는데, 이 설정은 병렬도의 상한을 지정해요. Flink는 내부적으로 상태를 키-그룹(key-group)으로 분할하기 때문에 이 상한이 필요해요. 성능에 해가 되므로 +Inf 개수의 키-그룹을 가질 수는 없어요.
병렬도 설정하기 (Setting the Parallelism)
작업의 병렬도는 Flink에서 여러 수준으로 지정할 수 있어요.
연산자 수준 (Operator Level)
개별 연산자, 데이터 소스 또는 데이터 싱크의 병렬도는 setParallelism() 메서드를 호출해 정의할 수 있어요. 예를 들면:
Java:
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<String> text = [...];
DataStream<Tuple2<String, Integer>> wordCounts = text
.flatMap(new LineSplitter())
.keyBy(value -> value.f0)
.window(TumblingEventTimeWindows.of(Duration.ofSeconds(5)))
.sum(1).setParallelism(5);
wordCounts.print();
env.execute("Word Count Example");
Python:
env = StreamExecutionEnvironment.get_execution_environment()
text = [...]
word_counts = text
.flat_map(lambda x: x.split(" ")) \
.map(lambda i: (i, 1), output_type=Types.TUPLE([Types.STRING(), Types.INT()])) \
.key_by(lambda i: i[0]) \
.window(TumblingEventTimeWindows.of(Duration.ofSeconds(5))) \
.reduce(lambda i, j: (i[0], i[1] + j[1])) \
.set_parallelism(5)
word_counts.print()
env.execute("Word Count Example")
실행 환경 수준 (Execution Environment Level)
여기에서 언급했듯이 Flink 프로그램은 실행 환경의 맥락에서 실행돼요. 실행 환경은 실행하는 모든 연산자, 데이터 소스, 데이터 싱크의 기본 병렬도를 정의해요. 실행 환경 병렬도는 연산자의 병렬도를 명시적으로 구성해 덮어쓸 수 있어요.
실행 환경의 기본 병렬도는 setParallelism() 메서드를 호출해 지정할 수 있어요. 모든 연산자, 데이터 소스, 데이터 싱크를 병렬도 3으로 실행하려면 실행 환경의 기본 병렬도를 다음과 같이 설정해요:
Java:
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(3);
DataStream<String> text = [...];
DataStream<Tuple2<String, Integer>> wordCounts = [...];
wordCounts.print();
env.execute("Word Count Example");
Python:
env = StreamExecutionEnvironment.get_execution_environment()
env.set_parallelism(3)
text = [...]
word_counts = text
.flat_map(lambda x: x.split(" ")) \
.map(lambda i: (i, 1), output_type=Types.TUPLE([Types.STRING(), Types.INT()])) \
.key_by(lambda i: i[0]) \
.window(TumblingEventTimeWindows.of(Duration.ofSeconds(5))) \
.reduce(lambda i, j: (i[0], i[1] + j[1]))
word_counts.print()
env.execute("Word Count Example")
클라이언트 수준 (Client Level)
Flink에 작업을 제출할 때 클라이언트에서 병렬도를 설정할 수 있어요. 그러한 클라이언트의 예로 Flink의 커맨드라인 인터페이스(CLI)가 있어요. CLI 클라이언트에서는 병렬도 파라미터를 -p로 지정할 수 있어요. 예를 들어:
./bin/flink run -p 10 ../examples/*WordCount-java*.jar
클라이언트 프로그램에서 병렬도는 다음과 같이 설정해요:
Java:
try {
PackagedProgram program = new PackagedProgram(file, args);
InetSocketAddress jobManagerAddress = RemoteExecutor.getInetFromHostport("localhost:6123");
Configuration config = new Configuration();
Client client = new Client(jobManagerAddress, config, program.getUserCodeClassLoader());
// set the parallelism to 10 here
client.run(program, 10, true);
} catch (ProgramInvocationException e) {
e.printStackTrace();
}
Python:
Still not supported in Python API.
시스템 수준 (System Level)
모든 실행 환경의 시스템 전체 기본 병렬도는 Flink 구성 파일에서 parallelism.default 프로퍼티를 설정해 정의할 수 있어요. 자세한 내용은 Configuration 문서를 참고해요.
최대 병렬도 설정하기 (Setting the Maximum Parallelism)
최대 병렬도는 병렬도를 설정할 수 있는 곳(클라이언트 수준과 시스템 수준 제외)에서 설정할 수 있어요. setParallelism()을 호출하는 대신 setMaxParallelism()을 호출해 최대 병렬도를 설정해요.
최대 병렬도의 기본 설정은 대략 operatorParallelism + (operatorParallelism / 2)이며 하한 128, 상한 32768이에요. 최대 병렬도를 매우 큰 값으로 설정하면 일부 상태 백엔드가 키-그룹 수에 따라 확장되는 내부 데이터 구조를 유지해야 하므로 성능에 해가 될 수 있어요.
원래 작업에서 복구할 때 최대 병렬도를 명시적으로 바꾸면 상태 비호환성이 발생할 수 있어요.