카탈로그

카탈로그 (Catalogs)

카탈로그는 데이터베이스, 테이블, 파티션, 뷰, 함수와 같은 메타데이터와, 데이터베이스나 다른 외부 시스템에 저장된 데이터에 접근하는 데 필요한 정보를 제공해요.

출처: 문서

본문

데이터 처리에서 가장 중요한 측면 중 하나는 메타데이터 관리예요. 그것은 임시 테이블 또는 테이블 환경에 등록된 UDF 같은 일시적(transient) 메타데이터일 수 있어요. 또는 Hive Metastore 같은 영구적(permanent) 메타데이터일 수 있어요. 카탈로그는 메타데이터를 관리하고 Table API와 SQL 쿼리에서 접근할 수 있게 하는 통합 API를 제공해요.

카탈로그는 사용자가 자신의 데이터 시스템에서 기존 메타데이터를 참조하고, 이를 Flink의 해당 메타데이터에 자동으로 매핑할 수 있게 해줘요. 예를 들어 Flink는 JDBC 테이블을 Flink 테이블에 자동으로 매핑할 수 있고, 사용자는 Flink에서 DDL을 수동으로 다시 작성할 필요가 없어요. 카탈로그는 사용자의 기존 시스템으로 Flink를 시작하는 데 필요한 단계를 크게 단순화하고, 사용자 경험을 크게 향상시켜요.

카탈로그 유형 (Catalog Types)

GenericInMemoryCatalog

GenericInMemoryCatalog는 카탈로그의 인메모리 구현이에요. 모든 객체는 세션의 수명 동안에만 사용할 수 있어요.

JdbcCatalog

JdbcCatalog는 사용자가 JDBC 프로토콜을 통해 Flink를 관계형 데이터베이스에 연결할 수 있게 해줘요. 현재 JdbcCatalog의 구현은 Postgres Catalog와 MySQL Catalog 둘뿐이에요. 카탈로그 설정에 대한 자세한 내용은 JdbcCatalog 문서를 참조하세요.

HiveCatalog

HiveCatalog는 두 가지 목적을 제공해요; 순수 Flink 메타데이터를 위한 영구 저장소, 그리고 기존 Hive 메타데이터를 읽고 쓰기 위한 인터페이스. Flink의 Hive 문서는 카탈로그 설정과 기존 Hive 설치와의 인터페이싱에 대한 전체 세부 사항을 제공해요.

Hive Metastore는 모든 메타 객체 이름을 소문자로 저장해요. 이것은 대소문자를 구분하는 GenericInMemoryCatalog와는 달라요.

사용자 정의 카탈로그 (User-Defined Catalog)

카탈로그는 플러그형이며, 사용자는 Catalog 인터페이스를 구현해 커스텀 카탈로그를 개발할 수 있어요.

Flink SQL에서 커스텀 카탈로그를 사용하려면, 사용자는 CatalogFactory 인터페이스를 구현해 해당 카탈로그 팩토리를 구현해야 해요. 팩토리는 Java의 Service Provider Interfaces (SPI)를 사용해 발견돼요. 이 인터페이스를 구현하는 클래스는 JAR 파일의 META_INF/services/org.apache.flink.table.factories.Factory에 추가될 수 있어요. 제공된 팩토리 식별자는 SQL CREATE CATALOG DDL 문의 필수 type 속성과 일치하는 데 사용돼요.

Flink v1.16부터 TableEnvironment는 테이블 프로그램, SQL Client 및 SQL Gateway에서 일관된 클래스 로딩 동작을 위해 사용자 클래스 로더(user class loader)를 도입했어요. 사용자 클래스로더는 ADD JAR 또는 CREATE FUNCTION .. USING JAR .. 문으로 추가된 jar 같은 모든 사용자 jar를 관리해요. 사용자 정의 카탈로그는 클래스를 로드하기 위해 Thread.currentThread().getContextClassLoader()를 사용자 클래스 로더로 대체해야 해요. 그렇지 않으면 ClassNotFoundException이 발생할 수 있어요. 사용자 클래스 로더는 CatalogFactory.Context#getClassLoader로 접근할 수 있어요.

시간 여행(time travel) 지원을 위한 Catalog 인터페이스

버전 1.18부터 Flink 프레임워크는 테이블의 과거 데이터를 질의하는 시간 여행(time travel)을 지원해요. 테이블의 과거 데이터를 질의하려면, 사용자는 테이블이 속한 카탈로그에 getTable(ObjectPath tablePath, long timestamp) 메서드를 구현해야 해요.

public class MyCatalogSupportTimeTravel implements Catalog {

    @Override
    public CatalogBaseTable getTable(ObjectPath tablePath, long timestamp)
            throws TableNotExistException {
        // Build a schema corresponding to the specific time point.
        Schema schema = buildSchema(timestamp);
        // Set parameters to read data at the corresponding time point.
        Map<String, String> options = buildOptions(timestamp);
        // Build CatalogTable
        CatalogTable catalogTable =
                CatalogTable.newBuilder()
                        .schema(schema)
                        .comment("")
                        .partitionKeys(Collections.emptyList())
                        .options(options)
                        .snapshot(timestamp)
                        .build();
        return catalogTable;
    }
}

public class MyDynamicTableFactory implements DynamicTableSourceFactory {
    @Override
    public DynamicTableSource createDynamicTableSource(Context context) {
        final ReadableConfig configuration =
                Configuration.fromMap(context.getCatalogTable().getOptions());

        // Get snapshot from CatalogTable
        final Optional<Long> snapshot = context.getCatalogTable().getSnapshot();

        // Build DynamicTableSource using snapshot options.
        final DynamicTableSource dynamicTableSource = buildDynamicSource(configuration, snapshot);

        return dynamicTableSource;
    }
}

