커넥터
커넥터 (Connectors, PyFlink)
이 페이지는 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이 소스와 싱크를 정의하는 권장 방식이며, TableEnvironment의 execute_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 문서에서 찾을 수 있습니다.