JDBC SQL 커넥터

JDBC SQL 커넥터 (JDBC SQL Connector)

스캔 소스: 유한(Bounded) / 룩업 소스: 동기 모드 / 싱크: 배치(Batch) / 싱크: 스트리밍 Append & Upsert 모드

JDBC 커넥터는 JDBC 드라이버가 있는 모든 관계형 데이터베이스에서 데이터를 읽고 쓸 수 있게 해줍니다. 이 문서는 관계형 데이터베이스에 대해 SQL 쿼리를 실행하도록 JDBC 커넥터를 설정하는 방법을 설명합니다.

JDBC 싱크는 DDL에 기본 키가 정의되어 있으면 외부 시스템과 UPDATE/DELETE 메시지를 교환하는 upsert 모드로 동작하고, 그렇지 않으면 append 모드로 동작하며 UPDATE/DELETE 메시지 소비를 지원하지 않습니다.

출처: 문서

본문

의존성 (Dependencies)

현재 Flink 2.3 버전용 커넥터는 아직 없습니다.

JDBC 커넥터는 바이너리 배포의 일부가 아닙니다. 클러스터 실행을 위해 연결하는 방법은 여기를 참고하세요.

지정된 데이터베이스에 연결하려면 드라이버 의존성도 필요합니다. 현재 지원되는 드라이버는 다음과 같습니다:

드라이버 Group Id Artifact Id JAR
MySQL mysql mysql-connector-java Download
Oracle com.oracle.database.jdbc ojdbc8 Download
PostgreSQL org.postgresql postgresql Download
Derby org.apache.derby derby Download
SQL Server com.microsoft.sqlserver mssql-jdbc Download

JDBC 커넥터와 드라이버는 Flink의 바이너리 배포의 일부가 아닙니다. 클러스터 실행을 위해 연결하는 방법은 여기를 참고하세요.

JDBC 테이블 생성 방법 (How to create a JDBC table)

JDBC 테이블은 다음과 같이 정의할 수 있습니다:

-- register a MySQL table 'users' in Flink SQL
CREATE TABLE MyUserTable (
  id BIGINT,
  name STRING,
  age INT,
  status BOOLEAN,
  PRIMARY KEY (id) NOT ENFORCED
) WITH (
   'connector' = 'jdbc',
   'url' = 'jdbc:mysql://localhost:3306/mydatabase',
   'table-name' = 'users'
);

-- write data into the JDBC table from the other table "T"
INSERT INTO MyUserTable
SELECT id, name, age, status FROM T;

-- scan data from the JDBC table
SELECT id, name, age, status FROM MyUserTable;

-- temporal join the JDBC table as a dimension table
SELECT * FROM myTopic
LEFT JOIN MyUserTable FOR SYSTEM_TIME AS OF myTopic.proctime
ON myTopic.key = MyUserTable.id;

커넥터 옵션 (Connector Options)

