커넥터

이 페이지는 PyFlink에서 커넥터를 사용하는 방법을 설명하고 Python 프로그램에서 Flink 커넥터를 사용할 때 주의해야 할 세부 사항을 강조합니다.

참고: 일반적인 커넥터 정보와 공통 구성은 해당 Java/Scala 문서를 참고하세요.

출처: 문서

본문

커넥터 및 포맷 JAR 다운로드 (Download connector and format jars)

Flink는 Java/Scala 기반 프로젝트이므로 커넥터와 포맷 모두 구현체가 jar로 제공되며, 이를 작업 의존성으로 지정해야 합니다(의존성 관리 참고).

table_env.get_config().set("pipeline.jars", "file:///my/jar/path/connector.jar;file:///my/jar/path/json.jar")

커넥터 사용법 (How to use connectors)

PyFlink의 Table API에서 DDL이 소스와 싱크를 정의하는 권장 방식이며, TableEnvironmentexecute_sql() 메서드로 실행됩니다. 이렇게 하면 해당 테이블을 애플리케이션에서 사용할 수 있게 됩니다.

source_ddl = """
        CREATE TABLE source_table(
            a VARCHAR,
            b INT
        ) WITH (
          'connector' = 'kafka',
          'topic' = 'source_topic',
          'properties.bootstrap.servers' = 'kafka:9092',
          'properties.group.id' = 'test_3',
          'scan.startup.mode' = 'latest-offset',
          'format' = 'json'
        )
        """

sink_ddl = """
        CREATE TABLE sink_table(
            a VARCHAR
        ) WITH (
          'connector' = 'kafka',
          'topic' = 'sink_topic',
          'properties.bootstrap.servers' = 'kafka:9092',
          'format' = 'json'
        )
        """

t_env.execute_sql(source_ddl)
t_env.execute_sql(sink_ddl)

t_env.sql_query("SELECT a FROM source_table") \
    .execute_insert("sink_table").wait()

다음은 PyFlink에서 Kafka 소스/싱크와 JSON 포맷을 사용하는 전체 예시입니다.

from pyflink.table import TableEnvironment, EnvironmentSettings

def log_processing():
    env_settings = EnvironmentSettings.in_streaming_mode()
    t_env = TableEnvironment.create(env_settings)
    # specify connector and format jars
    t_env.get_config().set("pipeline.jars", "file:///my/jar/path/connector.jar;file:///my/jar/path/json.jar")

    source_ddl = """
            CREATE TABLE source_table(
                a VARCHAR,
                b INT
            ) WITH (
              'connector' = 'kafka',
              'topic' = 'source_topic',
              'properties.bootstrap.servers' = 'kafka:9092',
              'properties.group.id' = 'test_3',
              'scan.startup.mode' = 'latest-offset',
              'format' = 'json'
            )
            """

    sink_ddl = """
            CREATE TABLE sink_table(
                a VARCHAR
            ) WITH (
              'connector' = 'kafka',
              'topic' = 'sink_topic',
              'properties.bootstrap.servers' = 'kafka:9092',
              'format' = 'json'
            )
            """

    t_env.execute_sql(source_ddl)
    t_env.execute_sql(sink_ddl)

    t_env.sql_query("SELECT a FROM source_table") \
        .execute_insert("sink_table").wait()

if __name__ == '__main__':
    log_processing()

사전 정의된 소스와 싱크 (Predefined Sources and Sinks)

일부 데이터 소스와 싱크는 Flink에 내장되어 있어 별도 설정 없이 사용할 수 있습니다. 이러한 사전 정의 데이터 소스에는 Pandas DataFrame에서 읽기, 컬렉션에서 데이터 수집이 포함됩니다. 사전 정의 데이터 싱크는 Pandas DataFrame에 쓰기를 지원합니다.

from/to Pandas

PyFlink 테이블은 Pandas DataFrame 간 변환을 지원합니다.

from pyflink.table.expressions import col

import pandas as pd
import numpy as np

# Create a PyFlink Table
pdf = pd.DataFrame(np.random.rand(1000, 2))
table = t_env.from_pandas(pdf, ["a", "b"]).filter(col('a') > 0.5)

# Convert the PyFlink Table to a Pandas DataFrame
pdf = table.to_pandas()

from_elements()

from_elements()는 요소 컬렉션에서 테이블을 만드는 데 사용됩니다. 요소 타입은 허용되는 원자 타입(atomic types) 또는 허용되는 복합 타입(composite types)이어야 합니다.

from pyflink.table import DataTypes

table_env.from_elements([(1, 'Hi'), (2, 'Hello')])

# use the second parameter to specify custom field names
table_env.from_elements([(1, 'Hi'), (2, 'Hello')], ['a', 'b'])

# use the second parameter to specify a custom table schema
table_env.from_elements([(1, 'Hi'), (2, 'Hello')],
                        DataTypes.ROW([DataTypes.FIELD("a", DataTypes.INT()),
                                       DataTypes.FIELD("b", DataTypes.STRING())]))

위 쿼리는 다음과 같은 테이블을 반환합니다:

+----+-------+
| a  |   b   |
+====+=======+
| 1  |  Hi   |
+----+-------+
| 2  | Hello |
+----+-------+

사용자 정의 소스 & 싱크 (User-defined sources & sinks)

어떤 경우에는 커스텀 소스와 싱크를 정의하고 싶을 수 있습니다. 현재 소스와 싱크는 Java/Scala로 구현되어야 하지만, DDL을 통한 사용을 지원하기 위해 TableFactory를 정의할 수 있습니다. 더 자세한 내용은 Java/Scala 문서에서 찾을 수 있습니다.

더 알아보기 (Learn more)