의존성 관리

의존성 관리 (Dependency Management)

PyFlink 프로그램 안에서 의존성을 사용하기 위한 요구사항이 있어요. 예를 들어 Python 사용자 정의 함수에서 타사 Python 라이브러리를 사용해야 할 수 있고, 머신러닝 예측 같은 시나리오에서는 Python 사용자 정의 함수 안에서 머신러닝 모델을 로드하고 싶을 수 있어요.

출처: 문서

본문

PyFlink job을 로컬에서 실행할 때는 타사 Python 라이브러리를 로컬 Python 환경에 설치하고, 머신러닝 모델을 로컬로 다운로드하는 등의 작업을 할 수 있어요. 하지만 PyFlink job을 원격 클러스터에 제출하려 할 때는 이 접근 방식이 잘 동작하지 않아요. 다음 절들에서는 이러한 요구사항을 위해 PyFlink에서 제공하는 옵션을 소개해요.

참고: Python DataStream API와 Python Table API 모두 각 종류의 의존성에 대한 API를 제공해요. 단일 job에서 Python DataStream API와 Python Table API를 혼합해서 사용한다면, Python DataStream API를 통해 의존성을 지정해 Python DataStream API와 Python Table API 양쪽 모두에서 동작하게 해야 해요.

JAR 의존성 (JAR Dependencies)

타사 JAR을 사용한다면 Python Table API에서 다음과 같이 JAR을 지정할 수 있어요.

# Specify a list of jar URLs via "pipeline.jars". The jars are separated by ";"
# and will be uploaded to the cluster.
# 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")

# It looks like the following on Windows:
table_env.get_config().set("pipeline.jars", "file:///E:/my/jar/path/connector.jar;file:///E:/my/jar/path/udf.jar")

# Specify a list of URLs via "pipeline.classpaths". The URLs are separated by ";"
# and will be added to the classpath during job execution.
# 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")

또는 Python DataStream API에서 다음과 같이 지정해요.

# Use the add_jars() to add local jars and the jars will be uploaded to the cluster.
# NOTE: Only local file URLs (start with "file://") are supported.
stream_execution_environment.add_jars("file:///my/jar/path/connector1.jar", "file:///my/jar/path/connector2.jar")

# It looks like the following on Windows:
stream_execution_environment.add_jars("file:///E:/my/jar/path/connector1.jar", "file:///E:/my/jar/path/connector2.jar")

# Use the add_classpaths() to add the dependent jars URLs into the classpath.
# The URLs will also be added to the classpath of both the client and the cluster.
# 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.
stream_execution_environment.add_classpaths("file:///my/jar/path/connector1.jar", "file:///my/jar/path/connector2.jar")

또는 job을 제출할 때 명령줄 인자 --jarfile을 통해 지정할 수 있어요.

참고: --jarfile 명령줄 인자로는 하나의 jar 파일만 지정할 수 있으므로, 여러 jar 파일이 있다면 fat jar를 빌드해야 해요.

Python 의존성 (Python Dependencies)

Python 라이브러리 (Python libraries)

Python 사용자 정의 함수에서 타사 Python 라이브러리를 사용하고 싶을 수 있어요. Python 라이브러리를 지정하는 방법은 여러 가지예요.

다음과 같이 코드 안에서 Python Table API를 사용해 지정할 수 있어요.

table_env.add_python_file(file_path)

또는 Python DataStream API를 사용해요.

stream_execution_environment.add_python_file(file_path)

Python 라이브러리를 구성(python.files, PyFlink Configuration 참고) 또는 job 제출 시 명령줄 인자 -pyfs 또는 --pyFiles로 지정할 수도 있어요.

참고: Python 라이브러리는 로컬 파일 또는 로컬 디렉터리일 수 있어요. 이들은 Python UDF worker의 PYTHONPATH에 추가돼요.

requirements.txt

