모듈
모듈 (Modules)
모듈은 사용자가 Flink의 내장 객체를 확장할 수 있게 합니다(예: Flink 내장 함수처럼 동작하는 함수 정의). 플러그 가능하며, Flink가 몇 가지 사전 빌드된 모듈을 제공하지만 사용자가 직접 작성할 수도 있습니다.
출처: 문서
본문
모듈은 사용자가 Flink의 내장 객체를 확장할 수 있게 합니다(예: Flink 내장 함수처럼 동작하는 함수 정의). 모듈은 플러그 가능하며, Flink가 몇 가지 사전 빌드된 모듈을 제공하지만 사용자가 직접 작성할 수 있습니다.
예를 들어 사용자는 자신의 geo 함수를 정의하고 Flink에 내장 함수로 플러그인하여 Flink SQL과 Table API에서 사용할 수 있습니다. 또 다른 예로 사용자는 기성품 Hive 모듈을 로드하여 Hive 내장 함수를 Flink 내장 함수로 사용할 수 있습니다.
또한 모듈은 내장 table source and sink factories를 제공할 수 있으며, 이는 Java의 Service Provider Interfaces(SPI)에 기반한 Flink의 기본 발견 메커니즘을 비활성화하거나, 해당 카탈로그 없이 임시 테이블의 커넥터가 만들어지는 방식에 영향을 줄 수 있습니다.
모듈 유형 (Module Types)
CoreModule
CoreModule은 Flink의 모든 시스템(내장) 함수를 포함하며 기본적으로 로드되고 활성화됩니다.
HiveModule
HiveModule은 Hive 내장 함수를 Flink의 시스템 함수로 SQL 및 Table API 사용자에게 제공합니다. 모듈 설정에 대한 전체 세부 사항은 Flink의 Hive 문서에서 제공합니다.
사용자 정의 모듈 (User-Defined Module)
사용자는 Module 인터페이스를 구현하여 사용자 지정 모듈을 개발할 수 있습니다. SQL CLI에서 사용자 지정 모듈을 사용하려면 ModuleFactory 인터페이스를 구현하여 모듈과 해당 모듈 팩토리를 모두 개발해야 합니다.
모듈 팩토리는 SQL CLI가 부트스트랩할 때 모듈을 구성하기 위한 속성 집합을 정의합니다. 속성은 발견 서비스(discovery service)로 전달되며, 서비스는 속성을 ModuleFactory와 일치시키고 해당 모듈 인스턴스를 인스턴스화하려 시도합니다.
모듈 수명주기와 해석 순서 (Module Lifecycle and Resolution Order)
모듈은 로드, 활성화, 비활성화, 언로드될 수 있습니다. TableEnvironment가 모듈을 처음 로드할 때 기본적으로 모듈을 활성화합니다. Flink는 여러 모듈을 지원하며 메타데이터를 해석하기 위해 로드 순서를 추적합니다. 또한 Flink는 활성화된 모듈 중에서만 함수를 해석합니다. 예: 두 모듈에 같은 이름의 함수가 두 개 있을 때 세 가지 조건이 있습니다.
- 두 모듈이 모두 활성화되면 Flink는 모듈의 해석 순서에 따라 함수를 해석합니다.
- 하나가 비활성화되면 Flink는 함수를 활성화된 모듈로 해석합니다.
- 둘 다 비활성화되면 Flink는 함수를 해석할 수 없습니다.
사용자는 다른 선언 순서로 모듈을 사용하여 해석 순서를 변경할 수 있습니다. 예: USE MODULES hive, core로 Hive에서 먼저 함수를 찾도록 Flink에 지정할 수 있습니다.
또한 사용자는 모듈을 선언하지 않아 비활성화할 수 있습니다. 예: USE MODULES hive로 core 모듈을 비활성화하도록 Flink에 지정할 수 있습니다(그러나 core 모듈 비활성화는 강력히 권장되지 않습니다). 모듈 비활성화는 언로드를 의미하지 않으며, 사용자가 사용함으로써 다시 활성화할 수 있습니다. 예: USE MODULES core, hive로 core 모듈을 다시 가져와 첫 번째에 둘 수 있습니다. 모듈은 이미 로드되었을 때만 활성화될 수 있습니다. 언로드된 모듈을 사용하면 Exception이 발생합니다. 결국 사용자는 모듈을 언로드할 수 있습니다.
비활성화와 언로드의 차이는 TableEnvironment가 비활성화된 모듈을 계속 유지한다는 점이며, 사용자는 모든 로드된 모듈을 나열하여 비활성화된 모듈을 볼 수 있다는 것입니다.
네임스페이스 (Namespace)
모듈이 제공하는 객체는 Flink의 시스템(내장) 객체의 일부로 간주되므로 네임스페이스가 없습니다.
모듈 로드, 언로드, 사용, 나열 (How to Load, Unload, Use and List Modules)
SQL 사용
다음 SQL 사용 예시입니다.
Java:
EnvironmentSettings settings = EnvironmentSettings.inStreamingMode();
TableEnvironment tableEnv = TableEnvironment.create(settings);
// Show initially loaded and enabled modules
tableEnv.executeSql("SHOW MODULES").print();
// +-------------+
// | module name |
// +-------------+
// | core |
// +-------------+
tableEnv.executeSql("SHOW FULL MODULES").print();
// +-------------+------+
// | module name | used |
// +-------------+------+
// | core | true |
// +-------------+------+
// Load a hive module
tableEnv.executeSql("LOAD MODULE hive WITH ('hive-version' = '...')");
// Show all enabled modules
tableEnv.executeSql("SHOW MODULES").print();
// +-------------+
// | module name |
// +-------------+
// | core |
// | hive |
// +-------------+
// Show all loaded modules with both name and use status
tableEnv.executeSql("SHOW FULL MODULES").print();
// +-------------+------+
// | module name | used |
// +-------------+------+
// | core | true |
// | hive | true |
// +-------------+------+
// Change resolution order
tableEnv.executeSql("USE MODULES hive, core");
tableEnv.executeSql("SHOW MODULES").print();
// +-------------+
// | module name |
// +-------------+
// | hive |
// | core |
// +-------------+
tableEnv.executeSql("SHOW FULL MODULES").print();
// +-------------+------+
// | module name | used |
// +-------------+------+
// | hive | true |
// | core | true |
// +-------------+------+
// Disable core module
tableEnv.executeSql("USE MODULES hive");
tableEnv.executeSql("SHOW MODULES").print();
// +-------------+
// | module name |
// +-------------+
// | hive |
// +-------------+
tableEnv.executeSql("SHOW FULL MODULES").print();
// +-------------+-------+
// | module name | used |
// +-------------+-------+
// | hive | true |
// | core | false |
// +-------------+-------+
// Unload hive module
tableEnv.executeSql("UNLOAD MODULE hive");
tableEnv.executeSql("SHOW MODULES").print();
// Empty set
tableEnv.executeSql("SHOW FULL MODULES").print();
// +-------------+-------+
// | module name | used |
// +-------------+-------+
// | hive | false |
// +-------------+-------+
Scala:
val settings = EnvironmentSettings.inStreamingMode()
val tableEnv = TableEnvironment.create(setting)
// Show initially loaded and enabled modules
tableEnv.executeSql("SHOW MODULES").print()
// +-------------+
// | module name |
// +-------------+
// | core |
// +-------------+
tableEnv.executeSql("SHOW FULL MODULES").print()
// +-------------+------+
// | module name | used |
// +-------------+------+
// | core | true |
// +-------------+------+
// Load a hive module
tableEnv.executeSql("LOAD MODULE hive WITH ('hive-version' = '...')")
// Show all enabled modules
tableEnv.executeSql("SHOW MODULES").print()
// +-------------+
// | module name |
// +-------------+
// | core |
// | hive |
// +-------------+
// Show all loaded modules with both name and use status
tableEnv.executeSql("SHOW FULL MODULES")
// +-------------+------+
// | module name | used |
// +-------------+------+
// | core | true |
// | hive | true |
// +-------------+------+
// Change resolution order
tableEnv.executeSql("USE MODULES hive, core")
tableEnv.executeSql("SHOW MODULES").print()
// +-------------+
// | module name |
// +-------------+
// | hive |
// | core |
// +-------------+
tableEnv.executeSql("SHOW FULL MODULES").print()
// +-------------+------+
// | module name | used |
// +-------------+------+
// | hive | true |
// | core | true |
// +-------------+------+
// Disable core module
tableEnv.executeSql("USE MODULES hive")
tableEnv.executeSql("SHOW MODULES").print()
// +-------------+
// | module name |
// +-------------+
// | hive |
// +-------------+
tableEnv.executeSql("SHOW FULL MODULES").print()
// +-------------+-------+
// | module name | used |
// +-------------+-------+
// | hive | true |
// | core | false |
// +-------------+-------+
// Unload hive module
tableEnv.executeSql("UNLOAD MODULE hive")
tableEnv.executeSql("SHOW MODULES").print()
// Empty set
tableEnv.executeSql("SHOW FULL MODULES").print()
// +-------------+-------+
// | module name | used |
// +-------------+-------+
// | hive | false |
// +-------------+-------+
Python:
from pyflink.table import *
# environment configuration
settings = EnvironmentSettings.inStreamingMode()
t_env = TableEnvironment.create(settings)
# Show initially loaded and enabled modules
t_env.execute_sql("SHOW MODULES").print()
# +-------------+
# | module name |
# +-------------+
# | core |
# +-------------+
t_env.execute_sql("SHOW FULL MODULES").print()
# +-------------+------+
# | module name | used |
# +-------------+------+
# | core | true |
# +-------------+------+
# Load a hive module
t_env.execute_sql("LOAD MODULE hive WITH ('hive-version' = '...')")
# Show all enabled modules
t_env.execute_sql("SHOW MODULES").print()
# +-------------+
# | module name |
# +-------------+
# | core |
# | hive |
# +-------------+
# Show all loaded modules with both name and use status
t_env.execute_sql("SHOW FULL MODULES").print()
# +-------------+------+
# | module name | used |
# +-------------+------+
# | core | true |
# | hive | true |
# +-------------+------+
# Change resolution order
t_env.execute_sql("USE MODULES hive, core")
t_env.execute_sql("SHOW MODULES").print()
# +-------------+
# | module name |
# +-------------+
# | hive |
# | core |
# +-------------+
t_env.execute_sql("SHOW FULL MODULES").print()
# +-------------+------+
# | module name | used |
# +-------------+------+
# | hive | true |
# | core | true |
# +-------------+------+
# Disable core module
t_env.execute_sql("USE MODULES hive")
t_env.execute_sql("SHOW MODULES").print()
# +-------------+
# | module name |
# +-------------+
# | hive |
# +-------------+
t_env.execute_sql("SHOW FULL MODULES").print()
# +-------------+-------+
# | module name | used |
# +-------------+-------+
# | hive | true |
# | core | false |
# +-------------+-------+
# Unload hive module
t_env.execute_sql("UNLOAD MODULE hive")
t_env.execute_sql("SHOW MODULES").print()
# Empty set
t_env.execute_sql("SHOW FULL MODULES").print()
# +-------------+-------+
# | module name | used |
# +-------------+-------+
# | hive | false |
# +-------------+-------+
SQL Client:
-- Show initially loaded and enabled modules
Flink SQL> SHOW MODULES;
+-------------+
| module name |
+-------------+
| core |
+-------------+
1 row in set
Flink SQL> SHOW FULL MODULES;
+-------------+------+
| module name | used |
+-------------+------+
| core | true |
+-------------+------+
1 row in set
-- Load a hive module
Flink SQL> LOAD MODULE hive WITH ('hive-version' = '...');
-- Show all enabled modules
Flink SQL> SHOW MODULES;
+-------------+
| module name |
+-------------+
| core |
| hive |
+-------------+
2 rows in set
-- Show all loaded modules with both name and use status
Flink SQL> SHOW FULL MODULES;
+-------------+------+
| module name | used |
+-------------+------+
| core | true |
| hive | true |
+-------------+------+
2 rows in set
-- Change resolution order
Flink SQL> USE MODULES hive, core ;
Flink SQL> SHOW MODULES;
+-------------+
| module name |
+-------------+
| hive |
| core |
+-------------+
2 rows in set
Flink SQL> SHOW FULL MODULES;
+-------------+------+
| module name | used |
+-------------+------+
| hive | true |
| core | true |
+-------------+------+
2 rows in set
-- Unload hive module
Flink SQL> UNLOAD MODULE hive;
Flink SQL> SHOW MODULES;
Empty set
Flink SQL> SHOW FULL MODULES;
+-------------+-------+
| module name | used |
+-------------+-------+
| hive | false |
+-------------+-------+
1 row in set
YAML:
YAML로 정의된 모든 모듈은 유형을 지정하는 type 속성을 제공해야 합니다. 다음 유형이 기본으로 지원됩니다.
| 모듈 | 타입 값 |
|---|---|
| CoreModule | core |
| HiveModule | hive |
modules:
- name: core
type: core
- name: hive
type: hive
경고: SQL을 사용할 때 모듈 이름은 모듈 발견을 수행하는 데 사용됩니다. 단순 식별자로 파싱되며 대소문자를 구분합니다.
Java, Scala 또는 Python 사용
사용자는 Java, Scala 또는 Python을 사용해 프로그래밍 방식으로 모듈을 로드/언로드/사용/나열할 수 있습니다.
Java:
EnvironmentSettings settings = EnvironmentSettings.inStreamingMode();
TableEnvironment tableEnv = TableEnvironment.create(settings);
// Show initially loaded and enabled modules
tableEnv.listModules();
// +-------------+
// | module name |
// +-------------+
// | core |
// +-------------+
tableEnv.listFullModules();
// +-------------+------+
// | module name | used |
// +-------------+------+
// | core | true |
// +-------------+------+
// Load a hive module
tableEnv.loadModule("hive", new HiveModule());
// Show all enabled modules
tableEnv.listModules();
// +-------------+
// | module name |
// +-------------+
// | core |
// | hive |
// +-------------+
// Show all loaded modules with both name and use status
tableEnv.listFullModules();
// +-------------+------+
// | module name | used |
// +-------------+------+
// | core | true |
// | hive | true |
// +-------------+------+
// Change resolution order
tableEnv.useModules("hive", "core");
tableEnv.listModules();
// +-------------+
// | module name |
// +-------------+
// | hive |
// | core |
// +-------------+
tableEnv.listFullModules();
// +-------------+------+
// | module name | used |
// +-------------+------+
// | hive | true |
// | core | true |
// +-------------+------+
// Disable core module
tableEnv.useModules("hive");
tableEnv.listModules();
// +-------------+
// | module name |
// +-------------+
// | hive |
// +-------------+
tableEnv.listFullModules();
// +-------------+-------+
// | module name | used |
// +-------------+-------+
// | hive | true |
// | core | false |
// +-------------+-------+
// Unload hive module
tableEnv.unloadModule("hive");
tableEnv.listModules();
// Empty set
tableEnv.listFullModules();
// +-------------+-------+
// | module name | used |
// +-------------+-------+
// | hive | false |
// +-------------+-------+
Scala:
val settings = EnvironmentSettings.inStreamingMode()
val tableEnv = TableEnvironment.create(setting)
// Show initially loaded and enabled modules
tableEnv.listModules()
// +-------------+
// | module name |
// +-------------+
// | core |
// +-------------+
tableEnv.listFullModules()
// +-------------+------+
// | module name | used |
// +-------------+------+
// | core | true |
// +-------------+------+
// Load a hive module
tableEnv.loadModule("hive", new HiveModule())
// Show all enabled modules
tableEnv.listModules()
// +-------------+
// | module name |
// +-------------+
// | core |
// | hive |
// +-------------+
// Show all loaded modules with both name and use status
tableEnv.listFullModules()
// +-------------+------+
// | module name | used |
// +-------------+------+
// | core | true |
// | hive | true |
// +-------------+------+
// Change resolution order
tableEnv.useModules("hive", "core")
tableEnv.listModules()
// +-------------+
// | module name |
// +-------------+
// | hive |
// | core |
// +-------------+
tableEnv.listFullModules()
// +-------------+------+
// | module name | used |
// +-------------+------+
// | hive | true |
// | core | true |
// +-------------+------+
// Disable core module
tableEnv.useModules("hive")
tableEnv.listModules()
// +-------------+
// | module name |
// +-------------+
// | hive |
// +-------------+
tableEnv.listFullModules()
// +-------------+-------+
// | module name | used |
// +-------------+-------+
// | hive | true |
// | core | false |
// +-------------+-------+
// Unload hive module
tableEnv.unloadModule("hive")
tableEnv.listModules()
// Empty set
tableEnv.listFullModules()
// +-------------+-------+
// | module name | used |
// +-------------+-------+
// | hive | false |
// +-------------+-------+
Python:
from pyflink.table import *
# environment configuration
settings = EnvironmentSettings.inStreamingMode()
t_env = TableEnvironment.create(settings)
# Show initially loaded and enabled modules
t_env.list_modules()
# +-------------+
# | module name |
# +-------------+
# | core |
# +-------------+
t_env.list_full_modules()
# +-------------+------+
# | module name | used |
# +-------------+------+
# | core | true |
# +-------------+------+
# Load a hive module
t_env.load_module("hive", HiveModule())
# Show all enabled modules
t_env.list_modules()
# +-------------+
# | module name |
# +-------------+
# | core |
# | hive |
# +-------------+
# Show all loaded modules with both name and use status
t_env.list_full_modules()
# +-------------+------+
# | module name | used |
# +-------------+------+
# | core | true |
# | hive | true |
# +-------------+------+
# Change resolution order
t_env.use_modules("hive", "core")
t_env.list_modules()
# +-------------+
# | module name |
# +-------------+
# | hive |
# | core |
# +-------------+
t_env.list_full_modules()
# +-------------+------+
# | module name | used |
# +-------------+------+
# | hive | true |
# | core | true |
# +-------------+------+
# Disable core module
t_env.use_modules("hive")
t_env.list_modules()
# +-------------+
# | module name |
# +-------------+
# | hive |
# +-------------+
t_env.list_full_modules()
# +-------------+-------+
# | module name | used |
# +-------------+-------+
# | hive | true |
# | core | false |
# +-------------+-------+
# Unload hive module
t_env.unload_module("hive")
t_env.list_modules()
# Empty set
t_env.list_full_modules()
# +-------------+-------+
# | module name | used |
# +-------------+-------+
# | hive | false |
# +-------------+-------+