TableEnvironment

TableEnvironment

이 문서는 Table API 의 핵심 개념인 TableEnvironment 에 대한 소개입니다. TableEnvironment 클래스의 공개 인터페이스에 대한 상세한 설명을 포함합니다.

출처: 문서

본문

이 문서는 Table API 의 핵심 개념인 TableEnvironment 에 대한 소개입니다. TableEnvironment 클래스의 공개 인터페이스에 대한 상세한 설명을 포함합니다.

TableEnvironment 만들기

TableEnvironment 를 만드는 권장 방법은 EnvironmentSettings 객체로부터 만드는 것입니다:

Java:

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

// create a streaming TableEnvironment
EnvironmentSettings settings = EnvironmentSettings
    .newInstance()
    .inStreamingMode()
    .build();

TableEnvironment tableEnv = TableEnvironment.create(settings);

Scala:

import org.apache.flink.table.api.{EnvironmentSettings, TableEnvironment}

// create a streaming TableEnvironment
val settings = EnvironmentSettings
    .newInstance()
    .inStreamingMode()
    .build()

val tableEnv = TableEnvironment.create(settings)

Python:

from pyflink.common import Configuration
from pyflink.table import EnvironmentSettings, TableEnvironment

# create a streaming TableEnvironment
config = Configuration()
config.set_string('execution.buffer-timeout', '1 min')
env_settings = EnvironmentSettings \
    .new_instance() \
    .in_streaming_mode() \
    .with_configuration(config) \
    .build()

table_env = TableEnvironment.create(env_settings)

또는 사용자는 DataStream API 와 상호 운용하기 위해 기존의 StreamExecutionEnvironment 로부터 StreamTableEnvironment 를 만들 수 있습니다.

Java:

import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;

// create a streaming TableEnvironment from a StreamExecutionEnvironment
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);

Scala:

import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment
import org.apache.flink.table.api.bridge.scala.StreamTableEnvironment

// create a streaming TableEnvironment from a StreamExecutionEnvironment
val env = StreamExecutionEnvironment.getExecutionEnvironment
val tableEnv = StreamTableEnvironment.create(env)

Python:

from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment

# create a streaming TableEnvironment from a StreamExecutionEnvironment
env = StreamExecutionEnvironment.get_execution_environment()
table_env = StreamTableEnvironment.create(env)

TableEnvironment API

Table/SQL 연산

이 API 는 Table API/SQL 테이블을 만들고 제거하며 쿼리를 작성하는 데 사용됩니다:

Java/Scala Python 설명
fromValues(values...) from_elements(elements, schema) 값 모음에서 테이블을 만듭니다.
N/A from_pandas(pdf, schema) pandas DataFrame 에서 테이블을 만듭니다.
from(path) from_path(path) 지정된 경로 아래에 등록된 테이블에서 테이블을 만듭니다.
createTemporaryView(path, table) create_temporary_view(path, table) SQL 임시 뷰와 유사하게 Table 을 임시 뷰로 등록합니다.
dropTemporaryView(path) drop_temporary_view(path) 주어진 경로 아래에 등록된 임시 뷰를 제거합니다.
createTemporaryTable(path, descriptor) create_temporary_table(path, descriptor) TableDescriptor 에서 임시 테이블을 만듭니다.
createTable(path, descriptor) create_table(path, descriptor) TableDescriptor 에서 카탈로그 테이블을 만듭니다.
dropTemporaryTable(path) drop_temporary_table(path) 주어진 경로 아래에 등록된 임시 테이블을 제거합니다.
dropTable(path) drop_table(path) 주어진 경로 아래에 등록된 테이블을 제거합니다.
executeSql(stmt) execute_sql(stmt) 주어진 SQL 문을 실행하고 실행 결과를 반환합니다. DDL/DML/DQL/SHOW/DESCRIBE/EXPLAIN/USE 문을 지원합니다.
sqlQuery(query) sql_query(query) SQL 쿼리를 평가하고 결과를 Table 로 검색합니다.

전체 API 참조는 Java API 문서Python API 문서 를 참고하세요.

더 이상 사용되지 않는 API (Deprecated):

