연산자
연산자 (Operators)
연산자(Operators)는 하나 이상의 DataStream을 새로운 DataStream으로 변환해요. 프로그램은 여러 변환을 결합해 정교한 dataflow 토폴로지를 만들 수 있어요.
출처: 문서
본문
DataStream 변환 (DataStream Transformations)
Flink의 DataStream 프로그램은 데이터 스트림에 변환(예: 매핑, 필터링, 리듀싱)을 구현하는 일반적인 프로그램이에요. Python DataStream API에서 사용 가능한 변환의 개요는 operators를 참고하세요.
함수 (Functions)
변환은 변환의 기능을 정의하는 사용자 정의 함수를 입력으로 받아요. 다음 절에서는 Python DataStream API에서 Python 사용자 정의 함수를 정의하는 다양한 방법을 설명해요.
함수 인터페이스 구현 (Implementing Function Interfaces)
Python DataStream API에는 서로 다른 변환을 위해 서로 다른 Function 인터페이스가 제공돼요. 예를 들어 map 변환에는 MapFunction이, filter 변환에는 FilterFunction 등이 제공돼요. 사용자는 변환의 유형에 따라 해당 Function 인터페이스를 구현할 수 있어요. MapFunction을 예로 들면,
# Implementing MapFunction
class MyMapFunction(MapFunction):
def map(self, value):
return value + 1
data_stream = env.from_collection([1, 2, 3, 4, 5], type_info=Types.INT())
mapped_stream = data_stream.map(MyMapFunction(), output_type=Types.INT())
Lambda 함수 (Lambda Function)
다음 예시처럼 변환은 변환의 기능을 정의하기 위해 lambda 함수도 받을 수 있어요.
data_stream = env.from_collection([1, 2, 3, 4, 5], type_info=Types.INT())
mapped_stream = data_stream.map(lambda x: x + 1, output_type=Types.INT())
참고:
ConnectedStream.map()과ConnectedStream.flat_map()은 lambda 함수를 지원하지 않으며, 각각CoMapFunction과CoFlatMapFunction을 받아야 해요.
Python 함수 (Python Function)
사용자는 변환의 기능을 정의하기 위해 Python 함수를 사용할 수도 있어요.
def my_map_func(value):
return value + 1
data_stream = env.from_collection([1, 2, 3, 4, 5], type_info=Types.INT())
mapped_stream = data_stream.map(my_map_func, output_type=Types.INT())
출력 타입 (Output Type)
사용자는 Python DataStream API에서 변환의 출력 타입 정보를 명시적으로 지정할 수 있어요. 지정하지 않으면 기본적으로 출력 타입이 Types.PICKLED_BYTE_ARRAY가 되고, 결과 데이터는 pickle serializer를 사용해 직렬화돼요. pickle serializer에 대한 자세한 내용은 Data Types 페이지의 Pickle Serialization을 참고하세요. 일반적으로 출력 타입은 다음 시나리오에서 지정해야 해요.
DataStream을 Table로 변환 (Convert DataStream into Table)
from pyflink.common.typeinfo import Types
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment
def data_stream_api_demo():
env = StreamExecutionEnvironment.get_execution_environment()
t_env = StreamTableEnvironment.create(stream_execution_environment=env)
t_env.execute_sql("""
CREATE TABLE my_source (
a INT,
b VARCHAR
) WITH (
'connector' = 'datagen',
'number-of-rows' = '10'
)
""")
ds = t_env.to_append_stream(
t_env.from_path('my_source'),
Types.ROW([Types.INT(), Types.STRING()]))
def split(s):
splits = s[1].split("|")
for sp in splits:
yield s[0], sp
ds = ds.map(lambda i: (i[0] + 1, i[1])) \
.flat_map(split, Types.TUPLE([Types.INT(), Types.STRING()])) \
.key_by(lambda i: i[1]) \
.reduce(lambda i, j: (i[0] + j[0], i[1]))
t_env.execute_sql("""
CREATE TABLE my_sink (
a INT,
b VARCHAR
) WITH (
'connector' = 'print'
)
""")
table = t_env.from_data_stream(ds)
table_result = table.execute_insert("my_sink")
# 1) wait for job finishes and only used in local execution, otherwise, it may happen that the script exits with the job is still running
# 2) should be removed when submitting the job to a remote cluster such as YARN, standalone, K8s etc in detach mode
table_result.wait()
if __name__ == '__main__':
data_stream_api_demo()
위 예시에서 flat_map 연산에 대해 출력 타입을 지정해야 해요. 이는 reduce 연산의 출력 타입으로 암시적으로 사용될 거예요. 이유는 t_env.from_data_stream(ds)가 ds의 출력 타입이 복합(composite) 타입이어야 하기 때문이에요.
DataStream을 Sink에 쓰기 (Write DataStream to Sink)
from pyflink.common.typeinfo import Types
def split(s):
splits = s[1].split("|")
for sp in splits:
yield s[0], sp
ds.map(lambda i: (i[0] + 1, i[1]), Types.TUPLE([Types.INT(), Types.STRING()])) \
.sink_to(...)
일반적으로 싱크가 특별한 종류의 데이터(예: Row)만 받아들인다면 위 예시의 map 연산에 대해 출력 타입을 지정해야 해요.
연산자 체이닝 (Operator Chaining)
기본적으로 여러 non-shuffle Python 함수는 직렬화·역직렬화를 피하고 성능을 향상시키기 위해 체이닝돼요. 체이닝을 비활성화하고 싶은 경우도 있어요. 예를 들어 각 입력 요소에 대해 많은 수의 요소를 생성하는 flatmap 함수가 있고, 체이닝을 비활성화하면 그 출력을 다른 parallelism으로 처리할 수 있어요.
연산자 체이닝은 다음 중 한 가지 방식으로 비활성화할 수 있어요.
- 현재 연산자 이후에
key_by,shuffle,rescale,rebalance또는partition_custom연산을 추가해 다음 연산자와의 체이닝을 비활성화 - 현재 연산자에
start_new_chain연산을 적용해 이전 연산자와의 체이닝을 비활성화 - 현재 연산자에
disable_chaining연산을 적용해 이전 및 다음 연산자와의 체이닝을 모두 비활성화 - 두 연산자에 서로 다른 parallelism 또는 서로 다른 slot sharing group을 설정해 두 연산자의 체이닝을 비활성화
구성 python.operator-chaining.enabled(PyFlink Configuration 참고)로 모든 연산자 체이닝을 비활성화할 수도 있어요.
Python 함수 번들링 (Bundling Python Functions)
로컬이 아닌 모드에서 Python 함수를 실행하려면, Python 함수가 main() 함수가 정의된 파일 밖에 있다면 구성 옵션 python-files(PyFlink Configuration 참고)로 Python 함수 정의를 번들링하는 것을 강력히 권장해요. 그렇지 않으면 my_function.py라는 파일에 Python 함수를 정의했을 때 ModuleNotFoundError: No module named 'my_function' 오류가 발생할 수 있어요.
Python 함수에서 리소스 로드 (Loading resources in Python Functions)
먼저 Python 함수에서 어떤 리소스를 로드한 다음, 리소스를 다시 로드하지 않고 계산을 반복해서 실행하고 싶은 시나리오가 있어요. 예를 들어 큰 딥러닝 모델을 한 번만 로드하고, 모델에 대해 배치 예측을 여러 번 실행하고 싶을 수 있어요.
기본 클래스 Function에서 상속된 open 메서드를 오버라이드하면 정확히 그것을 해요.
class Predict(MapFunction):
def open(self, runtime_context: RuntimeContext):
import pickle
with open("resources.zip/resources/model.pkl", "rb") as f:
self.model = pickle.load(f)
def eval(self, x):
return self.model.predict(x)