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을 실행하기 전에 환경을 활성화해야 해요.

로컬 (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_archiveset_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 또는 TableEnvironmentadd_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을 실행할 때는 이러한 코드를 제거하는 것을 기억하세요.

더 알아보기 (Learn more)