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")
사용 가능한 모든 옵션은 Configuration 및 Python Configuration 을 참고하세요.
카탈로그 API
이 API 는 카탈로그와 모듈에 접근하는 데 사용됩니다. 자세한 내용은 Modules 및 Catalogs 를 참고하세요.
| 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/")