FAQ
FAQ (자주 묻는 질문)
이 페이지는 PyFlink 사용자들을 위한 몇 가지 흔한 질문에 대한 해결책을 설명해요. Python 가상 환경 준비, JAR 파일 추가, Python 파일 추가, minicluster에서 job 실행 대기 방법 등을 다뤄요.
출처: 문서
본문
이 페이지는 PyFlink 사용자들을 위한 몇 가지 흔한 질문에 대한 해결책을 설명해요.
Python 가상 환경 준비 (Preparing Python Virtual Environment)
Mac OS와 대부분의 Linux 배포에서 사용할 수 있는 Python 가상 환경 zip을 준비하는 편리한 스크립트를 다운로드할 수 있어요. 해당 PyFlink 버전에 필요한 Python 가상 환경을 생성하도록 PyFlink 버전을 지정할 수 있으며, 지정하지 않으면 가장 최근 버전이 설치돼요.
$ sh setup-pyflink-virtual-env.sh
Python 가상 환경으로 PyFlink job 실행하기
이전 절에서 설명한 대로 Python 가상 환경을 설정한 뒤에는, PyFlink job을 실행하기 전에 환경을 활성화해야 해요.
로컬 (Local)
# activate the python virtual environment
$ source venv/bin/activate
$ python xxx.py
클러스터 (Cluster)
# specify the Python virtual environment
table_env.add_python_archive("venv.zip")
# specify the path of the python interpreter which is used to execute the python UDF workers
table_env.get_config().set_python_executable("venv.zip/venv/bin/python")
add_python_archive와 set_python_executable의 사용법에 대한 자세한 내용은 Dependency Management를 참고할 수 있어요.
JAR 파일 추가 (Adding Jar Files)
PyFlink job은 jar 파일(예: 커넥터, Java UDF 등)에 의존할 수 있어요. 다음 Python Table API로 의존성을 지정하거나, job을 제출할 때 명령줄 인자로 직접 지정할 수 있어요.
# NOTE: Only local file URLs (start with "file:") are supported.
table_env.get_config().set("pipeline.jars", "file:///my/jar/path/connector.jar;file:///my/jar/path/udf.jar")
# NOTE: The Paths must specify a protocol (e.g. "file") and users should ensure that the URLs are accessible on both the client and the cluster.
table_env.get_config().set("pipeline.classpaths", "file:///my/jar/path/connector.jar;file:///my/jar/path/udf.jar")
Java 의존성을 추가하는 API에 대한 자세한 내용은 Dependency Management를 참고할 수 있어요.
Python 파일 추가 (Adding Python Files)
명령줄 인자 pyfs 또는 TableEnvironment의 add_python_file API를 사용해 python 파일, python 패키지, 또는 로컬 디렉터리일 수 있는 python 파일 의존성을 추가할 수 있어요.
예를 들어 다음과 같은 계층을 가진 myDir라는 디렉터리가 있다면,
myDir
+-- utils
+-- __init__.py
+-- my_util.py
다음과 같이 myDir 디렉터리의 Python 파일을 추가할 수 있어요.
table_env.add_python_file('myDir')
def my_udf():
from utils import my_util
minicluster에서 job 실행 시 완료 대기 (Wait for jobs to finish when executing jobs in mini cluster)
minicluster에서 job을 실행할 때(예: IDE에서 job을 실행할 때)와 job에서 다음 API를 사용할 때(예: Python Table API의 TableEnvironment.execute_sql, StatementSet.execute; Python DataStream API의 StreamExecutionEnvironment.execute_async), 이러한 API는 비동기이므로 job 실행이 끝날 때까지 명시적으로 기다리는 것을 잊지 마세요. 그렇지 않으면 job 실행이 끝나기 전에 프로그램이 종료되어 실행 결과를 찾지 못할 수 있어요.
다음 예시를 참고하세요.
# execute SQL / Table API query asynchronously
t_result = table_env.execute_sql(...)
t_result.wait()
# execute DataStream Job asynchronously
job_client = stream_execution_env.execute_async('My DataStream Job')
job_client.get_job_execution_result().result()
참고: 원격 클러스터에서 job을 실행할 때는 job 실행이 끝날 때까지 기다릴 필요가 없으므로, 원격 클러스터에서 job을 실행할 때는 이러한 코드를 제거하는 것을 기억하세요.