옵션 필수 Forwarded 기본값 타입 설명
connector required no (none) String 사용할 커넥터. 여기서는 'jdbc'여야 합니다.
url required yes (none) String JDBC 데이터베이스 url.
table-name required yes (none) String 연결할 JDBC 테이블 이름.
driver optional yes (none) String 이 URL에 연결하기 위해 사용할 JDBC 드라이버의 클래스 이름. 설정하지 않으면 URL에서 자동으로 파생됩니다.
username optional yes (none) String JDBC 사용자 이름. 'username''password' 중 하나라도 지정되면 둘 다 지정해야 합니다.
password optional yes (none) String JDBC 비밀번호.
connection.max-retry-timeout optional yes 60s Duration 재시도 사이의 최대 시간 초과. 시간 초과는 초 단위여야 하며 1초보다 작을 수 없습니다.
scan.partition.column optional no (none) String 입력 파티셔닝에 사용되는 컬럼 이름. 자세한 내용은 아래 Partitioned Scan 섹션을 참고하세요.
scan.partition.num optional no (none) Integer 파티션 수.
scan.partition.lower-bound optional no (none) Integer 첫 번째 파티션의 가장 작은 값.
scan.partition.upper-bound optional no (none) Integer 마지막 파티션의 가장 큰 값.
scan.fetch-size optional yes 0 Integer 읽을 때 왕복당 데이터베이스에서 가져와야 하는 행 수. 지정된 값이 0이면 힌트가 무시됩니다.
scan.auto-commit optional yes true Boolean JDBC 드라이버에 auto-commit 플래그를 설정합니다. 각 문이 트랜잭션으로 자동 커밋되는지 여부를 결정합니다. 일부 JDBC 드라이버, 특히 Postgres는 결과를 스트리밍하기 위해 이것을 false로 설정해야 할 수 있습니다.
lookup.cache optional yes NONE Enum (NONE, PARTIAL) 룩업 테이블의 캐시 전략. 현재 NONE(캐시 없음)과 PARTIAL(외부 데이터베이스에서 룩업 연산 시 항목 캐싱)을 지원합니다.
lookup.partial-cache.max-rows optional yes (none) Long 룩업 캐시의 최대 행 수. 이 값을 초과하면 가장 오래된 행이 만료됩니다. 이 옵션을 사용하려면 "lookup.cache"가 "PARTIAL"이어야 합니다.
lookup.partial-cache.expire-after-write optional yes (none) Duration 캐시에 쓴 후 룩업 캐시의 각 행의 최대 수명. 이 옵션을 사용하려면 "lookup.cache"가 "PARTIAL"이어야 합니다.
lookup.partial-cache.expire-after-access optional yes (none) Duration 캐시의 항목에 접근한 후 룩업 캐시의 각 행의 최대 수명. 이 옵션을 사용하려면 "lookup.cache"가 "PARTIAL"이어야 합니다.
lookup.partial-cache.cache-missing-key optional yes true Boolean 룩업 키가 테이블의 어떤 행과도 일치하지 않을 때 빈 값을 캐시에 저장할지 여부. 이 옵션을 사용하려면 "lookup.cache"가 "PARTIAL"이어야 합니다.
lookup.max-retries optional yes 3 Integer 룩업 데이터베이스가 실패한 경우 최대 재시도 횟수.
sink.buffer-flush.max-rows optional yes 100 Integer 플러시 전 버퍼링되는 레코드의 최대 크기. 비활성화하려면 0으로 설정할 수 있습니다.
sink.buffer-flush.interval optional yes 1s Duration 플러시 간격(밀리초). 이 시간이 지나면 비동기 스레드가 데이터를 플러시합니다. 비활성화하려면 '0'으로 설정할 수 있습니다. 'sink.buffer-flush.max-rows''0'으로 설정하고 플러시 간격을 설정하면 버퍼링된 작업을 완전히 비동기 처리할 수 있음에 유의하세요.
sink.max-retries optional yes 3 Integer 데이터베이스에 레코드 쓰기가 실패한 경우 최대 재시도 횟수.
sink.parallelism optional no (none) Integer JDBC 싱크 연산자의 병렬도를 정의합니다. 기본적으로 병렬도는 업스트림 체인 연산자의 병렬도와 동일하게 프레임워크가 결정합니다.

더 이상 사용되지 않는 옵션 (Deprecated Options)

이러한 더 이상 사용되지 않는 옵션은 위에 나열된 새 옵션으로 교체되었으며 결국 제거될 것입니다. 새 옵션을 먼저 사용하는 것을 고려하세요.

옵션 필수 Forwarded 기본값 타입 설명
lookup.cache.max-rows optional yes (none) Integer "lookup.cache" = "PARTIAL"로 설정하고 대신 "lookup.partial-cache.max-rows"를 사용하세요.
lookup.cache.ttl optional yes (none) Duration "lookup.cache" = "PARTIAL"로 설정하고 대신 "lookup.partial-cache.expire-after-write"를 사용하세요.
lookup.cache.caching-missing-key optional yes true Boolean "lookup.cache" = "PARTIAL"로 설정하고 대신 "lookup.partial-cache.cache-missing-key"를 사용하세요.

