PyFlink Configuration
PyFlink Configuration (PyFlink 구성)
PyFlink는 다양한 구성 옵션을 제공해요. Python DataStream API 프로그램과 Python Table API 프로그램 모두에서 구성 옵션을 설정할 수 있어요.
본문
PyFlink는 다양한 구성 옵션을 제공해요:
- Python Configuration (Python 특정 설정)
- Execution Configuration (잡 실행 설정)
- State Backend Configuration (상태 저장 설정)
- Checkpointing Configuration (내결함성 설정)
- Network Configuration (네트워크 버퍼 설정)
최적화를 위해 사용돼요.
Python DataStream API 프로그램의 경우 구성 옵션은 다음과 같이 설정할 수 있어요:
from pyflink.common import Configuration
from pyflink.datastream import StreamExecutionEnvironment
config = Configuration()
config.set_integer("python.fn-execution.bundle.size", 1000)
env = StreamExecutionEnvironment.get_execution_environment(config)
Python Table API 프로그램의 경우 Java/Scala Table API 프로그램에 사용 가능한 모든 구성 옵션을 Python Table API 프로그램에서도 사용할 수 있어요. Table API 프로그램에 사용 가능한 모든 구성 옵션에 대한 자세한 내용은 Table API Configuration을 참조하세요. 구성 옵션은 Table API 프로그램에서 다음과 같이 설정할 수 있어요:
from pyflink.table import TableEnvironment, EnvironmentSettings
env_settings = EnvironmentSettings.in_streaming_mode()
t_env = TableEnvironment.create(env_settings)
t_env.get_config().set("python.fn-execution.bundle.size", "1000")
구성 옵션은 EnvironmentSettings를 만들 때 설정할 수도 있어요:
from pyflink.common import Configuration
from pyflink.table import TableEnvironment, EnvironmentSettings
# create a streaming TableEnvironment
config = Configuration()
config.set_string("python.fn-execution.bundle.size", "1000")
env_settings = EnvironmentSettings \
.new_instance() \
.in_streaming_mode() \
.with_configuration(config) \
.build()
table_env = TableEnvironment.create(env_settings)
# or directly pass config into create method
table_env = TableEnvironment.create(config)
Python Options (Python 옵션)
| Key | Default | Type | Description |
|---|---|---|---|
| python.archives | (none) | String | 잡에 Python 아카이브 파일을 추가해요. 아카이브 파일은 Python UDF 워커의 작업 디렉터리로 추출돼요. 각 아카이브 파일에 대해 대상 디렉터리가 지정돼요. 대상 디렉터리 이름이 지정되면 아카이브 파일은 지정된 이름의 디렉터리로 추출돼요. 그렇지 않으면 아카이브 파일은 아카이브 파일과 같은 이름의 디렉터리로 추출돼요. 이 옵션으로 업로드된 파일은 상대 경로로 접근할 수 있어요. 아카이브 파일 경로와 대상 디렉터리 이름의 구분자로 '#'를 사용할 수 있어요. 여러 아카이브 파일을 지정하려면 쉼표(',')를 구분자로 사용할 수 있어요. 이 옵션은 가상 환경, Python UDF에 사용되는 데이터 파일을 업로드하는 데 사용할 수 있어요. 데이터 파일은 Python UDF에서 접근할 수 있어요. 예: f = open('data/data.txt', 'r'). 이 옵션은 명령줄 옵션 "-pyarch"와 동일해요. |
| python.client.executable | "python" | String | "flink run"으로 Python 잡을 제출하거나 Python UDF를 포함하는 Java/Scala 잡을 컴파일할 때 Python 프로세스를 시작하는 데 사용되는 Python 인터프리터의 경로예요. 명령줄 옵션 "-pyclientexec" 또는 환경 변수 PYFLINK_CLIENT_EXECUTABLE과 동일해요. 우선순위는 다음과 같아요: - 소스 코드에 정의된 'python.client.executable' 구성 (Flink Java SQL/Table API 잡이 Python UDF를 호출할 때만 사용); - 명령줄 옵션 "-pyclientexec"; - config.yaml에 정의된 'python.client.executable' 구성; - 환경 변수 PYFLINK_CLIENT_EXECUTABLE; |
| python.executable | "python" | String | python UDF 워커를 실행하는 데 사용되는 Python 인터프리터의 경로를 지정해요. python UDF 워커는 Python 3.8+, Apache Beam (version >= 2.54.0, <= 2.61.0), Pip (version >= 20.3) 및 SetupTools (version >= 37.0.0)에 의존해요. 지정된 환경이 위 요구 사항을 충족하는지 확인하세요. 이 옵션은 명령줄 옵션 "-pyexec"와 동일해요. |
| python.execution-mode | "process" | String | Python 런타임 실행 모드를 지정해요. 선택 값은 process와 thread예요. process 모드는 Python 사용자 정의 함수가 별도의 Python 프로세스에서 실행된다는 뜻이에요. thread 모드는 Python 사용자 정의 함수가 Java 연산자와 같은 프로세스에서 실행된다는 뜻이에요. 현재 모든 위치에서 thread 모드로 Python 사용자 정의 함수를 실행하는 것을 지원하지는 않는다는 점에 유의하세요. 이러한 경우 process 모드로 폴백돼요. |
| python.files | (none) | String | 잡에 사용자 지정 파일을 첨부해요. .py/.egg/.zip/.whl 또는 디렉터리 같은 표준 리소스 파일 접미사를 모두 지원해요. 이 파일들은 로컬 클라이언트와 원격 python UDF 워커 모두의 PYTHONPATH에 추가돼요. .zip 접미사가 있는 파일은 추출되어 PYTHONPATH에 추가돼요. 여러 파일을 지정하려면 쉼표(',')를 구분자로 사용할 수 있어요. 이 옵션은 명령줄 옵션 "-pyfs"와 동일해요. |
| python.fn-execution.arrow.batch.size | 1000 | Integer | Python 사용자 정의 함수 실행을 위한 arrow 배치에 포함할 최대 요소 수예요. arrow 배치 크기는 번들 크기를 초과해서는 안 돼요. 그렇지 않으면 번들 크기가 arrow 배치 크기로 사용돼요. |
| python.fn-execution.bundle.size | 1000 | Integer | Python 사용자 정의 함수 실행을 위한 번들에 포함할 최대 요소 수예요. 요소는 비동기로 처리돼요. 하나의 번들 요소가 처리된 후 다음 번들 요소가 처리돼요. 더 큰 값은 처리량을 개선할 수 있지만 더 많은 메모리 사용과 더 높은 대기 시간의 비용이 들어요. |
| python.fn-execution.bundle.time | 1000 | Long | Python 사용자 정의 함수 실행을 위해 번들을 처리하기 전에 대기할 타임아웃(밀리초)을 설정해요. 타임아웃은 번들의 요소가 처리되기 전에 얼마나 오래 버퍼링되는지 정의해요. 더 낮은 타임아웃은 더 낮은 꼬리 대기 시간을 초래하지만 처리량에 영향을 줄 수 있어요. |
| python.fn-execution.memory.managed | true | Boolean | 설정하면 Python 워커가 task slot의 managed memory 예산을 사용하도록 구성돼요. 그렇지 않으면 task slot의 Off-Heap Memory를 사용해요. 이 경우 사용자는 구성 키 taskmanager.memory.task.off-heap.size를 사용해 Task Off-Heap Memory를 설정해야 해요. |
| python.map-state.iterate-response-batch-size | 1000 | Integer | Python MapState를 반복할 때 각 배치에서 Python UDF 워커로 보내는 MapState 키/항목의 최대 수예요. 이는 실험적 플래그이며 향후 릴리스에서 제공되지 않을 수 있다는 점에 유의하세요. |
| python.map-state.read-cache-size | 1000 | Integer | 단일 Python MapState에 대한 캐시된 항목의 최대 수예요. 이는 실험적 플래그이며 향후 릴리스에서 제공되지 않을 수 있다는 점에 유의하세요. |
| python.map-state.write-cache-size | 1000 | Integer | 단일 Python MapState에 대한 캐시된 쓰기 요청의 최대 수예요. 캐시된 쓰기 요청 수가 이 한도를 초과하면 쓰기 요청은 상태 백엔드(Java 연산자에서 관리)로 플러시돼요. 이는 실험적 플래그이며 향후 릴리스에서 제공되지 않을 수 있다는 점에 유의하세요. |
| python.metric.enabled | true | Boolean | false이면 Python에 대한 메트릭이 비활성화돼요. 어떤 상황에서는 더 나은 성능을 위해 메트릭을 비활성화할 수 있어요. |
| python.operator-chaining.enabled | true | Boolean | Python 연산자 체이닝은 non-shuffle 연산이 같은 스레드에 함께 배치되어 직렬화와 역직렬화를 완전히 피할 수 있게 해줘요. |
| python.profile.enabled | false | Boolean | Python 워커 프로파일링을 활성화할지 여부를 지정해요. 프로파일 결과는 TaskManager의 로그 파일에 주기적으로 표시돼요. 각 프로파일링 사이의 간격은 config 옵션 python.fn-execution.bundle.size와 python.fn-execution.bundle.time에 의해 결정돼요. |
| python.pythonpath | (none) | String | Worker Node에서 Flink Python 종속성이 설치된 경로를 지정하며, Python Worker의 PYTHONPATH에 추가돼요. 이 옵션은 명령줄 옵션 "-pypath"와 동일해요. |
| python.requirements | (none) | String | 타사 종속성을 정의하는 requirements.txt 파일을 지정해요. 이 종속성은 설치되어 python UDF 워커의 PYTHONPATH에 추가돼요. 선택적으로 이러한 종속성의 설치 패키지를 포함하는 디렉터리를 지정할 수 있어요. 선택 매개변수가 존재하면 '#'를 구분자로 사용해요. 이 옵션은 명령줄 옵션 "-pyreq"와 동일해요. |
| python.state.cache-size | 1000 | Integer | Python UDF 워커에 캐시된 상태의 최대 수예요. 이는 실험적 플래그이며 향후 릴리스에서 제공되지 않을 수 있다는 점에 유의하세요. |
| python.systemenv.enabled | true | Boolean | Python 워커를 시작할 때 System Environment를 로드할지 여부를 지정해요. |