실행 구성
실행 구성 (Execution Configuration)
StreamExecutionEnvironment에는 런타임에 대한 작업 고유의 구성 값을 설정할 수 있는 ExecutionConfig가 포함되어 있습니다. 모든 작업에 영향을 주는 기본값을 바꾸려면 Configuration을 참고하세요.
출처: 문서
본문
Java
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
ExecutionConfig executionConfig = env.getConfig();
Python
env = StreamExecutionEnvironment.get_execution_environment()
execution_config = env.get_config()
다음과 같은 구성 옵션을 사용할 수 있습니다(기본값은 볼드체):
setClosureCleanerLevel(). 클로저 클리너 레벨은 기본적으로ClosureCleanerLevel.RECURSIVE로 설정됩니다. 클로저 클리너는 Flink 프로그램 내 익명 함수가 둘러싼(surrounding) 클래스에 대한 불필요한 참조를 제거합니다. 클로저 클리너가 비활성화되면 익명 사용자 함수가 보통 Serializable이 아닌 둘러싼 클래스를 참조하는 상황이 발생할 수 있으며, 이는 직렬화기에서 예외를 일으킵니다. 설정은 다음과 같습니다:NONE— 클로저 클리너를 완전히 비활성화,TOP_LEVEL— 필드로 재귀하지 않고 최상위 클래스만 정리,RECURSIVE— 모든 필드를 재귀적으로 정리.getParallelism()/setParallelism(int parallelism)작업의 기본 병렬도를 설정합니다.getMaxParallelism()/setMaxParallelism(int parallelism)작업의 기본 최대 병렬도를 설정합니다. 이 설정은 최대 병렬도 수준을 결정하고 동적 확장의 상한을 지정합니다.getNumberOfExecutionRetries()/setNumberOfExecutionRetries(int numberOfExecutionRetries)실패한 태스크가 재실행되는 횟수를 설정합니다. 값이 0이면 사실상 장애 허용을 비활성화합니다.-1값은 시스템 기본값(구성에서 정의된)을 사용해야 함을 나타냅니다. 이는 더 이상 사용되지 않으며(deprecated) 재시작 전략(restart strategies)을 대신 사용하세요.getExecutionRetryDelay()/setExecutionRetryDelay(long executionRetryDelay)작업이 실패한 후 시스템이 재실행 전에 대기하는 지연 시간(밀리초)을 설정합니다. 지연은 모든 태스크가 TaskManager에서 성공적으로 중지된 후 시작되며, 지연이 지나면 태스크가 재시작됩니다. 이 파라미터는 특정 시간 초과 관련 실패(예: 완전히 시간 초과되지 않은 끊어진 연결)가 재실행을 시도하고 같은 문제로 즉시 다시 실패하기 전에 완전히 드러나도록 재실행을 지연시키는 데 유용합니다. 이 파라미터는 실행 재시도 횟수가 1 이상인 경우에만 효과가 있습니다. 이는 더 이상 사용되지 않으며 재시작 전략을 대신 사용하세요.getExecutionMode()/setExecutionMode(). 기본 실행 모드는 PIPELINED입니다. 프로그램을 실행할 실행 모드를 설정합니다. 실행 모드는 데이터 교환이 배치 방식으로 수행되는지 파이프라인 방식으로 수행되는지를 정의합니다.enableForceKryo()/disableForceKryo. Kryo는 기본적으로 강제되지 않습니다. POJO로 분석할 수 있음에도GenericTypeInformation이 POJO에 Kryo 직렬화기를 사용하도록 강제합니다. 어떤 경우에는 이것이 선호될 수 있습니다. 예를 들어 Flink의 내부 직렬화기가 POJO를 제대로 처리하지 못할 때입니다.enableForceAvro()/disableForceAvro(). Avro는 기본적으로 강제되지 않습니다. FlinkAvroTypeInfo가 Avro POJO 직렬화 시 Kryo 대신 Avro 직렬화기를 사용하도록 강제합니다.enableObjectReuse()/disableObjectReuse(). 기본적으로 Flink에서는 객체가 재사용되지 않습니다. 객체 재사용 모드를 활성화하면 런타임이 더 나은 성능을 위해 사용자 객체를 재사용하도록 지시합니다. 연산의 사용자 코드 함수가 이 동작을 인지하지 못하면 버그가 발생할 수 있음에 유의하세요.getGlobalJobParameters()/setGlobalJobParameters()이 메서드는 사용자가 작업의 전역 구성으로 커스텀 객체를 설정할 수 있게 해줍니다.ExecutionConfig는 모든 사용자 정의 함수에서 접근할 수 있으므로, 작업에서 구성을 전역적으로 사용할 수 있게 하는 쉬운 방법입니다.addDefaultKryoSerializer(Class<?> type, Serializer<?> serializer)주어진type에 대해 Kryo 직렬화기 인스턴스를 등록합니다.addDefaultKryoSerializer(Class<?> type, Class<? extends Serializer<?>> serializerClass)주어진type에 대해 Kryo 직렬화기 클래스를 등록합니다.registerTypeWithKryoSerializer(Class<?> type, Serializer<?> serializer)주어진 타입을 Kryo에 등록하고 그 직렬화기를 지정합니다. 타입을 Kryo에 등록하면 타입의 직렬화가 훨씬 더 효율적입니다.registerKryoType(Class<?> type)타입이 결국 Kryo로 직렬화되면 Kryo에 등록되어 태그(정수 ID)만 쓰이도록 보장합니다. 타입이 Kryo에 등록되지 않으면 전체 클래스 이름이 모든 인스턴스와 함께 직렬화되어 I/O 비용이 훨씬 높아집니다.registerPojoType(Class<?> type)주어진 타입을 직렬화 스택에 등록합니다. 타입이 결국 POJO로 직렬화되면 POJO 직렬화기에 등록됩니다. 타입이 결국 Kryo로 직렬화되면 Kryo에 등록되어 태그만 쓰이도록 보장합니다. 타입이 Kryo에 등록되지 않으면 전체 클래스 이름이 모든 인스턴스와 함께 직렬화되어 I/O 비용이 훨씬 높아집니다.
registerKryoType()으로 등록된 타입은 Flink의 POJO 직렬화기 인스턴스에서는 사용할 수 없음에 유의하세요.
disableAutoTypeRegistration()자동 타입 등록은 기본적으로 활성화되어 있습니다. 자동 타입 등록은 사용자 코드가 사용하는 모든 타입(하위 타입 포함)을 Kryo와 POJO 직렬화기에 등록합니다.setTaskCancellationInterval(long interval)연속적인 실행 중인 태스크 취소 시도 사이에 대기할 간격(밀리초)을 설정합니다. 태스크가 취소되면 태스크 스레드가 일정 시간 내에 종료되지 않을 경우 주기적으로interrupt()를 호출하는 새 스레드가 생성됩니다. 이 파라미터는 연속적인interrupt()호출 사이의 시간을 말하며 기본값은 30000밀리초, 즉 30초입니다.
Rich* 함수에서 getRuntimeContext() 메서드로 접근할 수 있는 RuntimeContext는 모든 사용자 정의 함수에서 ExecutionConfig에 접근할 수도 있게 해줍니다.