기능 (Features)

키 처리 (Key handling)

Flink는 외부 데이터베이스에 데이터를 쓸 때 DDL에 정의된 기본 키를 사용합니다. 기본 키가 정의되면 커넥터는 upsert 모드로 동작하고, 그렇지 않으면 append 모드로 동작합니다.

upsert 모드에서 Flink는 기본 키에 따라 새 행을 삽입하거나 기존 행을 갱신하며, 이로써 멱등성을 보장할 수 있습니다. 출력 결과가 예상대로 나오도록 보장하려면 테이블에 기본 키를 정의하고, 기본 키가 기본 데이터베이스 테이블의 고유 키 집합 중 하나이거나 기본 키인지 확인하는 것을 권장합니다. append 모드에서 Flink는 모든 레코드를 INSERT 메시지로 해석하며, 기본 데이터베이스에서 기본 키나 고유 제약 위반이 발생하면 INSERT 연산이 실패할 수 있습니다.

PRIMARY KEY 문법에 대한 자세한 내용은 CREATE TABLE DDL을 참고하세요.

파티션 스캔 (Partitioned Scan)

병렬 Source 작업 인스턴스에서 데이터 읽기를 가속화하기 위해 Flink는 JDBC 테이블에 대한 파티션 스캔 기능을 제공합니다.

다음 스캔 파티션 옵션 중 하나라도 지정되면 모두 지정해야 합니다. 이들은 여러 태스크에서 병렬로 읽을 때 테이블을 파티셔닝하는 방법을 설명합니다. scan.partition.column은 해당 테이블의 숫자, 날짜 또는 타임스탬프 컬럼이어야 합니다. scan.partition.lower-boundscan.partition.upper-bound는 파티션 보폭(stride)을 결정하고 테이블의 행을 필터링하는 데 사용됩니다. 배치 작업이라면 Flink 작업 제출 전에 최대·최소값을 먼저 얻을 수도 있습니다.

  • scan.partition.column: 입력 파티셔닝에 사용되는 컬럼 이름.
  • scan.partition.num: 파티션 수.
  • scan.partition.lower-bound: 첫 번째 파티션의 가장 작은 값.
  • scan.partition.upper-bound: 마지막 파티션의 가장 큰 값.

룩업 캐시 (Lookup Cache)

JDBC 커넥터는 시간 조인에서 룩업 소스(일명 차원 테이블)로 사용될 수 있습니다. 현재 동기 룩업 모드만 지원됩니다.

기본적으로 룩업 캐시는 활성화되지 않습니다. lookup.cachePARTIAL로 설정해 활성화할 수 있습니다.

룩업 캐시는 JDBC 커넥터의 시간 조인 성능을 향상시키는 데 사용됩니다. 기본적으로 룩업 캐시가 활성화되지 않으므로 모든 요청이 외부 데이터베이스로 전송됩니다. 룩업 캐시가 활성화되면 각 프로세스(즉 TaskManager)가 캐시를 보유합니다. Flink는 먼저 캐시를 조회하고, 캐시 미스일 때만 외부 데이터베이스에 요청을 보내며, 반환된 행으로 캐시를 갱신합니다. 캐시가 최대 캐시 행 수 lookup.partial-cache.max-rows에 도달하거나 행이 lookup.partial-cache.expire-after-write 또는 lookup.partial-cache.expire-after-access가 지정한 최대 수명을 초과하면 캐시의 가장 오래된 행이 만료됩니다. 캐시된 행은 최신이 아닐 수 있으며, 사용자는 만료 옵션을 더 작은 값으로 조정해 더 신선한 데이터를 얻을 수 있지만 데이터베이스에 보내는 요청 수가 늘어날 수 있습니다. 따라서 이는 처리량과 정확성 사이의 균형입니다.

