Operators

Operator는 하나 이상의 DataStream을 새 DataStream으로 변환해요. 프로그램은 여러 변환을 결합해 정교한 dataflow 토폴로지를 만들 수 있어요.

출처: Operators (PyFlink)

본문

Operator는 하나 이상의 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 Function)

다음 예제와 같이 변환은 변환의 기능을 정의하기 위해 람다 함수도 받을 수 있어요:

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())

Note

ConnectedStream.map()과 ConnectedStream.flat_map()은 람다 함수를 지원하지 않으며 반드시 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 직렬 변환기를 사용해 직렬화돼요. pickle 직렬 변환기에 대한 자세한 내용은 Data Types 페이지의 Pickle Serialization을 참조하세요.

일반적으로 다음 시나리오에서 출력 타입을 지정해야 해요.

DataStream을 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의 출력 타입이 복합 타입이어야 하기 때문이에요.

DataStream을 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(...)

일반적으로 위 예제의 map 연산에 대해 출력 타입을 지정해야 하는데, 이는 싱크가 Row 등과 같은 특별한 종류의 데이터만 받아들일 때 그렇답니다.

연산자 체이닝 (Operator Chaining)

기본적으로 여러 non-shuffle Python 함수는 직렬화와 역직렬화를 피하고 성능을 개선하기 위해 함께 연결돼요. 체이닝을 비활성화하고 싶은 경우도 있어요. 예를 들어 각 입력 요소에 대해 많은 수의 요소를 생성하는 flatmap 함수가 있고, 체이닝을 비활성화하면 그 출력을 다른 병렬도로 처리할 수 있게 해줘요.

연산자 체이닝은 다음 중 한 가지 방법으로 비활성화할 수 있어요:

  • 현재 연산자 뒤에 key_by 연산, shuffle 연산, rescale 연산, rebalance 연산 또는 partition_custom 연산을 추가하여 이후 연산자와의 체이닝을 비활성화해요.
  • 현재 연산자에 start_new_chain 연산을 적용하여 이전 연산자와의 체이닝을 비활성화해요.
  • 현재 연산자에 disable_chaining 연산을 적용하여 이전과 이후 연산자와의 체이닝을 비활성화해요.
  • 두 연산자에 서로 다른 병렬도 또는 서로 다른 slot sharing group을 설정하여 체이닝을 비활성화해요.
  • 구성 python.operator-chaining.enabled (PyFlink Configuration 참조)로 모든 연산자 체이닝을 비활성화할 수도 있어요.

Python 함수 번들링 (Bundling Python Functions)

Python 함수를 non-local 모드에서 실행하려면 Python 함수 정의가 main() 함수가 정의된 파일 밖에 있다면 config 옵션 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)

더 알아보기 (Learn more)