JDBC Sink 커넥터
JDBC Sink 커넥터
JDBC sink 커넥터는 Pulsar 토픽에서 메시지를 가져와 ClickHouse, MariaDB, PostgreSQL, SQLite에 저장하는 커넥터예요. JDBC를 통해 관계형 데이터베이스 및 그 호환 데이터베이스로 Pulsar 데이터를 써 넣을 때 사용해요.
참고: 모든 Pulsar 커넥터는 download page에서 내려받을 수 있어요.
출처: 문서
본문
JDBC sink 커넥터는 Pulsar 토픽에서 메시지를 가져와 ClickHouse, MariaDB, PostgreSQL, SQLite에 저장해요.
현재 INSERT, DELETE, UPDATE 연산이 지원돼요. SQLite, MariaDB, PostgreSQL은 UPSERT 연산과 멱등 쓰기(idempotent writes)도 지원해요.
구성 (Configuration)
모든 JDBC sink 커넥터의 구성에는 다음과 같은 프로퍼티가 있어요.
프로퍼티 (Property)
| 이름 | 타입 | 필수 | 기본값 | 설명 |
|---|---|---|---|---|
| userName | String | false | " " (빈 문자열) | jdbcUrl이 지정한 데이터베이스에 연결하는 데 사용하는 사용자 이름이에요. 참고: userName은 대소문자를 구분해요. |
| password | String | false | " " (빈 문자열) | jdbcUrl이 지정한 데이터베이스에 연결하는 데 사용하는 비밀번호예요. 참고: password는 대소문자를 구분해요. |
| jdbcUrl | String | true | " " (빈 문자열) | 커넥터가 연결하는 데이터베이스의 JDBC URL이에요. |
| tableName | String | true | " " (빈 문자열) | 커넥터가 쓰는 테이블 이름이에요. |
| nonKey | String | false | " " (빈 문자열) | 업데이트 이벤트에서 사용하는 필드를 담은 쉼표 구분 목록이에요. |
| key | String | false | " " (빈 문자열) | 업데이트 및 삭제 이벤트의 where 조건에서 사용하는 필드를 담은 쉼표 구분 목록이에요. |
| timeoutMs | int | false | 500 | JDBC 작업 타임아웃(밀리초)이에요. |
| batchSize | int | false | 200 | 데이터베이스에 적용되는 업데이트의 배치 크기예요. |
| insertMode | enum(INSERT, UPSERT, UPDATE) | false | INSERT | UPSERT로 설정하면 sink는 단순 INSERT/UPDATE 문 대신 upsert 의미를 사용해요. Upsert 의미는 기본 키 제약 위반이 있을 때 새 행을 원자적으로 추가하거나 기존 행을 업데이트하는 것으로, 멱등성을 제공해요. |
| nullValueAction | enum(FAIL, DELETE) | false | FAIL | NULL 값을 가진 레코드를 어떻게 처리할지예요. 가능한 옵션은 DELETE 또는 FAIL이에요. |
| useTransactions | boolean | false | true | 데이터베이스의 트랜잭션을 활성화할지 여부예요. |
| excludeNonDeclaredFields | boolean | false | false | 테이블의 모든 필드는 자동으로 발견돼요. excludeNonDeclaredFields는 nonKey와 key에 명시적으로 나열되지 않은 테이블 필드를 쿼리에 포함할지 여부를 나타내요. 기본적으로 모든 테이블 필드가 포함돼요. 삽입 중에 테이블 필드 기본값을 활용하려면 이 값을 false로 설정하는 것을 권장해요. |
| useJdbcBatch | boolean | false | false | JDBC 배치 API를 사용할지 여부예요. 쓰기 성능을 높이기 위해 이 옵션을 권장해요. |
ClickHouse 예시
- JSON
{
"configs": {
"userName": "clickhouse",
"password": "password",
"jdbcUrl": "jdbc:clickhouse://localhost:8123/pulsar_clickhouse_jdbc_sink",
"tableName": "pulsar_clickhouse_jdbc_sink"
"useTransactions": "false"
}
}
- YAML
tenant: "public"
namespace: "default"
name: "jdbc-clickhouse-sink"
inputs: [ "persistent://public/default/jdbc-clickhouse-topic" ]
sinkType: "jdbc-clickhouse"
configs:
userName: "clickhouse"
password: "password"
jdbcUrl: "jdbc:clickhouse://localhost:8123/pulsar_clickhouse_jdbc_sink"
tableName: "pulsar_clickhouse_jdbc_sink"
useTransactions: "false"
MariaDB 예시
- JSON
{
"configs": {
"userName": "mariadb",
"password": "password",
"jdbcUrl": "jdbc:mariadb://localhost:3306/pulsar_mariadb_jdbc_sink",
"tableName": "pulsar_mariadb_jdbc_sink"
}
}
- YAML
tenant: "public"
namespace: "default"
name: "jdbc-mariadb-sink"
inputs: [ "persistent://public/default/jdbc-mariadb-topic" ]
sinkType: "jdbc-mariadb"
configs:
userName: "mariadb"
password: "password"
jdbcUrl: "jdbc:mariadb://localhost:3306/pulsar_mariadb_jdbc_sink"
tableName: "pulsar_mariadb_jdbc_sink"
OpenMLDB 예시
OpenMLDB는 DELETE와 UPDATE 연산을 지원하지 않아요.
- JSON
{
"configs": {
"jdbcUrl": "jdbc:openmldb:///pulsar_openmldb_db?zk=localhost:6181&zkPath=/openmldb",
"tableName": "pulsar_openmldb_jdbc_sink"
}
}
- YAML
tenant: "public"
namespace: "default"
name: "jdbc-openmldb-sink"
inputs: [ "persistent://public/default/jdbc-openmldb-topic" ]
sinkType: "jdbc-openmldb"
configs:
jdbcUrl: "jdbc:openmldb:///pulsar_openmldb_db?zk=localhost:6181&zkPath=/openmldb"
tableName: "pulsar_openmldb_jdbc_sink"
PostgreSQL 예시
JDBC PostgreSQL sink 커넥터를 사용하기 전에 다음 방법 중 하나로 구성 파일을 만들어야 해요.
- JSON
{
"configs": {
"userName": "postgres",
"password": "password",
"jdbcUrl": "jdbc:postgresql://localhost:5432/pulsar_postgres_jdbc_sink",
"tableName": "pulsar_postgres_jdbc_sink"
}
}
- YAML
tenant: "public"
namespace: "default"
name: "jdbc-postgres-sink"
inputs: [ "persistent://public/default/jdbc-postgres-topic" ]
sinkType: "jdbc-postgres"
configs:
userName: "postgres"
password: "password"
jdbcUrl: "jdbc:postgresql://localhost:5432/pulsar_postgres_jdbc_sink"
tableName: "pulsar_postgres_jdbc_sink"
이 JDBC sink 커넥터를 사용하는 방법에 대한 더 자세한 내용은 connect Pulsar to PostgreSQL을 참고하세요.
SQLite 예시
- JSON
{
"configs": {
"jdbcUrl": "jdbc:sqlite:db.sqlite",
"tableName": "pulsar_sqlite_jdbc_sink"
}
}
- YAML
tenant: "public"
namespace: "default"
name: "jdbc-sqlite-sink"
inputs: [ "persistent://public/default/jdbc-sqlite-topic" ]
sinkType: "jdbc-sqlite"
configs:
jdbcUrl: "jdbc:sqlite:db.sqlite"
tableName: "pulsar_sqlite_jdbc_sink"