기본적으로 Flink는 기본 키에 대한 빈 쿼리 결과를 캐시합니다. lookup.partial-cache.cache-missing-key를 false로 설정해 이 동작을 전환할 수 있습니다.

멱등 쓰기 (Idempotent Writes)

JDBC 싱크는 DDL에 기본 키가 정의되어 있으면 일반 INSERT 문이 아닌 upsert 의미론을 사용합니다. upsert 의미론은 기본 데이터베이스에서 고유 제약 위반이 있을 때 새 행을 원자적으로 추가하거나 기존 행을 갱신하는 것을 말하며, 이는 멱등성을 제공합니다.

실패가 있으면 Flink 작업은 복구되어 마지막 성공한 체크포인트부터 다시 처리하며, 이로 인해 복구 중 메시지가 재처리될 수 있습니다. 레코드를 재처리해야 하는 경우 제약 위반이나 중복 데이터를 피하는 데 도움이 되므로 upsert 모드를 강력히 권장합니다.

실패 복구 외에도 소스 토픽은 시간에 따라 같은 기본 키를 가진 여러 레코드를 자연스럽게 포함할 수 있으므로 upsert가 바람직합니다.

upsert에 표준 문법이 없으므로 다음 표는 사용되는 데이터베이스별 DML을 설명합니다.

데이터베이스 Upsert 문법
MySQL INSERT .. ON DUPLICATE KEY UPDATE ..
Oracle MERGE INTO .. USING (..) ON (..) WHEN MATCHED THEN UPDATE SET (..) WHEN NOT MATCHED THEN INSERT (..) VALUES (..)
PostgreSQL INSERT .. ON CONFLICT .. DO UPDATE SET ..
MS SQL Server MERGE INTO .. USING (..) ON (..) WHEN MATCHED THEN UPDATE SET (..) WHEN NOT MATCHED THEN INSERT (..) VALUES (..)

JDBC 카탈로그 (JDBC Catalog)

JdbcCatalog는 사용자가 JDBC 프로토콜을 통해 Flink를 관계형 데이터베이스에 연결할 수 있게 해줍니다.

현재 Postgres Catalog와 MySQL Catalog 두 가지 JDBC 카탈로그 구현이 있습니다. 이들은 다음 카탈로그 메서드를 지원합니다. 다른 메서드는 현재 지원되지 않습니다.

// The supported methods by Postgres & MySQL Catalog.
databaseExists(String databaseName);
listDatabases();
getDatabase(String databaseName);
listTables(String databaseName);
getTable(ObjectPath tablePath);
tableExists(ObjectPath tablePath);

다른 Catalog 메서드는 현재 지원되지 않습니다.

JDBC 카탈로그 사용 (Usage of JDBC Catalog)

이 섹션은 주로 Postgres Catalog 또는 MySQL Catalog를 생성하고 사용하는 방법을 설명합니다. JDBC 커넥터와 해당 드라이버를 설정하는 방법은 의존성 섹션을 참고하세요.

JDBC 카탈로그는 다음 옵션을 지원합니다:

  • name: 필수, 카탈로그 이름.
  • default-database: 필수, 연결할 기본 데이터베이스.
  • username: 필수, Postgres/MySQL 계정의 사용자 이름.
  • password: 필수, 계정의 비밀번호.
  • base-url: 필수(데이터베이스 이름을 포함하지 않아야 함)
    • Postgres Catalog의 경우 "jdbc:postgresql://<ip>:<port>"
    • MySQL Catalog의 경우 "jdbc:mysql://<ip>:<port>"

SQL

CREATE CATALOG my_catalog WITH(
    'type' = 'jdbc',
    'default-database' = '...',
    'username' = '...',
    'password' = '...',
    'base-url' = '...'
);

USE CATALOG my_catalog;

Java

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

String name            = "my_catalog";
String defaultDatabase = "mydb";
String username        = "...";
String password        = "...";
String baseUrl         = "..."