Java/Scala Python 설명 대체
scan(tablePath...) scan(*table_path) 카탈로그에서 등록된 테이블을 스캔합니다. from / from_path
registerTable(name, table) register_table(name, table) Table 을 고유한 이름으로 등록합니다. createTemporaryView / create_temporary_view
insertInto(targetPath, table) insert_into(target_path, table) Table 내용을 sink 에 씁니다. Table.executeInsert() / execute_insert()
N/A sql_update(stmt) SQL 문을 평가합니다. execute_sql

작업 실행/설명 (Execute/Explain)

이 API 는 작업을 설명하거나 실행하는 데 사용됩니다. executeSql / execute_sql 도 작업을 실행하는 데 사용할 수 있습니다.

Java/Scala Python 설명
explainSql(stmt, extraDetails...) explain_sql(stmt, *extra_details) 지정된 문의 AST 와 실행 계획을 반환합니다.
createStatementSet() create_statement_set() 여러 문을 단일 작업으로 실행하기 위한 StatementSet 을 만듭니다.

사용자 정의 함수 만들기/제거

이 API 는 UDF 를 등록하거나 등록된 UDF 를 제거하는 데 사용됩니다. executeSql / execute_sql 도 SQL DDL 을 통해 UDF 를 등록/제거하는 데 사용할 수 있습니다. UDF 에 대한 자세한 내용은 User Defined Functions 를 참고하세요.

Java/Scala Python 설명
createTemporaryFunction(path, class) create_temporary_function(path, function) 함수를 임시 카탈로그 함수로 등록합니다.
createTemporarySystemFunction(name, class) create_temporary_system_function(name, function) 함수를 임시 시스템 함수로 등록합니다.
createFunction(path, class) create_java_function(path, function_class_name) 함수 클래스를 카탈로그 함수로 등록합니다.
N/A create_java_temporary_function(path, function_class_name) Java UDF 클래스를 임시 카탈로그 함수로 등록합니다.
N/A create_java_temporary_system_function(name, function_class_name) Java UDF 클래스를 임시 시스템 함수로 등록합니다.
dropFunction(path) drop_function(path) 주어진 경로 아래에 등록된 카탈로그 함수를 제거합니다.
dropTemporaryFunction(path) drop_temporary_function(path) 주어진 경로 아래에 등록된 임시 함수를 제거합니다.
dropTemporarySystemFunction(name) drop_temporary_system_function(name) 주어진 이름으로 등록된 임시 시스템 함수를 제거합니다.

의존성 관리 (Python 전용)

이 API 는 Python UDF 가 필요로 하는 Python 의존성을 관리하는 데 사용됩니다. 자세한 내용은 Dependency Management 를 참고하세요.

메서드 설명
add_python_file(file_path) UDF worker 의 PYTHONPATH 에 Python 의존성(파일, 패키지 또는 디렉터리)을 추가합니다.
set_python_requirements(requirements_file_path, cache_dir) 타사 의존성을 위한 requirements.txt 파일을 지정합니다.
add_python_archive(archive_path, target_dir) UDF worker 의 작업 디렉터리에서 추출할 Python 아카이브 파일을 추가합니다.

구성 (Configuration)

Java:

// get the TableConfig
TableConfig config = tableEnv.getConfig();

// set configuration options
config.set("parallelism.default", "8");
config.set("pipeline.name", "my_first_job");

사용 가능한 모든 옵션은 Configuration 을 참고하세요.

Scala:

// get the TableConfig
val config = tableEnv.getConfig

// set configuration options
config.set("parallelism.default", "8")
config.set("pipeline.name", "my_first_job")

사용 가능한 모든 옵션은 Configuration 을 참고하세요.

Python:

# get the TableConfig
config = table_env.get_config()

# set configuration options
config.set("parallelism.default", "8")
config.set("pipeline.name", "my_first_job")

사용 가능한 모든 옵션은 ConfigurationPython Configuration 을 참고하세요.

카탈로그 API

이 API 는 카탈로그와 모듈에 접근하는 데 사용됩니다. 자세한 내용은 ModulesCatalogs 를 참고하세요.

