Flink JDBC 드라이버
Flink JDBC 드라이버 (Flink JDBC Driver)
Flink JDBC 드라이버는 클라이언트가 SQL Gateway를 통해 Flink 클러스터에 Flink SQL을 보낼 수 있게 해주는 Java 라이브러리예요. Beeline, SQLLine, Tableau 같은 다양한 JDBC 클라이언트에서 사용할 수 있어요.
출처: 문서
본문
Flink JDBC 드라이버는 클라이언트가 SQL Gateway를 통해 Flink 클러스터에 Flink SQL을 보낼 수 있게 해주는 Java 라이브러리예요.
Hive JDBC 드라이버를 Flink와 함께 사용할 수도 있어요. Hive dialect SQL을 실행하고 Hive Catalog를 활용하고 싶다면 유용해요. Flink와 함께 Hive JDBC를 사용하려면 HiveServer2 엔드포인트로 SQL Gateway를 실행해야 해요.
사용법 (Usage)
Flink JDBC 드라이버를 사용하기 전에 REST 엔드포인트로 SQL Gateway를 시작해야 해요. 이는 JDBC 서버 역할을 하고 Flink 클러스터에 바인딩해요. 아래 예시들은 게이트웨이가 시작되어 실행 중인 Flink 클러스터에 연결되어 있다고 가정해요.
의존성 (Dependency)
JDBC 드라이버의 모든 의존성은 flink-sql-jdbc-driver-bundle에 패키징되어 있어요. 다운로드해 프로젝트에 jar 파일을 추가할 수 있어요.
| Group Id | Artifact Id | JAR |
|---|---|---|
| org.apache.flink | flink-sql-jdbc-driver-bundle | Download |
maven 또는 gradle 프로젝트에 Flink JDBC 드라이버의 의존성을 추가할 수도 있어요.
Maven Dependency
org.apache.flink
flink-sql-jdbc-driver-bundle
{VERSION}
JDBC 클라이언트 (JDBC Clients)
Flink JDBC 드라이버는 Flink 배포에 포함되지 않아요. Maven에서 다운로드할 수 있어요. SLF4J(slf4j-api-{slf4j.version}.jar) jar도 필요할 수 있어요.
Beeline
Beeline은 Apache Hive에 접근하기 위한 명령줄 도구이지만, 일반 JDBC 드라이버도 지원해요. Hive와 beeline을 설치하려면 Hive 문서를 참고하세요.
- 다운로드 페이지에서
flink-jdbc-driver-bundle-{VERSION}.jar을 다운로드해$HIVE_HOME/lib에 추가해요. - beeline을 실행해 Flink SQL 게이트웨이에 연결해요. Flink SQL 게이트웨이는 현재 사용자 이름과 비밀번호를 무시하므로, 그냥 비워 두면 돼요.
- 원하는 문을 실행해요.
샘플 명령:
Beeline version 3.1.3 by Apache Hive
beeline> !connect jdbc:flink://localhost:8083
Connecting to jdbc:flink://localhost:8083
Enter username for jdbc:flink://localhost:8083:
Enter password for jdbc:flink://localhost:8083:
Connected to: Flink JDBC Driver (version 1.18-SNAPSHOT)
Driver: org.apache.flink.table.jdbc.FlinkDriver (version 1.18-SNAPSHOT)
0: jdbc:flink://localhost:8083> CREATE TABLE T(
. . . . . . . . . . . . . . . > a INT,
. . . . . . . . . . . . . . . > b VARCHAR(10)
. . . . . . . . . . . . . . . > ) WITH (
. . . . . . . . . . . . . . . > 'connector' = 'filesystem',
. . . . . . . . . . . . . . . > 'path' = 'file:///tmp/T.csv',
. . . . . . . . . . . . . . . > 'format' = 'csv'
. . . . . . . . . . . . . . . > );
No rows affected (0.108 seconds)
0: jdbc:flink://localhost:8083> INSERT INTO T VALUES (1, 'Hi'), (2, 'Hello');
+-----------------------------------+
| job id |
+-----------------------------------+
| da22010cf1c962b377493fc4fc509527 |
+-----------------------------------+
1 row selected (0.952 seconds)
0: jdbc:flink://localhost:8083> SELECT * FROM T;
+----+--------+
| a | b |
+----+--------+
| 1 | Hi |
| 2 | Hello |
+----+--------+
2 rows selected (1.142 seconds)
0: jdbc:flink://localhost:8083>
SQLLine
SQLLine은 일반 JDBC 드라이버를 지원하는 가벼운 JDBC 명령줄 도구예요. SQLLine을 사용하려면 GitHub 저장소를 클론하고 먼저 프로젝트를 컴파일해야 해요 (./mvnw package -DskipTests).
- Flink JDBC 드라이버(
flink-jdbc-driver-bundle-{VERSION}.jar) - SLF4J(
slf4j-api-{slf4j.version}.jar) - 명령
./bin/sqlline으로 SQLLine을 실행해요.
sqlline version 1.13.0-SNAPSHOT
sqlline> !connect jdbc:flink://localhost:8083
Enter username for jdbc:flink://localhost:8083:
Enter password for jdbc:flink://localhost:8083:
0: jdbc:flink://localhost:8083>
- SQLLine에서
!connect명령으로 Flink SQL 게이트웨이에 연결해요. Flink SQL 게이트웨이는 현재 사용자 이름과 비밀번호를 무시하므로 비워 두면 돼요. - 이제 원하는 Flink SQL 문을 실행할 수 있어요.
샘플 명령:
0: jdbc:flink://localhost:8083> CREATE TABLE T(
. . . . . . . . . . . . . . .)> a INT,
. . . . . . . . . . . . . . .)> b VARCHAR(10)
. . . . . . . . . . . . . . .)> ) WITH (
. . . . . . . . . . . . . . .)> 'connector' = 'filesystem',
. . . . . . . . . . . . . . .)> 'path' = 'file:///tmp/T.csv',
. . . . . . . . . . . . . . .)> 'format' = 'csv'
. . . . . . . . . . . . . . .)> );
No rows affected (0.122 seconds)
0: jdbc:flink://localhost:8083> INSERT INTO T VALUES (1, 'Hi'), (2, 'Hello');
+----------------------------------+
| job id |
+----------------------------------+
| fbade1ab4450fc57ebd5269fdf60dcfd |
+----------------------------------+
1 row selected (1.282 seconds)
0: jdbc:flink://localhost:8083> SELECT * FROM T;
+---+-------+
| a | b |
+---+-------+
| 1 | Hi |
| 2 | Hello |
+---+-------+
2 rows selected (1.955 seconds)
0: jdbc:flink://localhost:8083>
Tableau
Tableau는 대화형 데이터 시각화 소프트웨어예요. 2018.3 버전부터 Other Database (JDBC) 연결을 지원해요. Flink JDBC 드라이버를 사용하려면 Tableau 버전 >= 2018.3이 필요해요. Tableau에서 Other Database (JDBC)의 일반적인 사용법은 Tableau 문서를 참고하세요.
- 드라이버를 설치해요. (Windows:
C:\Program Files\Tableau\Drivers, Mac:~/Library/Tableau/Drivers, Linux:/opt/tableau/tableau_driver/jdbc) - Connect에서 Other Database (JDBC)를 선택하고 Flink SQL 게이트웨이의 url을 입력해요. SQL92 dialect를 선택하고 사용자 이름과 비밀번호는 비워 두어요.
- Login 버튼을 누르고 평소처럼 Tableau를 사용해요.
다른 JDBC 도구와 함께 사용 (Use with other JDBC Tools)
JDBC API를 지원하는 어떤 도구든 Flink JDBC 드라이버와 Flink SQL 게이트웨이와 함께 사용할 수 있어요. 사용자 정의 JDBC 드라이버를 사용하는 방법은 원하는 도구의 문서를 참고하세요.
애플리케이션과 함께 사용 (Use with Application)
Java
Flink JDBC 드라이버는 JDBC API를 통해 Flink 클러스터에 접근하기 위한 라이브러리예요. Java에서 JDBC의 일반적인 사용법은 JDBC 튜토리얼을 참고하세요.
- 프로젝트의 pom.xml에 다음 의존성을 추가하거나
flink-jdbc-driver-bundle-{VERSION}.jar을 다운로드해 클래스패스에 추가해요. - Java 코드에서 특정 url로 Flink SQL 게이트웨이에 연결해요.
- 원하는 문을 실행해요.
Sample.java
public class Sample {
public static void main(String[] args) throws Exception {
try (Connection connection = DriverManager.getConnection("jdbc:flink://localhost:8083")) {
try (Statement statement = connection.createStatement()) {
statement.execute("CREATE TABLE T(\n" +
" a INT,\n" +
" b VARCHAR(10)\n" +
") WITH (\n" +
" 'connector' = 'filesystem',\n" +
" 'path' = 'file:///tmp/T.csv',\n" +
" 'format' = 'csv'\n" +
")");
statement.execute("INSERT INTO T VALUES (1, 'Hi'), (2, 'Hello')");
try (ResultSet rs = statement.executeQuery("SELECT * FROM T")) {
while (rs.next()) {
System.out.println(rs.getInt(1) + ", " + rs.getString(2));
}
}
}
}
}
}
출력:
1, Hi
2, Hello
DriverManager 외에도 Flink JDBC 드라이버는 DataSource를 지원하며, 그로부터 연결을 만들 수도 있어요.
DataSource.java
public class Sample {
public static void main(String[] args) throws Exception {
DataSource dataSource = new FlinkDataSource("jdbc:flink://localhost:8083", new Properties());
try (Connection connection = dataSource.getConnection()) {
try (Statement statement = connection.createStatement()) {
statement.execute("CREATE TABLE T(\n" +
" a INT,\n" +
" b VARCHAR(10)\n" +
") WITH (\n" +
" 'connector' = 'filesystem',\n" +
" 'path' = 'file:///tmp/T.csv',\n" +
" 'format' = 'csv'\n" +
")");
statement.execute("INSERT INTO T VALUES (1, 'Hi'), (2, 'Hello')");
try (ResultSet rs = statement.executeQuery("SELECT * FROM T")) {
while (rs.next()) {
System.out.println(rs.getInt(1) + ", " + rs.getString(2));
}
}
}
}
}
}
다른 언어 (Other languages)
Java 외에도 Flink JDBC 드라이버는 Scala, Kotlin 등 어떤 JVM 언어로든 사용할 수 있어요. 프로젝트에 Flink JDBC 드라이버의 의존성을 추가하고 직접 사용해요.
많은 애플리케이션이 SQL 데이터베이스의 데이터에 직접 또는 JOOQ, MyBatis, Spring Data 같은 프레임워크를 통해 접근해요. 이러한 애플리케이션과 프레임워크를 Flink JDBC 드라이버를 사용하도록 구성하면, 일반 데이터베이스 대신 Flink 클러스터에서 SQL 쿼리를 수행해요.