JdbcCatalog catalog = new JdbcCatalog(name, defaultDatabase, username, password, baseUrl);
tableEnv.registerCatalog("my_catalog", catalog);

// set the JdbcCatalog as the current catalog of the session
tableEnv.useCatalog("my_catalog");

Scala

val settings = EnvironmentSettings.inStreamingMode()
val tableEnv = TableEnvironment.create(settings)

val name            = "my_catalog"
val defaultDatabase = "mydb"
val username        = "..."
val password        = "..."
val baseUrl         = "..."

val catalog = new JdbcCatalog(name, defaultDatabase, username, password, baseUrl)
tableEnv.registerCatalog("my_catalog", catalog)

// set the JdbcCatalog as the current catalog of the session
tableEnv.useCatalog("my_catalog")

Python

from pyflink.table.catalog import JdbcCatalog

environment_settings = EnvironmentSettings.in_streaming_mode()
t_env = TableEnvironment.create(environment_settings)

name = "my_catalog"
default_database = "mydb"
username = "..."
password = "..."
base_url = "..."

catalog = JdbcCatalog(name, default_database, username, password, base_url)
t_env.register_catalog("my_catalog", catalog)

# set the JdbcCatalog as the current catalog of the session
t_env.use_catalog("my_catalog")

YAML

execution:
    ...
    current-catalog: my_catalog  # set the target JdbcCatalog as the current catalog of the session
    current-database: mydb

catalogs:
   - name: my_catalog
     type: jdbc
     default-database: mydb
     username: ...
     password: ...
     base-url: ...

PostgreSQL용 JDBC 카탈로그 (JDBC Catalog for PostgreSQL)

PostgreSQL 메타스페이스 매핑 (PostgreSQL Metaspace Mapping)

PostgreSQL은 데이터베이스 외에 schema라는 추가 네임스페이스를 가집니다. Postgres 인스턴스는 여러 데이터베이스를 가질 수 있고, 각 데이터베이스는 기본 이름 "public"을 가진 여러 스키마를 가질 수 있으며, 각 스키마는 여러 테이블을 가질 수 있습니다. Flink에서 Postgres 카탈로그로 등록된 테이블을 쿼리할 때 사용자는 schema_name.table_name 또는 그냥 table_name을 사용할 수 있습니다. schema_name은 선택 사항이며 기본값은 "public"입니다.

따라서 Flink Catalog와 Postgres 사이의 메타스페이스 매핑은 다음과 같습니다:

Flink Catalog 메타스페이스 구조 Postgres 메타스페이스 구조
카탈로그 이름(Flink에서만 정의) N/A
데이터베이스 이름 데이터베이스 이름
테이블 이름 [schema_name.]table_name

스키마가 지정되면 Flink에서 Postgres 테이블의 전체 경로는 "<catalog>.<db>.<schema.table>"여야 하며, <schema.table>은 이스케이프되어야 함에 유의하세요.

Postgres 테이블에 접근하는 몇 가지 예시:

-- scan table 'test_table' of 'public' schema (i.e. the default schema), the schema name can be omitted
SELECT * FROM mypg.mydb.test_table;
SELECT * FROM mydb.test_table;
SELECT * FROM test_table;

-- scan table 'test_table2' of 'custom_schema' schema,
-- the custom schema can not be omitted and must be escaped with table.
SELECT * FROM mypg.mydb.`custom_schema.test_table2`
SELECT * FROM mydb.`custom_schema.test_table2`;
SELECT * FROM `custom_schema.test_table2`;

MySQL용 JDBC 카탈로그 (JDBC Catalog for MySQL)

MySQL 메타스페이스 매핑 (MySQL Metaspace Mapping)