SQL DDL 사용

Table API와 SQL 모두에서 SQL DDL을 사용해 카탈로그에 테이블을 생성할 수 있어요.

TableEnvironment tableEnv = ...;

// Create a HiveCatalog 
Catalog catalog = new HiveCatalog("myhive", null, "");

// Register the catalog
tableEnv.registerCatalog("myhive", catalog);

// Create a catalog database
tableEnv.executeSql("CREATE DATABASE mydb WITH (...)");

// Create a catalog table
tableEnv.executeSql("CREATE TABLE mytable (name STRING, age INT) WITH (...)");

tableEnv.listTables(); // should return the tables in current catalog and database.
val tableEnv = ...

// Create a HiveCatalog 
val catalog = new HiveCatalog("myhive", null, "")

// Register the catalog
tableEnv.registerCatalog("myhive", catalog)

// Create a catalog database
tableEnv.executeSql("CREATE DATABASE mydb WITH (...)")

// Create a catalog table
tableEnv.executeSql("CREATE TABLE mytable (name STRING, age INT) WITH (...)")

tableEnv.listTables() // should return the tables in current catalog and database.
from pyflink.table.catalog import HiveCatalog

# Create a HiveCatalog
catalog = HiveCatalog("myhive", None, "")

# Register the catalog
t_env.register_catalog("myhive", catalog)

# Create a catalog database
t_env.execute_sql("CREATE DATABASE mydb WITH (...)")

# Create a catalog table
t_env.execute_sql("CREATE TABLE mytable (name STRING, age INT) WITH (...)")

# should return the tables in current catalog and database.
t_env.list_tables()
// the catalog should have been registered via yaml file
Flink SQL> CREATE DATABASE mydb WITH (...);

Flink SQL> CREATE TABLE mytable (name STRING, age INT) WITH (...);

Flink SQL> SHOW TABLES;
mytable

자세한 정보는 Flink SQL CREATE DDL을 확인하세요.

Java, Scala 또는 Python 사용

사용자는 Java, Scala 또는 Python을 사용해 프로그래밍 방식으로 카탈로그 테이블을 생성할 수 있어요.

import org.apache.flink.table.api.*;
import org.apache.flink.table.catalog.*;
import org.apache.flink.table.catalog.hive.HiveCatalog;

TableEnvironment tableEnv = TableEnvironment.create(EnvironmentSettings.inStreamingMode());

// Create a HiveCatalog 
Catalog catalog = new HiveCatalog("myhive", null, "");

// Register the catalog
tableEnv.registerCatalog("myhive", catalog);

// Create a catalog database 
catalog.createDatabase("mydb", new CatalogDatabaseImpl(...));

// Create a catalog table
final Schema schema = Schema.newBuilder()
    .column("name", DataTypes.STRING())
    .column("age", DataTypes.INT())
    .build();

tableEnv.createTable("myhive.mydb.mytable", TableDescriptor.forConnector("kafka")
    .schema(schema)
    // …
    .build());

List<String> tables = catalog.listTables("mydb"); // tables should contain "mytable"
from pyflink.table import *
from pyflink.table.catalog import HiveCatalog, CatalogDatabase, ObjectPath, CatalogBaseTable

settings = EnvironmentSettings.in_batch_mode()
t_env = TableEnvironment.create(settings)

# Create a HiveCatalog
catalog = HiveCatalog("myhive", None, "")

# Register the catalog
t_env.register_catalog("myhive", catalog)

# Create a catalog database
database = CatalogDatabase.create_instance({"k1": "v1"}, None)
catalog.create_database("mydb", database)

# Create a catalog table
schema = Schema.new_builder() \
    .column("name", DataTypes.STRING()) \
    .column("age", DataTypes.INT()) \
    .build()
    
catalog_table = t_env.create_table("myhive.mydb.mytable", TableDescriptor.for_connector("kafka")
    .schema(schema)
    # …
    .build())

# tables should contain "mytable"
tables = catalog.list_tables("mydb")

카탈로그 API (Catalog API)

참고: 여기에는 카탈로그 프로그램 API만 나열돼요. 사용자는 SQL DDL로 동일한 많은 기능을 달성할 수 있어요. 자세한 DDL 정보는 SQL CREATE DDL을 참조하세요.

데이터베이스 연산 (Database operations)

// create database
catalog.createDatabase("mydb", new CatalogDatabaseImpl(...), false);

// drop database
catalog.dropDatabase("mydb", false);

// alter database
catalog.alterDatabase("mydb", new CatalogDatabaseImpl(...), false);

// get database
catalog.getDatabase("mydb");

// check if a database exist
catalog.databaseExists("mydb");

// list databases in a catalog
catalog.listDatabases();
from pyflink.table.catalog import CatalogDatabase

# create database
catalog_database = CatalogDatabase.create_instance({"k1": "v1"}, None)
catalog.create_database("mydb", catalog_database, False)

# drop database
catalog.drop_database("mydb", False)

# alter database
catalog.alter_database("mydb", catalog_database, False)

카탈로그는 데이터베이스 메타데이터를 관리하는 것 외에도 테이블뿐 아니라 파티션, 뷰, 함수 같은 객체들을 관리하는 다양한 API를 제공해요. 자세한 내용은 Java Document(Scala는 Java API를 기반) 및 Python API 문서를 참고하세요.

더 알아보기 (Learn more)