Java/Scala Python 설명
registerCatalog(name, catalog) register_catalog(name, catalog) Catalog 를 고유한 이름으로 등록합니다.
getCatalog(name) get_catalog(name) 이름으로 등록된 Catalog 를 가져옵니다.
useCatalog(name) use_catalog(name) 현재 카탈로그를 설정합니다.
getCurrentCatalog() get_current_catalog() 현재 기본 카탈로그 이름을 가져옵니다.
useDatabase(name) use_database(name) 현재 기본 데이터베이스를 설정합니다.
getCurrentDatabase() get_current_database() 현재 기본 데이터베이스 이름을 가져옵니다.
loadModule(name, module) load_module(name, module) Module 을 고유한 이름으로 로드합니다.
unloadModule(name) unload_module(name) 주어진 이름의 Module 을 언로드합니다.
useModules(names...) use_modules(*names) 로드된 모듈의 해석 순서를 활성화하고 변경합니다.
listCatalogs() list_catalogs() 모든 등록된 카탈로그의 이름을 가져옵니다.
listModules() list_modules() 모든 활성화된 모듈의 이름을 가져옵니다.
N/A list_full_modules() (비활성화 포함) 모든 로드된 모듈의 이름을 가져옵니다.
listDatabases() list_databases() 현재 카탈로그의 모든 데이터베이스 이름을 가져옵니다.
listTables() list_tables() 현재 데이터베이스의 모든 테이블과 뷰 이름을 가져옵니다.
listViews() list_views() 현재 데이터베이스의 모든 뷰 이름을 가져옵니다.
listFunctions() list_functions() 이 환경의 모든 함수 이름을 가져옵니다.
N/A list_user_defined_functions() 모든 사용자 정의 함수의 이름을 가져옵니다.
listTemporaryTables() list_temporary_tables() 모든 임시 테이블과 뷰의 이름을 가져옵니다.
listTemporaryViews() list_temporary_views() 모든 임시 뷰의 이름을 가져옵니다.

Statebackend, Checkpoint 및 재시작 전략

TableConfig 에서 key-value 옵션을 설정해 statebackend, 체크포인팅, 재시작 전략을 구성할 수 있습니다. 자세한 내용은 Fault Tolerance, State Backends, Checkpointing 을 참고하세요.

Java:

TableConfig config = tableEnv.getConfig();

// set the restart strategy to "fixed-delay"
config.set("restart-strategy.type", "fixed-delay");
config.set("restart-strategy.fixed-delay.attempts", "3");
config.set("restart-strategy.fixed-delay.delay", "30s");

// set the checkpoint mode to EXACTLY_ONCE
config.set("execution.checkpointing.mode", "EXACTLY_ONCE");
config.set("execution.checkpointing.interval", "3min");

// set the statebackend type to "rocksdb"
config.set("state.backend.type", "rocksdb");

// set the checkpoint directory
config.set("execution.checkpointing.dir", "file:///tmp/checkpoints/");

Scala:

val config = tableEnv.getConfig

// set the restart strategy to "fixed-delay"
config.set("restart-strategy.type", "fixed-delay")
config.set("restart-strategy.fixed-delay.attempts", "3")
config.set("restart-strategy.fixed-delay.delay", "30s")

// set the checkpoint mode to EXACTLY_ONCE
config.set("execution.checkpointing.mode", "EXACTLY_ONCE")
config.set("execution.checkpointing.interval", "3min")

// set the statebackend type to "rocksdb"
config.set("state.backend.type", "rocksdb")

// set the checkpoint directory
config.set("execution.checkpointing.dir", "file:///tmp/checkpoints/")

Python:

config = table_env.get_config()

# set the restart strategy to "fixed-delay"
config.set("restart-strategy.type", "fixed-delay")
config.set("restart-strategy.fixed-delay.attempts", "3")
config.set("restart-strategy.fixed-delay.delay", "30s")

# set the checkpoint mode to EXACTLY_ONCE
config.set("execution.checkpointing.mode", "EXACTLY_ONCE")
config.set("execution.checkpointing.interval", "3min")

# set the statebackend type to "rocksdb"
config.set("state.backend.type", "rocksdb")

# set the checkpoint directory
config.set("execution.checkpointing.dir", "file:///tmp/checkpoints/")

더 알아보기 (Learn more)