타사 Python 의존성을 정의하는 requirements.txt 파일을 지정하는 것도 허용돼요. 이러한 Python 의존성은 작업 디렉터리에 설치되고 Python UDF worker의 PYTHONPATH에 추가돼요.

requirements.txt를 수동으로 준비할 수 있어요.

echo numpy==1.16.5 >> requirements.txt
echo pandas==1.0.0 >> requirements.txt

또는 현재 Python 환경에 설치된 모든 패키지를 나열하는 pip freeze를 사용해요.

pip freeze > requirements.txt

requirements.txt 파일의 내용은 다음과 같을 수 있어요.

numpy==1.16.5
pandas==1.0.0

불필요한 항목을 제거하거나 추가 항목을 추가하는 등 수동으로 편집할 수 있어요. 그다음 requirements.txt 파일을 코드 안에서 Python Table API로 지정할 수 있어요.

# requirements_cache_dir is optional
table_env.set_python_requirements(
    requirements_file_path="/path/to/requirements.txt",
    requirements_cache_dir="cached_dir")

또는 Python DataStream API로 지정해요.

# requirements_cache_dir is optional
stream_execution_environment.set_python_requirements(
    requirements_file_path="/path/to/requirements.txt",
    requirements_cache_dir="cached_dir")

참고: 클러스터에서 접근할 수 없는 의존성을 위해, 이 의존성들의 설치 패키지를 포함하는 디렉터리를 requirements_cached_dir 매개변수로 지정할 수 있어요. 이 디렉터리는 오프라인 설치를 지원하기 위해 클러스터에 업로드돼요. requirements_cache_dir을 다음과 같이 준비할 수 있어요.

pip download -d cached_dir -r requirements.txt --no-binary :all:

참고: 준비한 패키지가 클러스터의 플랫폼과 사용되는 Python 버전과 일치하는지 확인하세요.

requirements.txt 파일을 구성(python.requirements, PyFlink Configuration 참고) 또는 job 제출 시 명령줄 인자 -pyreq 또는 --pyRequirements로 지정할 수도 있어요.

참고: requirements.txt 파일에 지정된 패키지를 pip로 설치하므로, pip (버전 >= 20.3)와 setuptools (버전 >= 37.0.0)가 사용 가능한지 확인하세요.

아카이브 (Archives)

아카이브 파일을 지정하고 싶을 수도 있어요. 아카이브 파일은 사용자 정의 Python 가상 환경, 데이터 파일 등을 지정하는 데 사용할 수 있어요.

코드 안에서 Python Table API로 아카이브 파일을 지정할 수 있어요.

table_env.add_python_archive(archive_path="/path/to/archive_file", target_dir=None)

또는 Python DataStream API로 지정해요.

stream_execution_environment.add_python_archive(archive_path="/path/to/archive_file", target_dir=None)

참고: target_dir 매개변수는 선택 사항이에요. 지정하면 아카이브 파일이 실행 중에 target_dir의 지정된 이름을 가진 디렉터리로 추출돼요. 그렇지 않으면 아카이브 파일이 아카이브 파일과 같은 이름의 디렉터리로 추출돼요.

다음과 같이 아카이브 파일을 지정했다고 가정해요.

table_env.add_python_archive("/path/to/py_env.zip", "myenv")

그러면 Python 사용자 정의 함수에서 다음과 같이 아카이브 파일의 내용에 접근할 수 있어요.

def my_udf():
    with open("myenv/py_env/data/data.txt") as f:
        ...

target_dir 매개변수를 지정하지 않았다면,

table_env.add_python_archive("/path/to/py_env.zip")

Python 사용자 정의 함수에서 다음과 같이 아카이브 파일의 내용에 접근할 수 있어요.

def my_udf():
    with open("py_env.zip/py_env/data/data.txt") as f:
        ...

참고: 아카이브 파일은 Python UDF worker의 작업 디렉터리로 추출되므로 상대 경로를 사용해 아카이브 파일 안의 파일에 접근할 수 있어요.

