애플리케이션 매개변수 처리
애플리케이션 매개변수 처리 (Handling Application Parameters)
거의 모든 Flink 애플리케이션(배치 및 스트리밍 모두)은 외부 구성 매개변수에 의존해요. 이 매개변수들은 입력·출력 소스(경로나 주소 같은), 시스템 매개변수(병렬도, 런타임 구성), 애플리케이션 특정 매개변수(보통 사용자 함수 내에서 사용)를 지정하는 데 사용돼요.
출처: 문서
본문
Handling Application Parameters
Flink는 이러한 문제를 해결하기 위한 기본 도구를 제공하는 ParameterTool이라는 간단한 유틸리티를 제공해요.
여기에 설명된 ParameterTool을 꼭 사용할 필요는 없다는 점에 유의하세요. Commons CLI나 argparse4j 같은 다른 프레임워크도 Flink와 잘 함께 동작해요.
구성 값을 ParameterTool에 넣기 (Getting your configuration values into the ParameterTool)
ParameterTool은 구성을 읽기 위한 미리 정의된 정적 메서드 집합을 제공해요. 이 도구는 내부적으로 Map을 기대하므로 자신만의 구성 스타일과 통합하기가 매우 쉬워요.
.properties 파일에서 (From .properties files)
다음 메서드는 Properties 파일을 읽고 key/value 쌍을 제공해요.
String propertiesFilePath = "/home/sam/flink/myjob.properties";
ParameterTool parameters = ParameterTool.fromPropertiesFile(propertiesFilePath);
File propertiesFile = new File(propertiesFilePath);
ParameterTool parameters = ParameterTool.fromPropertiesFile(propertiesFile);
InputStream propertiesFileInputStream = new FileInputStream(file);
ParameterTool parameters = ParameterTool.fromPropertiesFile(propertiesFileInputStream);
명령줄 인자에서 (From the command line arguments)
이를 통해 명령줄에서 --input hdfs:///mydata --elements 42 같은 인자를 가져올 수 있어요.
public static void main(String[] args) {
ParameterTool parameters = ParameterTool.fromArgs(args);
// .. regular code ..
시스템 속성에서 (From system properties)
JVM을 시작할 때 -Dinput=hdfs:///mydata처럼 시스템 속성을 전달할 수 있어요. 이러한 시스템 속성에서 ParameterTool을 초기화할 수도 있어요.
ParameterTool parameters = ParameterTool.fromSystemProperties();
Flink 프로그램에서 매개변수 사용하기
이제 (위에서) 어딘가에서 매개변수를 얻었으니, 다양한 방식으로 사용할 수 있어요.
ParameterTool에서 직접
ParameterTool 자체에 값에 접근하기 위한 메서드가 있어요.
ParameterTool parameters = // ...
parameters.getRequired("input");
parameters.get("output", "myDefaultValue");
parameters.getLong("expectedCount", -1L);
parameters.getNumberOfParameters();
// .. there are more methods available.
이 메서드들의 반환값을 애플리케이션을 제출하는 클라이언트의 main() 메서드에서 직접 사용할 수 있어요. 예를 들어 연산자의 parallelism을 이렇게 설정할 수 있어요.
ParameterTool parameters = ParameterTool.fromArgs(args);
int parallelism = parameters.get("mapParallelism", 2);
DataStream<Tuple2<String, Integer>> counts = text.flatMap(new Tokenizer()).setParallelism(parallelism);
ParameterTool은 직렬화 가능하므로 함수 자체에 전달할 수 있어요.
ParameterTool parameters = ParameterTool.fromArgs(args);
DataStream<Tuple2<String, Integer>> counts = text.flatMap(new Tokenizer(parameters));
그리고 함수 내부에서 명령줄에서 값을 가져오는 데 사용할 수 있어요.
매개변수를 전역으로 등록하기 (Register the parameters globally)
ExecutionConfig에서 전역 job 매개변수로 등록된 매개변수는 JobManager web 인터페이스와 사용자가 정의한 모든 함수에서 구성 값으로 접근할 수 있어요.
매개변수를 전역으로 등록해요.
ParameterTool parameters = ParameterTool.fromArgs(args);
// set up the execution environment
final ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
env.getConfig().setGlobalJobParameters(parameters);
모든 rich 사용자 함수에서 접근해요.
public static final class Tokenizer extends RichFlatMapFunction<String, Tuple2<String, Integer>> {
@Override
public void flatMap(String value, Collector<Tuple2<String, Integer>> out) {
ParameterTool parameters = ParameterTool.fromMap(getRuntimeContext().getGlobalJobParameters());
parameters.getRequired("input");
// .. do more ..