MySQL 인스턴스의 데이터베이스는 MySQL Catalog로 등록된 카탈로그 아래의 데이터베이스와 같은 매핑 수준에 있습니다. MySQL 인스턴스는 여러 데이터베이스를 가질 수 있고, 각 데이터베이스는 여러 테이블을 가질 수 있습니다. Flink에서 MySQL 카탈로그로 등록된 테이블을 쿼리할 때 사용자는 database.table_name 또는 그냥 table_name을 사용할 수 있습니다. 기본값은 MySQL Catalog 생성 시 지정된 기본 데이터베이스입니다.

따라서 Flink Catalog와 MySQL Catalog 사이의 메타스페이스 매핑은 다음과 같습니다:

Flink Catalog 메타스페이스 구조 MySQL 메타스페이스 구조
카탈로그 이름(Flink에서만 정의) N/A
데이터베이스 이름 데이터베이스 이름
테이블 이름 table_name

Flink에서 MySQL 테이블의 전체 경로는 "..

"이어야 합니다.

MySQL 테이블에 접근하는 몇 가지 예시:

-- scan table 'test_table', the default database is 'mydb'.
SELECT * FROM mysql_catalog.mydb.test_table;
SELECT * FROM mydb.test_table;
SELECT * FROM test_table;

-- scan table 'test_table' with the given database.
SELECT * FROM mysql_catalog.given_database.test_table2;
SELECT * FROM given_database.test_table2;

데이터 타입 매핑 (Data Type Mapping)

Flink는 MySQL, Oracle, PostgreSQL, Derby 같은 방언을 사용하는 여러 데이터베이스 연결을 지원합니다. Derby 방언은 보통 테스트 목적으로 사용됩니다. 관계형 데이터베이스 데이터 타입에서 Flink SQL 데이터 타입으로의 필드 데이터 타입 매핑은 다음 표에 나열되며, 이 매핑 테이블은 Flink에서 JDBC 테이블을 쉽게 정의하는 데 도움이 됩니다.

MySQL 타입 Oracle 타입 PostgreSQL 타입 SQL Server 타입 Flink SQL 타입
TINYINT TINYINT TINYINT TINYINT
SMALLINT, TINYINT UNSIGNED, MEDIUMINT SMALLINT INT2, SMALLSERIAL, SERIAL2 SMALLINT SMALLINT
INT, MEDIUMINT, SMALLINT UNSIGNED, INTEGER INTEGER SERIAL INT INT
BIGINT, INT UNSIGNED BIGINT BIGINT, BIGSERIAL BIGINT BIGINT
BIGINT UNSIGNED DECIMAL(20, 0)
FLOAT BINARY_FLOAT REAL, FLOAT4 REAL FLOAT
DOUBLE, DOUBLE PRECISION BINARY_DOUBLE FLOAT8, DOUBLE PRECISION FLOAT DOUBLE
NUMERIC(p, s), DECIMAL(p, s), SMALLINT, FLOAT(s), DOUBLE PRECISION, REAL NUMBER(p, s) NUMERIC(p, s), DECIMAL(p, s) DECIMAL(p, s) DECIMAL(p, s)
BOOLEAN, TINYINT(1) BOOLEAN BIT BOOLEAN
DATE DATE DATE DATE DATE
TIME [(p)] DATE TIME [(p)] [WITHOUT TIMEZONE] TIME(0) TIME [(p)] [WITHOUT TIMEZONE]
DATETIME [(p)] TIMESTAMP [(p)] [WITHOUT TIMEZONE] TIMESTAMP [(p)] [WITHOUT TIMEZONE] DATETIME, DATETIME2 TIMESTAMP [(p)] [WITHOUT TIMEZONE]
CHAR(n), VARCHAR(n), TEXT CHAR(n), VARCHAR(n), CLOB CHAR(n), CHARACTER(n), VARCHAR(n), CHARACTER VARYING(n), TEXT CHAR(n), NCHAR(n), VARCHAR(n), NVARCHAR(n), TEXT, NTEXT STRING
BINARY, VARBINARY, BLOB RAW(s), BLOB BYTEA BINARY(n), VARBINARY(n) BYTES
ARRAY ARRAY ARRAY

더 알아보기 (Learn more)