아카이브 파일을 구성(python.archives, PyFlink Configuration 참고) 또는 job 제출 시 명령줄 인자 -pyarch 또는 --pyArchives로 지정할 수도 있어요.

참고: 아카이브 파일이 Python 가상 환경을 포함한다면, Python 가상 환경이 클러스터가 실행되는 플랫폼과 일치하는지 확인하세요.

참고: 현재는 zip 파일(즉 zip, jar, whl, egg 등)과 tar 파일(즉 tar, tar.gz, tgz)만 지원돼요.

Python 인터프리터 (Python interpreter)

Python worker를 실행할 Python 인터프리터의 경로를 지정하는 것을 지원해요. 코드 안에서 Python Table API로 Python 인터프리터를 지정할 수 있어요.

table_env.get_config().set_python_executable("/path/to/python")

또는 Python DataStream API로 지정해요.

stream_execution_environment.set_python_executable("/path/to/python")

아카이브 파일 안의 Python 인터프리터를 사용하는 것도 지원돼요.

# Python Table API
table_env.add_python_archive("/path/to/py_env.zip", "venv")
table_env.get_config().set_python_executable("venv/py_env/bin/python")

# Python DataStream API
stream_execution_environment.add_python_archive("/path/to/py_env.zip", "venv")
stream_execution_environment.set_python_executable("venv/py_env/bin/python")

Python 인터프리터를 구성(python.executable, PyFlink Configuration 참고) 또는 job 제출 시 명령줄 인자 -pyexec 또는 --pyExecutable로 지정할 수도 있어요.

참고: Python 인터프리터의 경로가 Python 아카이브 파일을 참조한다면 절대 경로 대신 상대 경로를 사용해야 해요.

클라이언트의 Python 인터프리터 (Python interpreter of client)

job을 컴파일하는 동안 Python 사용자 정의 함수를 파싱하려면 클라이언트 측에서 Python이 필요해요. 현재 세션에서 활성화해 클라이언트 측에서 사용할 사용자 정의 Python 인터프리터를 지정할 수 있어요.

source my_env/bin/activate

또는 구성(python.client.executable, PyFlink Configuration 참고), 명령줄 인자 -pyclientexec 또는 --pyClientExecutable, 환경 변수 PYFLINK_CLIENT_EXECUTABLE(Environment Variables 참고)로 지정할 수 있어요.

Java/Scala 프로그램에서 Python 의존성 지정하기

Java Table API 프로그램이나 순수 SQL 프로그램에서 Python 사용자 정의 함수를 사용하는 것도 지원돼요. 다음 코드는 Java Table API 프로그램에서 Python 사용자 정의 함수를 사용하는 간단한 예시를 보여줘요.

import org.apache.flink.configuration.CoreOptions;
import org.apache.flink.table.api.EnvironmentSettings;
import org.apache.flink.table.api.TableEnvironment;

TableEnvironment tEnv = TableEnvironment.create(
    EnvironmentSettings.inBatchMode());
tEnv.getConfig().set(CoreOptions.DEFAULT_PARALLELISM, 1);

// register the Python UDF
tEnv.executeSql("create temporary system function add_one as 'add_one.add_one' language python");

tEnv.createTemporaryView("source", tEnv.fromValues(1L, 2L, 3L).as("a"));

// use Python UDF in the Java Table API program
tEnv.executeSql("select add_one(a) as a from source").collect();

SQL 문으로 Python 사용자 정의 함수를 생성하는 방법에 대한 자세한 내용은 CREATE FUNCTION SQL 문을 참고할 수 있어요. 그다음 Python 의존성은 python.archives, python.files, python.requirements, python.client.executable, python.executable 같은 Python 구성 옵션(PyFlink Configuration 참고) 또는 job 제출 시 명령줄 인자를 통해 지정할 수 있어요.

더 알아보기 (Learn more)