Pulsar를 데이터베이스에 연결하는 방법
Pulsar를 데이터베이스에 연결하는 방법 (How to connect Pulsar to database)
이 튜토리얼은 코드를 한 줄도 작성하지 않고 Pulsar에서 데이터를 외부로 내보내는 방법을 직접 실습해 보는 내용이에요. Pulsar I/O의 개념을 더 깊게 이해하고 싶다면, 이 가이드의 단계를 직접 실행해 보면서 익히는 게 큰 도움이 돼요. 이 튜토리얼을 끝내면 다음 작업을 할 수 있게 돼요.
- Pulsar를 Cassandra와 ScyllaDB에 연결하기
- Pulsar를 PostgreSQL에 연결하기
팁
- 이 지침은 Pulsar를 standalone 모드로 실행 중이라고 가정해요. 하지만 튜토리얼에서 쓰는 모든 명령은 수정 없이 다중 노드 Pulsar 클러스터에서도 사용할 수 있어요.
- 모든 지침은 Pulsar 바이너리 배포판의 루트 디렉터리에서 실행한다고 가정해요.
출처: 문서
본문
Pulsar와 내장 커넥터 설치하기 (Install Pulsar and built-in connector)
Pulsar를 데이터베이스에 연결하기 전에, Pulsar와 원하는 내장 커넥터를 먼저 설치해야 해요.
Pulsar 배포판을 다운로드하는 방법은 Run a standalone Pulsar cluster locally 문서를 참고하세요.
Pulsar 커넥터를 사용하려면 download page에서 커넥터 타르볼(tarball) 릴리스를 다운로드해야 해요.
NAR 파일을 다운로드한 뒤에는 그 파일을 Pulsar 디렉터리의 connectors 디렉터리로 복사해요. 예를 들어 pulsar-io-aerospike-5.0.0-M2.nar 커넥터 파일을 다운로드했다면 다음 명령을 입력하세요.
mkdir connectors
mv pulsar-io-aerospike-5.0.0-M2.nar connectors
ls connectors
# pulsar-io-aerospike-5.0.0-M2.nar
# ...
참고
- 베어메탈(bare metal) 클러스터에서 Pulsar를 실행 중이라면, 모든 브로커의 Pulsar 디렉터리마다
connectors타르볼의 압축이 풀려 있어야 해요. (Functions용으로 별도의 워커 클러스터를 운영 중이라면 모든 function-worker의 Pulsar 디렉터리에도 마찬가지로요.)- Docker에서 Pulsar를 실행하거나 Docker 이미지로(예: K8S) 배포한다면,
apachepulsar/pulsar이미지 대신apachepulsar/pulsar-all이미지를 사용할 수 있어요.apachepulsar/pulsar-all이미지에는 모든 내장 커넥터가 이미 번들로 포함되어 있어요.
Pulsar standalone 시작하기 (Start Pulsar standalone)
- Pulsar를 로컬에서 시작해요.
bin/pulsar standalone
Pulsar 서비스의 모든 구성 요소가 순서대로 시작돼요.
Pulsar 서비스가 정상적으로 기동됐는지 확인하려면 아래 엔드포인트를 curl로 확인하면 돼요.
- Pulsar 바이너리 프로토콜 포트를 확인해요.
telnet localhost 6650
- Pulsar Function 클러스터를 확인해요.
curl -s http://localhost:8080/admin/v2/worker/cluster
예시 출력:
[{"workerId":"c-standalone-fw-localhost-6750","workerHostname":"localhost","port":6750}]
- public 테넌트와 default 네임스페이스가 존재하는지 확인해요.
curl -s http://localhost:8080/admin/v2/namespaces/public
예시 출력:
["public/default","public/functions"]
- 모든 내장 커넥터가 사용 가능한 것으로 표시되는지 확인해요.
curl -s http://localhost:8080/admin/v2/functions/connectors
예시 출력:
[{"name":"aerospike","description":"Aerospike database sink","sinkClass":"org.apache.pulsar.io.aerospike.AerospikeStringSink"},{"name":"cassandra","description":"Writes data into Cassandra","sinkClass":"org.apache.pulsar.io.cassandra.CassandraStringSink"},{"name":"kafka","description":"Kafka source and sink connector","sourceClass":"org.apache.pulsar.io.kafka.KafkaStringSource","sinkClass":"org.apache.pulsar.io.kafka.KafkaBytesSink"},{"name":"kinesis","description":"Kinesis sink connector","sinkClass":"org.apache.pulsar.io.kinesis.KinesisSink"},{"name":"rabbitmq","description":"RabbitMQ source connector","sourceClass":"org.apache.pulsar.io.rabbitmq.RabbitMQSource"}]
Pulsar 서비스를 시작할 때 오류가 발생했다면, pulsar/standalone을 실행 중인 터미널에서 예외를 확인할 수 있어요. 아니면 Pulsar 디렉터리 아래 logs 디렉터리로 이동해서 로그를 확인해도 돼요.
Pulsar를 Cassandra와 ScyllaDB에 연결하기 (Connect Pulsar to Cassandra and ScyllaDB)
이 섹션은 Pulsar를 Cassandra와 ScyllaDB에 연결하는 방법을 보여줘요.
팁
- Docker가 설치되어 있는지 확인하세요. 설치되어 있지 않다면 install Docker를 참고하세요. Docker 명령어에 대한 더 자세한 내용은 Docker CLI 문서를 참고하세요.
- Cassandra sink 커넥터는 Pulsar 토픽에서 메시지를 읽어 Cassandra나 ScyllaDB 테이블에 써 넣어요. 더 자세한 내용은 Cassandra sink connector 문서를 참고하세요.
- ScyllaDB 호환성: ScyllaDB는 CQL 프로토콜 전체 호환을 유지하면서 Cassandra를 그대로 대체할 수 있는(drop-in replacement) 제품이에요. 같은
pulsar-io-cassandra커넥터가 수정 없이 두 데이터베이스 모두에서 동작해요. 더 자세한 내용은 Streaming Real-Time Chat Messages into ScyllaDB with Apache Pulsar를 참고하세요.
Pulsar를 Cassandra나 ScyllaDB에 연결하려면 아래 단계를 따라 하면 돼요.
1단계: Cassandra 또는 ScyllaDB 클러스터 설정하기 (Set up a Cassandra or ScyllaDB cluster)
이 예시는 cassandra Docker 이미지로 Docker 안에 단일 노드 Cassandra 클러스터를 시작해요. 또는 아래처럼 ScyllaDB를 사용할 수도 있어요.
옵션 A: Cassandra 사용하기
- Cassandra 클러스터를 시작해요.
docker run -d --rm --name=cassandra -p 9042:9042 cassandra:3.11
참고: 다음 단계로 넘어가기 전에 Cassandra 클러스터가 실행 중인지 확인하세요.
- Docker 프로세스가 실행 중인지 확인해요.
docker ps
- Cassandra 로그를 확인해서 Cassandra 프로세스가 정상적으로 동작 중인지 확인해요.
docker logs cassandra
- Cassandra 클러스터의 상태를 확인해요.
docker exec cassandra nodetool status
예시 출력:
Datacenter: datacenter1
=======================
Status=Up/Down
|/ State=Normal/Leaving/Joining/Moving
-- Address Load Tokens Owns (effective) Host ID Rack
UN 172.17.0.2 103.67 KiB 256 100.0% af0e4b2f-84e0-4f0b-bb14-bd5f9070ff26 rack1
cqlsh를 사용해 Cassandra 클러스터에 연결해요.
docker exec -ti cassandra cqlsh localhost
출력:
Connected to Test Cluster at localhost:9042.
[cqlsh 5.0.1 | Cassandra 3.11.2 | CQL spec 3.4.4 | Native protocol v4]
Use HELP for help.
cqlsh>
- keyspace
pulsar_test_keyspace를 만들어요.
cqlsh> CREATE KEYSPACE pulsar_test_keyspace WITH replication = {'class':'SimpleStrategy', 'replication_factor':1};
- 테이블
pulsar_test_table을 만들어요.
cqlsh> USE pulsar_test_keyspace;
cqlsh:pulsar_test_keyspace> CREATE TABLE pulsar_test_table (key text PRIMARY KEY, col text);
옵션 B: ScyllaDB 사용하기
- ScyllaDB 클러스터를 시작해요.
docker run -d --rm --name=scylladb -p 9042:9042 \
scylladb/scylla:latest \
--smp 1 --memory 750M --overprovisioned 1
참고: ScyllaDB는 컨테이너에서 실행할 때 특정 플래그가 필요해요. 위 플래그는 제한된 리소스로 단일 코어 연산용으로 구성하는 값이에요.
- Docker 프로세스가 실행 중인지 확인해요.
docker ps
- ScyllaDB 로그를 확인해서 ScyllaDB 프로세스가 정상적으로 동작 중인지 확인해요.
docker logs scylladb
ScyllaDB가 준비될 때까지 기다려요(대개 30~60초 걸려요). 데이터베이스가 CQL 연결을 받을 준비가 됐다는 메시지를 확인하세요.
- ScyllaDB 클러스터의 상태를 확인해요.
docker exec scylladb nodetool status
cqlsh를 사용해 ScyllaDB 클러스터에 연결해요.
docker exec -ti scylladb cqlsh
- keyspace
pulsar_test_keyspace를 만들어요.
cqlsh> CREATE KEYSPACE pulsar_test_keyspace WITH replication = {'class':'SimpleStrategy', 'replication_factor':1};
- 테이블
pulsar_test_table을 만들어요.
cqlsh> USE pulsar_test_keyspace;
cqlsh:pulsar_test_keyspace> CREATE TABLE pulsar_test_table (key text PRIMARY KEY, col text);
2단계: Cassandra sink 구성하기 (Configure a Cassandra sink)
이제 로컬에 Cassandra 또는 ScyllaDB 클러스터가 실행되고 있어요.
이 섹션에서는 Cassandra sink 커넥터를 구성할 거예요. 같은 커넥터가 Cassandra와 ScyllaDB 양쪽 모두에서 동작해요.
Cassandra sink 커넥터를 실행하려면 Pulsar 커넥터 런타임이 알아야 하는 정보가 담긴 구성 파일을 준비해야 해요. 예를 들어, Pulsar 커넥터가 Cassandra 클러스터를 어떻게 찾을지, Pulsar 메시지를 쓸 keyspace와 테이블이 무엇인지 같은 정보요.
구성 파일은 다음 방법 중 하나로 만들 수 있어요.
- JSON
{
"roots": "localhost:9042",
"keyspace": "pulsar_test_keyspace",
"columnFamily": "pulsar_test_table",
"keyname": "key",
"columnName": "col"
}
- YAML
configs:
roots: "localhost:9042"
keyspace: "pulsar_test_keyspace"
columnFamily: "pulsar_test_table"
keyname: "key"
columnName: "col"
참고: ScyllaDB의 경우 구성은 동일해요. 다른 컨테이너 이름이나 호스트명을 썼다면
roots파라미터를 그에 맞게 수정하세요(예: Docker 네트워크 안에서 연결한다면"scylladb:9042").
더 자세한 내용은 Cassandra sink connector 문서를 참고하세요.
3단계: Cassandra sink 만들기 (Create a Cassandra sink)
Connector Admin CLI를 사용해서 sink 커넥터를 만들고 그에 대한 다양한 작업을 수행할 수 있어요.
sink 타입 cassandra와, 앞서 만든 구성 파일 /examples/cassandra-sink.yml을 사용해 Cassandra sink 커넥터를 만들려면 다음 명령을 실행하세요.
참고: 현재 내장 커넥터의
sink-type파라미터는pulsar-io.yaml파일에 지정된name파라미터 설정에 따라 결정돼요.
bin/pulsar-admin sinks create \
--tenant public \
--namespace default \
--name cassandra-test-sink \
--sink-type cassandra \
--sink-config-file $PWD/examples/cassandra-sink.yml \
--inputs test_cassandra
명령이 실행되면 Pulsar가 cassandra-test-sink라는 sink 커넥터를 만들어요.
이 sink 커넥터는 Pulsar Function으로 실행되며, test_cassandra 토픽에서 생성된 메시지를 Cassandra 테이블 pulsar_test_table에 써 넣어요.
4단계: Cassandra sink 살펴보기 (Inspect a Cassandra sink)
Connector Admin CLI를 사용해서 커넥터를 모니터링하고 그에 대한 다양한 작업을 수행할 수 있어요.
- Cassandra sink의 정보를 가져와요.
bin/pulsar-admin sinks get \
--tenant public \
--namespace default \
--name cassandra-test-sink
예시 출력:
{"tenant": "public","namespace": "default","name": "cassandra-test-sink","className": "org.apache.pulsar.io.cassandra.CassandraStringSink","inputSpecs": { "test_cassandra": { "isRegexPattern": false }},"configs": { "roots": "localhost:9042", "keyspace": "pulsar_test_keyspace", "columnFamily": "pulsar_test_table", "keyname": "key", "columnName": "col"},"parallelism": 1,"processingGuarantees": "ATLEAST_ONCE","retainOrdering": false,"autoAck": true,"archive": "builtin://cassandra"}
- Cassandra sink의 상태를 확인해요.
bin/pulsar-admin sinks status \
--tenant public \
--namespace default \
--name cassandra-test-sink
예시 출력:
{"numInstances" : 1,"numRunning" : 1,"instances" : [ { "instanceId" : 0, "status" : { "running" : true, "error" : "", "numRestarts" : 0, "numReadFromPulsar" : 0, "numSystemExceptions" : 0, "latestSystemExceptions" : [ ], "numSinkExceptions" : 0, "latestSinkExceptions" : [ ], "numWrittenToSink" : 0, "lastReceivedTime" : 0, "workerId" : "c-standalone-fw-localhost-8080" }} ]}
5단계: Cassandra sink 검증하기 (Verify a Cassandra sink)
- Cassandra sink인 test_cassandra의 입력 토픽에 메시지를 몇 개 생성해요.
for i in {0..9}; do bin/pulsar-client produce -m "key-$i" -n 1 test_cassandra; done
- Cassandra sink인 test_cassandra의 상태를 확인해요.
bin/pulsar-admin sinks status \
--tenant public \
--namespace default \
--name cassandra-test-sink
Cassandra sink인 test_cassandra가 10개의 메시지를 처리한 것을 볼 수 있어요.
예시 출력:
{
"numInstances" : 1,
"numRunning" : 1,
"instances" : [ {
"instanceId" : 0,
"status" : {
"running" : true,
"error" : "",
"numRestarts" : 0,
"numReadFromPulsar" : 10,
"numSystemExceptions" : 0,
"latestSystemExceptions" : [ ],
"numSinkExceptions" : 0,
"latestSinkExceptions" : [ ],
"numWrittenToSink" : 10,
"lastReceivedTime" : 1551685489136,
"workerId" : "c-standalone-fw-localhost-8080"
}
} ]
}
cqlsh를 사용해 Cassandra 클러스터에 연결해요.
docker exec -ti cassandra cqlsh localhost
- Cassandra 테이블 pulsar_test_table의 데이터를 확인해요.
cqlsh> use pulsar_test_keyspace;
cqlsh:pulsar_test_keyspace> select * from pulsar_test_table;
key | col
--------+--------
key-5 | key-5
key-0 | key-0
key-9 | key-9
key-2 | key-2
key-1 | key-1
key-3 | key-3
key-6 | key-6
key-7 | key-7
key-4 | key-4
key-8 | key-8
6단계: Cassandra sink 삭제하기 (Delete a Cassandra Sink)
Connector Admin CLI를 사용해서 커넥터를 삭제하고 그에 대한 다양한 작업을 수행할 수 있어요.
bin/pulsar-admin sinks delete \
--tenant public \
--namespace default \
--name cassandra-test-sink
Pulsar를 PostgreSQL에 연결하기 (Connect Pulsar to PostgreSQL)
이 섹션은 Pulsar를 PostgreSQL에 연결하는 방법을 보여줘요.
팁
- Docker가 설치되어 있는지 확인하세요. 설치되어 있지 않다면 install Docker를 참고하세요. Docker 명령어에 대한 더 자세한 내용은 Docker CLI 문서를 참고하세요.
- JDBC sink 커넥터는 Pulsar 토픽에서 메시지를 가져와 ClickHouse, MariaDB, PostgreSQL, SQLite에 저장해요. 더 자세한 내용은 JDBC sink connector 문서를 참고하세요.
Pulsar를 PostgreSQL에 연결하려면 아래 단계를 따라 하면 돼요.
1단계: PostgreSQL 클러스터 설정하기 (Set up a PostgreSQL cluster)
이 예시는 PostgreSQL 12 Docker 이미지로 Docker 안에 단일 노드 PostgreSQL 클러스터를 시작해요.
- Docker에서 PostgreSQL 12 이미지를 내려받아요.
docker pull postgres:12
- PostgreSQL을 시작해요.
docker run -d -it --rm \
--name pulsar-postgres \
-p 5432:5432 \
-e POSTGRES_PASSWORD=password \
-e POSTGRES_USER=postgres \
postgres:12
- PostgreSQL이 정상적으로 시작됐는지 확인해요.
docker logs -f pulsar-postgres
다음 메시지가 나타나면 PostgreSQL이 정상적으로 시작된 거예요.
2020-05-11 20:09:24.492 UTC [1] LOG: starting PostgreSQL 12.2 (Debian 12.2-2.pgdg100+1) on x86_64-pc-linux-gnu, compiled by gcc (Debian 8.3.0-6) 8.3.0, 64-bit
2020-05-11 20:09:24.492 UTC [1] LOG: listening on IPv4 address "0.0.0.0", port 5432
2020-05-11 20:09:24.492 UTC [1] LOG: listening on IPv6 address "::", port 5432
2020-05-11 20:09:24.499 UTC [1] LOG: listening on Unix socket "/var/run/postgresql/.s.PGSQL.5432"
2020-05-11 20:09:24.523 UTC [55] LOG: database system was shut down at 2020-05-11 20:09:24 UTC
2020-05-11 20:09:24.533 UTC [1] LOG: database system is ready to accept connections
- PostgreSQL 컨테이너에 접근해요.
docker exec -it pulsar-postgres /bin/bash
- 기본 사용자 이름과 비밀번호로 PostgreSQL에 로그인해요.
psql -U postgres postgres
- 다음 명령으로
pulsar_postgres_jdbc_sink테이블을 만들어요.
create table if not exists pulsar_postgres_jdbc_sink(id serial PRIMARY KEY,name VARCHAR(255) NOT NULL);
2단계: JDBC sink 구성하기 (Configure a JDBC sink)
이제 로컬에 PostgreSQL이 실행되고 있어요.
이 섹션에서는 JDBC sink 커넥터를 구성할 거예요.
- 구성 파일을 추가해요.
JDBC sink 커넥터를 실행하려면 Pulsar 커넥터 런타임이 알아야 하는 정보가 담긴 YAML 구성 파일을 준비해야 해요. 예를 들어, Pulsar 커넥터가 PostgreSQL 클러스터를 어떻게 찾을지, JDBC URL과 메시지를 쓸 테이블이 무엇인지 같은 정보요.
pulsar-postgres-jdbc-sink.yaml 파일을 만들고 다음 내용을 복사한 뒤 pulsar/connectors 폴더에 넣으세요.
configs:
userName: "postgres"
password: "password"
jdbcUrl: "jdbc:postgresql://localhost:5432/postgres"
tableName: "pulsar_postgres_jdbc_sink"
- 스키마를 만들어요.
avro-schema 파일을 만들고 다음 내용을 복사한 뒤 pulsar/connectors 폴더에 넣으세요.
{
"type": "AVRO",
"schema": "{\"type\":\"record\",\"name\":\"Test\",\"fields\":[{\"name\":\"id\",\"type\":[\"null\",\"int\"]},{\"name\":\"name\",\"type\":[\"null\",\"string\"]}]}",
"properties": {}
}
팁: AVRO에 대한 더 자세한 내용은 Apache Avro를 참고하세요.
- 토픽에 스키마를 업로드해요.
이 예시는 avro-schema 스키마를 pulsar-postgres-jdbc-sink-topic 토픽에 업로드해요.
bin/pulsar-admin schemas upload pulsar-postgres-jdbc-sink-topic -f ./connectors/avro-schema
- 스키마가 정상적으로 업로드됐는지 확인해요.
bin/pulsar-admin schemas get pulsar-postgres-jdbc-sink-topic
다음 메시지가 나타나면 스키마가 정상적으로 업로드된 거예요.
{"name":"pulsar-postgres-jdbc-sink-topic","schema":"{\"type\":\"record\",\"name\":\"Test\",\"fields\":[{\"name\":\"id\",\"type\":[\"null\",\"int\"]},{\"name\":\"name\",\"type\":[\"null\",\"string\"]}]}","type":"AVRO","properties":{}}
3단계: JDBC sink 만들기 (Create a JDBC sink)
Connector Admin CLI를 사용해서 sink 커넥터를 만들고 그에 대한 다양한 작업을 수행할 수 있어요.
이 예시는 sink 커넥터를 만들면서 원하는 정보를 지정해요.
bin/pulsar-admin sinks create \
--archive $PWD/connectors/pulsar-io-jdbc-postgres-5.0.0-M2.nar \
--inputs pulsar-postgres-jdbc-sink-topic \
--name pulsar-postgres-jdbc-sink \
--sink-config-file $PWD/connectors/pulsar-postgres-jdbc-sink.yaml \
--parallelism 1
명령이 실행되면 Pulsar가 pulsar-postgres-jdbc-sink라는 sink 커넥터를 만들어요.
이 sink 커넥터는 Pulsar Function으로 실행되며, pulsar-postgres-jdbc-sink-topic 토픽에서 생성된 메시지를 PostgreSQL 테이블 pulsar_postgres_jdbc_sink에 써 넣어요.
팁 (Tip)
| Flag | 설명 | 예시 |
|---|---|---|
| --archive | sink의 아카이브 파일에 대한 절대 경로예요. | $PWD/pulsar-io-jdbc-postgres-5.0.0-M2.nar |
| --inputs | sink의 입력 토픽(들)이에요. 여러 토픽은 쉼표로 구분해 지정할 수 있어요. | |
| --name | sink의 이름이에요. | pulsar-postgres-jdbc-sink |
| --sink-config-file | sink의 구성을 지정하는 YAML 구성 파일의 절대 경로예요. | $PWD/pulsar-postgres-jdbc-sink.yaml |
| --parallelism | sink의 병렬 처리 계수예요. 예를 들어 실행할 sink 인스턴스 수를 말해요. | 1 |
팁:
pulsar-admin sinks create options에 대한 더 자세한 내용은 Pulsar admin 문서를 참고하세요.
다음 메시지가 나타나면 sink가 정상적으로 생성된 거예요.
Created successfully
4단계: JDBC sink 살펴보기 (Inspect a JDBC sink)
Connector Admin CLI를 사용해서 커넥터를 모니터링하고 그에 대한 다양한 작업을 수행할 수 있어요.
- 실행 중인 모든 JDBC sink를 나열해요.
bin/pulsar-admin sinks list \
--tenant public \
--namespace default
팁:
pulsar-admin sinks list options에 대한 더 자세한 내용은 Pulsar admin 문서를 참고하세요.
결과를 보면 postgres-jdbc-sink sink만 실행 중인 것을 알 수 있어요.
["pulsar-postgres-jdbc-sink"]
- JDBC sink의 정보를 가져와요.
bin/pulsar-admin sinks get \
--tenant public \
--namespace default \
--name pulsar-postgres-jdbc-sink
팁:
pulsar-admin sinks get options에 대한 더 자세한 내용은 Pulsar admin 문서를 참고하세요.
결과는 sink 커넥터의 정보(tenant, namespace, topic 등)를 보여줘요.
{"tenant": "public","namespace": "default","name": "pulsar-postgres-jdbc-sink","className": "org.apache.pulsar.io.jdbc.PostgresJdbcAutoSchemaSink","inputSpecs": { "pulsar-postgres-jdbc-sink-topic": { "isRegexPattern": false }},"configs": { "password": "password", "jdbcUrl": "jdbc:postgresql://localhost:5432/pulsar_postgres_jdbc_sink", "userName": "postgres", "tableName": "pulsar_postgres_jdbc_sink"},"parallelism": 1,"processingGuarantees": "ATLEAST_ONCE","retainOrdering": false,"autoAck": true}
- JDBC sink의 상태를 가져와요.
bin/pulsar-admin sinks status \
--tenant public \
--namespace default \
--name pulsar-postgres-jdbc-sink
팁:
pulsar-admin sinks status options에 대한 더 자세한 내용은 Pulsar admin 문서를 참고하세요.
결과는 sink 커넥터의 현재 상태(인스턴스 수, 실행 상태, 워커 ID 등)를 보여줘요.
{"numInstances" : 1,"numRunning" : 1,"instances" : [ { "instanceId" : 0, "status" : { "running" : true, "error" : "", "numRestarts" : 0, "numReadFromPulsar" : 0, "numSystemExceptions" : 0, "latestSystemExceptions" : [ ], "numSinkExceptions" : 0, "latestSinkExceptions" : [ ], "numWrittenToSink" : 0, "lastReceivedTime" : 0, "workerId" : "c-standalone-fw-192.168.2.52-8080" }} ]}
5단계: JDBC sink 중지하기 (Stop a JDBC sink)
Connector Admin CLI를 사용해서 커넥터를 중지하고 그에 대한 다양한 작업을 수행할 수 있어요.
bin/pulsar-admin sinks stop \
--tenant public \
--namespace default \
--name pulsar-postgres-jdbc-sink
팁:
pulsar-admin sinks stop options에 대한 더 자세한 내용은 Pulsar admin 문서를 참고하세요.
다음 메시지가 나타나면 sink 인스턴스가 정상적으로 중지된 거예요.
Stopped successfully
6단계: JDBC sink 재시작하기 (Restart a JDBC sink)
Connector Admin CLI를 사용해서 커넥터를 재시작하고 그에 대한 다양한 작업을 수행할 수 있어요.
bin/pulsar-admin sinks restart \
--tenant public \
--namespace default \
--name pulsar-postgres-jdbc-sink
팁:
pulsar-admin sinks restart options에 대한 더 자세한 내용은 Pulsar admin 문서를 참고하세요.
다음 메시지가 나타나면 sink 인스턴스가 정상적으로 시작된 거예요.
Started successfully
팁
- 원한다면
pulsar-admin sinks localrun options으로 sink 커넥터를 standalone으로 실행할 수도 있어요.pulsar-admin sinks localrun options은 sink 커넥터를 로컬에서 실행하는 반면,pulsar-admin sinks start options은 sink 커넥터를 클러스터에서 시작한다는 점을 기억하세요.pulsar-admin sinks localrun options에 대한 더 자세한 내용은 Pulsar admin 문서를 참고하세요.
7단계: JDBC sink 업데이트하기 (Update a JDBC sink)
Connector Admin CLI를 사용해서 커넥터를 업데이트하고 그에 대한 다양한 작업을 수행할 수 있어요.
이 예시는 pulsar-postgres-jdbc-sink sink 커넥터의 병렬 처리 계수를 2로 업데이트해요.
bin/pulsar-admin sinks update \
--name pulsar-postgres-jdbc-sink \
--parallelism 2
팁:
pulsar-admin sinks update options에 대한 더 자세한 내용은 Pulsar admin 문서를 참고하세요.
다음 메시지가 나타나면 sink 커넥터가 정상적으로 업데이트된 거예요.
Updated successfully
이 예시는 정보를 다시 한번 확인해 봐요.
bin/pulsar-admin sinks get \
--tenant public \
--namespace default \
--name pulsar-postgres-jdbc-sink
결과를 보면 병렬 처리 계수가 2인 것을 확인할 수 있어요.
{
"tenant": "public",
"namespace": "default",
"name": "pulsar-postgres-jdbc-sink",
"className": "org.apache.pulsar.io.jdbc.PostgresJdbcAutoSchemaSink",
"inputSpecs": {
"pulsar-postgres-jdbc-sink-topic": {
"isRegexPattern": false
}
},
"configs": {
"password": "password",
"jdbcUrl": "jdbc:postgresql://localhost:5432/pulsar_postgres_jdbc_sink",
"userName": "postgres",
"tableName": "pulsar_postgres_jdbc_sink"
},
"parallelism": 2,
"processingGuarantees": "ATLEAST_ONCE",
"retainOrdering": false,
"autoAck": true
}
8단계: JDBC sink 삭제하기 (Delete a JDBC sink)
Connector Admin CLI를 사용해서 커넥터를 삭제하고 그에 대한 다양한 작업을 수행할 수 있어요.
이 예시는 pulsar-postgres-jdbc-sink sink 커넥터를 삭제해요.
bin/pulsar-admin sinks delete \
--tenant public \
--namespace default \
--name pulsar-postgres-jdbc-sink
팁:
pulsar-admin sinks delete options에 대한 더 자세한 내용은 Pulsar admin 문서를 참고하세요.
다음 메시지가 나타나면 sink 커넥터가 정상적으로 삭제된 거예요.
Deleted successfully
이 예시는 sink 커넥터의 상태를 다시 한번 확인해 봐요.
bin/pulsar-admin sinks get \
--tenant public \
--namespace default \
--name pulsar-postgres-jdbc-sink
결과를 보면 sink 커넥터가 존재하지 않는 것을 알 수 있어요.
HTTP 404 Not Found
Reason: Sink pulsar-postgres-jdbc-sink doesn't exist