Debezium 소스 커넥터

Debezium 소스 커넥터 (Debezium source connector)

Debezium source 커넥터는 MySQL 또는 PostgreSQL에서 메시지를 가져와 Pulsar 토픽에 저장하는 커넥터예요. 데이터베이스의 변경 데이터(CDC)를 Pulsar로 가져올 때 사용해요.

참고: 모든 Pulsar 커넥터는 download page에서 내려받을 수 있어요.

출처: 문서

본문

구성 (Configuration)

Debezium source 커넥터의 구성에는 다음과 같은 프로퍼티가 있어요.

이름 필수 기본값 설명
task.class true null Debezium에서 구현된 소스 태스크 클래스예요.
database.hostname true null 데이터베이스 서버의 주소예요.
database.port true null 데이터베이스 서버의 포트 번호예요.
database.user true null 필요한 권한을 가진 데이터베이스 사용자 이름이에요.
database.password true null 필요한 권한을 가진 데이터베이스 사용자의 비밀번호예요.
database.server.id true null 데이터베이스 클러스터 내에서 고유해야 하고 데이터베이스의 server-id 구성 프로퍼티와 비슷한 커넥터 식별자예요.
topic.prefix true null 데이터베이스 서버/클러스터의 논리적 이름이에요. 네임스페이스를 형성하며, 커넥터가 쓰는 모든 Kafka 토픽 이름, Kafka Connect 스키마 이름, 그리고 Avro 커넥터를 쓸 때 해당 Avro 스키마의 네임스페이스에 사용돼요.
database.include.list false null 이 서버가 호스팅하며 커넥터가 모니터링하는 모든 데이터베이스 목록이에요. 선택 사항이며, 모니터링에 포함하거나 제외할 데이터베이스와 테이블을 나열하는 다른 프로퍼티도 있어요.
key.converter true null 레코드 키를 변환하는 데 Kafka Connect가 제공하는 컨버터예요.
value.converter true null 레코드 값을 변환하는 데 Kafka Connect가 제공하는 컨버터예요.
database.history true null 데이터베이스 히스토리 클래스의 이름이에요.
database.history.pulsar.topic true null 커넥터가 DDL 문장을 쓰고 복구하는 데이터베이스 히스토리 토픽의 이름이에요. 참고: 이 토픽은 내부용으로만 쓰이며 컨슈머가 사용하면 안 돼요.
database.history.pulsar.service.url false null 히스토리 토픽용 Pulsar 클러스터 서비스 URL이에요. 참고: database.history.pulsar.service.url을 설정하지 않으면, 데이터베이스 히스토리 Pulsar 클라이언트는 client_auth_plugin과 client_auth_params 같은 소스 커넥터와 동일한 클라이언트 설정을 사용해요.
offset.storage.topic true null 커넥터가 성공적으로 완료한 마지막 커밋 오프셋을 기록해요.
json-with-envelope false false 페이로드만으로 구성된 메시지를 표현해요.
database.history.pulsar.reader.config false null 데이터베이스 스키마 히스토리 토픽용 리더(reader)의 구성으로, 키-값 쌍의 JSON 문자열 형태예요.
offset.storage.reader.config false null kafka 커넥터 오프셋 토픽용 리더의 구성으로, 키-값 쌍의 JSON 문자열 형태예요.

컨버터 옵션 (Converter Options)

  • org.apache.kafka.connect.json.JsonConverter

json-with-envelope 구성은 JsonConverter에만 유효해요. 기본값은 false이며, 컨슈머는 스키마 Schema.KeyValue(Schema.AUTO_CONSUME(), Schema.AUTO_CONSUME(), KeyValueEncodingType.SEPARATED)를 사용하고 메시지는 페이로드로만 구성돼요.

json-with-envelope 구성 값이 true이면, 컨슈머는 스키마 Schema.KeyValue(Schema.BYTES, Schema.BYTES를 사용하고 메시지는 스키마와 페이로드로 구성돼요.

  • org.apache.pulsar.kafka.shade.io.confluent.connect.avro.AvroConverter

사용자가 AvroConverter를 선택하면, pulsar 컨슈머는 스키마 Schema.KeyValue(Schema.AUTO_CONSUME(), Schema.AUTO_CONSUME(), KeyValueEncodingType.SEPARATED)를 사용해야 하며 메시지는 페이로드로 구성돼요.

MongoDB 구성 (MongoDB Configuration)

이름 필수 기본값 설명
mongodb.hosts true null 레플리카 셋에 있는 MongoDB 서버들의 호스트네임/포트 쌍('host' 또는 'host:port' 형식)의 쉼표 구분 목록이에요. 목록은 단일 호스트네임/포트 쌍을 포함해요. mongodb.members.auto.discover가 false로 설정되면 호스트/포트 쌍 앞에 레플리카 셋 이름이 붙어요(예: rs0/localhost:27017).
mongodb.name true null 이 커넥터가 모니터링하는 커넥터 및/또는 MongoDB 레플리카 셋 또는 공유 클러스터를 식별하는 고유 이름이에요. 각 서버는 최대 한 개의 Debezium 커넥터로 모니터링되어야 해요. 왜냐하면 이 서버 이름이 MongoDB 레플리카 셋 또는 클러스터에서 파생되는 모든 영속 Kafka 토픽에 접두사로 붙기 때문이에요.
mongodb.user true null MongoDB에 연결할 때 사용할 데이터베이스 사용자 이름이에요. MongoDB가 인증을 사용하도록 구성된 경우에만 필요해요.
mongodb.password true null MongoDB에 연결할 때 사용할 비밀번호예요. MongoDB가 인증을 사용하도록 구성된 경우에만 필요해요.
mongodb.task.id true null 각 레플리카 셋에 대해 별도의 태스크를 사용하려는 MongoDB 커넥터의 taskId예요.

메타데이터 토픽용 Reader 구성 커스터마이즈하기 (Customize the Reader config for the metadata topics)

Debezium 커넥터는 database.history.pulsar.reader.configoffset.storage.reader.config를 노출해 데이터베이스 스키마 히스토리 토픽과 Kafka 커넥터 오프셋 토픽의 리더(reader)를 구성할 수 있게 해줘요. 예를 들어 구독 이름과 다른 리더 구성을 설정하는 데 쓸 수 있어요. 사용 가능한 구성은 ReaderConfigurationData에서 찾을 수 있어요.

예를 들어 두 Reader 모두에 구독 이름을 구성하려면 다음 구성을 추가할 수 있어요.

  • JSON
{
  "configs": {
     "database.history.pulsar.reader.config": "{\"subscriptionName\":\"history-reader\"}",
     "offset.storage.reader.config": "{\"subscriptionName\":\"offset-reader\"}",
  }
}
  • YAML
configs:
   database.history.pulsar.reader.config: "{\"subscriptionName\":\"history-reader\"}"
   offset.storage.reader.config: "{\"subscriptionName\":\"offset-reader\"}"

MySQL 예시 (Example of MySQL)

Pulsar Debezium 커넥터를 사용하기 전에 구성 파일을 만들어야 해요.

구성 (Configuration)

다음 방법 중 하나로 구성 파일을 만들 수 있어요.

  • JSON
{
   "configs": {
      "database.hostname": "localhost",
      "database.port": "3306",
      "database.user": "debezium",
      "database.password": "dbz",
      "database.server.id": "184054",
      "topic.prefix": "dbserver1",
      "table.include.list": "inventory",
      "database.history": "org.apache.pulsar.io.debezium.PulsarDatabaseHistory",
      "database.history.pulsar.topic": "history-topic",
      "database.history.pulsar.service.url": "pulsar://127.0.0.1:6650",
      "key.converter": "org.apache.kafka.connect.storage.StringConverter",
      "value.converter": "org.apache.kafka.connect.storage.StringConverter",
      "offset.storage.topic": "offset-topic"
   }
}
  • YAML

debezium-mysql-source-config.yaml 파일을 만들고 아래 내용을 그 파일에 복사할 수 있어요.

tenant: "public"
namespace: "default"
name: "debezium-mysql-source"
topicName: "debezium-mysql-topic"
archive: "connectors/pulsar-io-debezium-mysql-5.0.0-M2.nar"
parallelism: 1
configs:
    ## config for mysql, docker image: debezium/example-mysql:0.8
    database.hostname: "localhost"
    database.port: "3306"
    database.user: "debezium"
    database.password: "dbz"
    database.server.id: "184054"
    topic.prefix: "dbserver1"
    database.include.list: "inventory"
    database.history: "org.apache.pulsar.io.debezium.PulsarDatabaseHistory"
    database.history.pulsar.topic: "history-topic"
    database.history.pulsar.service.url: "pulsar://127.0.0.1:6650"
    ## KEY_CONVERTER_CLASS_CONFIG, VALUE_CONVERTER_CLASS_CONFIG
    key.converter: "org.apache.kafka.connect.storage.StringConverter"
    value.converter: "org.apache.kafka.connect.storage.StringConverter"
    ## OFFSET_STORAGE_TOPIC_CONFIG
    offset.storage.topic: "offset-topic"

사용법 (Usage)

이 예시는 Pulsar Debezium 커넥터를 사용해 MySQL 테이블의 데이터를 변경하는 방법을 보여줘요.

  • Debezium이 변경을 캡처할 데이터베이스가 있는 MySQL 서버를 시작해요.
docker run -it --rm \
--name mysql \
-p 3306:3306 \
-e MYSQL_ROOT_PASSWORD=debezium \
-e MYSQL_USER=mysqluser \
-e MYSQL_PASSWORD=mysqlpw debezium/example-mysql:0.8
  • Pulsar 서비스를 로컬 standalone 모드로 시작해요.
bin/pulsar standalone
  • 다음 방법 중 하나로 Pulsar Debezium 커넥터를 로컬 실행 모드로 시작해요.

참고: Debezium 커넥터는 데이터를 다음 4가지 유형의 토픽에 저장해요.

  • 데이터베이스 메타데이터 메시지를 저장하는 database.server.name(데이터베이스 서버 이름)으로 이름이 지정된 토픽 하나. 예: public/default/database.server.name.
  • 데이터베이스 히스토리 정보를 저장하는 토픽(database.history.pulsar.topic) 하나. 커넥터는 이 토픽에 DDL 문장을 쓰고 복구해요.
  • 오프셋 메타데이터 메시지를 저장하는 토픽(offset.storage.topic) 하나. 커넥터는 마지막으로 성공적으로 커밋된 오프셋을 이 토픽에 저장해요.
  • 테이블별 토픽 하나. 커넥터는 테이블에서 발생하는 모든 연산의 변경 이벤트를 해당 테이블 전용의 단일 Pulsar 토픽에 써요.

브로커에서 자동 토픽 생성이 비활성화되어 있다면, 위 4가지 유형의 토픽을 수동으로 만들어야 해요.

  • 앞서 보여준 JSON 구성 파일을 사용해요.

NAR 파일이 connectors/pulsar-io-debezium-mysql-5.0.0-M2.nar 경로에 있는지 확인하세요.

bin/pulsar-admin source localrun \
    --archive $PWD/connectors/pulsar-io-debezium-mysql-5.0.0-M2.nar \
    --name debezium-mysql-source \
    --tenant public \
    --namespace default \
    --source-config '{"database.hostname": "localhost","database.port": "3306","database.user": "debezium","database.password": "dbz","database.server.id": "184054","topic.prefix": "dbserver1","database.include.list": "inventory","database.history": "org.apache.pulsar.io.debezium.PulsarDatabaseHistory","database.history.pulsar.topic": "history-topic","database.history.pulsar.service.url": "pulsar://127.0.0.1:6650","key.converter": "org.apache.kafka.connect.storage.StringConverter","value.converter": "org.apache.kafka.connect.storage.StringConverter","pulsar.service.url": "pulsar://127.0.0.1:6650","offset.storage.topic": "offset-topic"}'
  • 앞서 보여준 YAML 구성 파일을 사용해요.
bin/pulsar-admin source localrun \
   --source-config-file $PWD/debezium-mysql-source-config.yaml
  • inventory.products 테이블에 대한 sub-products 토픽을 구독해요.
bin/pulsar-client consume -s "sub-products" public/default/dbserver1.inventory.products -n 0
  • docker에서 MySQL 클라이언트를 시작해요.
docker run -it --rm \
    --name mysqlterm \
    --link mysql \
    --rm mysql:5.7 sh \
    -c 'exec mysql -h"$MYSQL_PORT_3306_TCP_ADDR" -P"$MYSQL_PORT_3306_TCP_PORT" -uroot -p"$MYSQL_ENV_MYSQL_ROOT_PASSWORD"'
  • MySQL 클라이언트가 나타나요. 연결 모드를 mysql_native_password로 변경해요.
mysql> show variables like "caching_sha2_password_auto_generate_rsa_keys";
+----------------------------------------------+-------+
| Variable_name                                | Value |
+----------------------------------------------+-------+
| caching_sha2_password_auto_generate_rsa_keys | ON    |
+----------------------------------------------+-------+
# If the value of "caching_sha2_password_auto_generate_rsa_keys" is ON, ensure the plugin of mysql.user is "mysql_native_password".
mysql> SELECT Host, User, plugin from mysql.user where user={user};
+-----------+------+-----------------------+
| Host      | User | plugin                |
+-----------+------+-----------------------+
| localhost | root | caching_sha2_password |
+-----------+------+-----------------------+
# If the plugin of mysql.user is is "caching_sha2_password", set it to "mysql_native_password".
alter user '{user}'@'{host}' identified with mysql_native_password by {password};
# Check the plugin of mysql.user.
mysql> SELECT Host, User, plugin from mysql.user where user={user};
+-----------+------+-----------------------+
| Host      | User | plugin                |
+-----------+------+-----------------------+
| localhost | root | mysql_native_password |
+-----------+------+-----------------------+

다음 명령들로 products 테이블의 데이터를 변경해요.

mysql> use inventory;
mysql> show tables;
mysql> SELECT * FROM  products;
mysql> UPDATE products SET name='1111111111' WHERE id=101;
mysql> UPDATE products SET name='1111111111' WHERE id=107;

구독 토픽의 터미널 창에서 데이터 변경이 sub-products 토픽에 유지된 것을 확인할 수 있어요.

PostgreSQL 예시 (Example of PostgreSQL)

Pulsar Debezium 커넥터를 사용하기 전에 구성 파일을 만들어야 해요.

구성 (Configuration)

  • JSON
{
    "database.hostname": "localhost",
    "database.port": "5432",
    "database.user": "postgres",
    "database.password": "changeme",
    "database.dbname": "postgres",
    "topic.prefix": "dbserver1",
    "plugin.name": "pgoutput",
    "schema.include.list": "public",
    "table.include.list": "public.users",
    "database.history.pulsar.service.url": "pulsar://127.0.0.1:6650"
}
  • YAML

debezium-postgres-source-config.yaml 파일을 만들고 아래 내용을 그 파일에 복사할 수 있어요.

tenant: "public"
namespace: "default"
name: "debezium-postgres-source"
topicName: "debezium-postgres-topic"
archive: "connectors/pulsar-io-debezium-postgres-5.0.0-M2.nar"
parallelism: 1
configs:
    ## config for postgres version 10+, official docker image: postgres:<10+>
    database.hostname: "localhost"
    database.port: "5432"
    database.user: "postgres"
    database.password: "changeme"
    database.dbname: "postgres"
    topic.prefix: "dbserver1"
    plugin.name: "pgoutput"
    schema.include.list: "public"
    table.include.list: "public.users"
    key.converter: "org.apache.kafka.connect.storage.StringConverter"
    value.converter: "org.apache.kafka.connect.storage.StringConverter"
    ## PULSAR_SERVICE_URL_CONFIG
    database.history.pulsar.service.url: "pulsar://127.0.0.1:6650"

사용법 (Usage)

이 예시는 Pulsar Debezium 커넥터를 사용해 PostgreSQL 테이블의 데이터를 변경하는 방법을 보여줘요.

  • Debezium이 변경을 캡처할 데이터베이스가 있는 PostgreSQL 서버를 시작해요.
docker run -d -it --rm \
    --name pulsar-postgres \
    -p 5432:5432 \
    -e POSTGRES_PASSWORD=changeme \
    postgres:13.3 -c wal_level=logical
  • Pulsar 서비스를 로컬 standalone 모드로 시작해요.
bin/pulsar standalone
  • 다음 방법 중 하나로 Pulsar Debezium 커넥터를 로컬 실행 모드로 시작해요.

NAR 파일이 connectors/pulsar-io-debezium-postgres-5.0.0-M2.nar 경로에 있는지 확인하세요.

  • 앞서 보여준 JSON 구성 파일을 사용해요.
bin/pulsar-admin source localrun \
    --archive $PWD/connectors/pulsar-io-debezium-postgres-5.0.0-M2.nar \
    --name debezium-postgres-source \
    --tenant public \
    --namespace default \
    --source-config '{"database.hostname": "localhost","database.port": "5432","database.user": "postgres","database.password": "changeme","database.dbname": "postgres","topic.prefix": "dbserver1","plugin.name": "pgoutput","schema.include.list": "public","table.include.list": "public.users","key.converter": "org.apache.kafka.connect.storage.StringConverter","value.converter": "org.apache.kafka.connect.storage.StringConverter","pulsar.service.url": "pulsar://127.0.0.1:6650"}'
  • 앞서 보여준 YAML 구성 파일을 사용해요.
bin/pulsar-admin source localrun  \
   --source-config-file $PWD/debezium-postgres-source-config.yaml
  • public.users 테이블에 대한 sub-users 토픽을 구독해요.
bin/pulsar-client consume -s "sub-users" public/default/dbserver1.public.users -n 0
  • docker에서 PostgreSQL 클라이언트를 시작해요.
docker exec -it pulsar-postgres /bin/bash
  • PostgreSQL 클라이언트가 나타나요. 다음 명령들로 users 테이블에 샘플 데이터를 만들어요.
psql -U postgres -h localhost -p 5432
Password for user postgres:
CREATE TABLE users(
  id BIGINT GENERATED ALWAYS AS IDENTITY, PRIMARY KEY(id),
  hash_firstname TEXT NOT NULL,
  hash_lastname TEXT NOT NULL,
  gender VARCHAR(6) NOT NULL CHECK (gender IN ('male', 'female')));
INSERT INTO users(hash_firstname, hash_lastname, gender)
  SELECT md5(RANDOM()::TEXT), md5(RANDOM()::TEXT), CASE WHEN RANDOM() < 0.5 THEN 'male' ELSE 'female' END FROM generate_series(1, 100);
postgres=# select * from users;
  id   |          hash_firstname          |          hash_lastname           | gender
-------+----------------------------------+----------------------------------+--------
     1 | 02bf7880eb489edc624ba637f5ab42bd | 3e742c2cc4217d8e3382cc251415b2fb | female
     2 | dd07064326bb9119189032316158f064 | 9c0e938f9eddbd5200ba348965afbc61 | male
     3 | 2c5316fdd9d6595c1cceb70eed12e80c | 8a93d7d8f9d76acfaaa625c82a03ea8b | female
     4 | 3dfa3b4f70d8cd2155567210e5043d2b | 32c156bc28f7f03ab5d28e2588a3dc19 | female
postgres=# UPDATE users SET hash_firstname='maxim' WHERE id=1;
UPDATE 1

구독 토픽의 터미널 창에서 다음 메시지를 받을 수 있어요.

----- got message -----
{"before":null,"after":{"id":1,"hash_firstname":"maxim","hash_lastname":"292113d30a3ccee0e19733dd7f88b258","gender":"male"},"source":{"version":"1.0.0.Final","connector":"postgresql","name":"foobar","ts_ms":1624045862644,"snapshot":"false","db":"postgres","schema":"public","table":"users","txId":595,"lsn":24419784,"xmin":null},"op":"u","ts_ms":1624045862648}

MongoDB 예시 (Example of MongoDB)

Pulsar Debezium 커넥터를 사용하기 전에 구성 파일을 만들어야 해요.

구성 (Configuration)

  • JSON
{
    "mongodb.hosts": "rs0/mongodb:27017",
    "mongodb.name": "dbserver1",
    "mongodb.user": "debezium",
    "mongodb.password": "dbz",
    "mongodb.task.id": "1",
    "database.include.list": "inventory",
    "database.history.pulsar.service.url": "pulsar://127.0.0.1:6650"
}
  • YAML

debezium-mongodb-source-config.yaml 파일을 만들고 아래 내용을 그 파일에 복사할 수 있어요.

tenant: "public"
namespace: "default"
name: "debezium-mongodb-source"
topicName: "debezium-mongodb-topic"
archive: "connectors/pulsar-io-debezium-mongodb-5.0.0-M2.nar"
parallelism: 1
configs:
    ## config for pg, docker image: debezium/example-mongodb:0.10
    mongodb.hosts: "rs0/mongodb:27017"
    mongodb.name: "dbserver1"
    mongodb.user: "debezium"
    mongodb.password: "dbz"
    mongodb.task.id: "1"
    database.include.list: "inventory"
    database.history.pulsar.service.url: "pulsar://127.0.0.1:6650"

사용법 (Usage)

이 예시는 Pulsar Debezium 커넥터를 사용해 MongoDB 데이터를 변경하는 방법을 보여줘요.

  • Debezium이 변경을 캡처할 데이터베이스가 있는 MongoDB 서버를 시작해요.
docker pull debezium/example-mongodb:0.10
docker run -d -it --rm --name pulsar-mongodb -e MONGODB_USER=mongodb -e MONGODB_PASSWORD=mongodb -p 27017:27017  debezium/example-mongodb:0.10

다음 명령들로 데이터를 초기화해요.

./usr/local/bin/init-inventory.sh

로컬 호스트가 컨테이너 네트워크에 접근할 수 없다면, /etc/hosts 파일을 업데이트해 127.0.0.1 f114527a95f 규칙을 추가할 수 있어요. f114527a95f는 컨테이너 ID이며 docker ps -a로 구할 수 있어요.

  • Pulsar 서비스를 로컬 standalone 모드로 시작해요.
bin/pulsar standalone
  • 다음 방법 중 하나로 Pulsar Debezium 커넥터를 로컬 실행 모드로 시작해요.

NAR 파일이 connectors/pulsar-io-mongodb-5.0.0-M2.nar 경로에 있는지 확인하세요.

  • 앞서 보여준 JSON 구성 파일을 사용해요.
bin/pulsar-admin source localrun \
    --archive $PWD/connectors/pulsar-io-debezium-mongodb-5.0.0-M2.nar \
    --name debezium-mongodb-source \
    --tenant public \
    --namespace default \
    --source-config '{"mongodb.hosts": "rs0/mongodb:27017","mongodb.name": "dbserver1","mongodb.user": "debezium","mongodb.password": "dbz","mongodb.task.id": "1","database.include.list": "inventory","database.history.pulsar.service.url": "pulsar://127.0.0.1:6650"}'
  • 앞서 보여준 YAML 구성 파일을 사용해요.
bin/pulsar-admin source localrun  \
--source-config-file $PWD/debezium-mongodb-source-config.yaml
  • inventory.products 테이블에 대한 sub-products 토픽을 구독해요.
bin/pulsar-client consume -s "sub-products" public/default/dbserver1.inventory.products -n 0
  • docker에서 MongoDB 클라이언트를 시작해요.
docker exec -it pulsar-mongodb /bin/bash
  • MongoDB 클라이언트가 나타나요. 다음 명령으로 products 도큐먼트의 무게를 업데이트해요.
mongo -u debezium -p dbz --authenticationDatabase admin localhost:27017/inventory
db.products.update({"_id":NumberLong(104)},{$set:{weight:1.25}})

구독 토픽의 터미널 창에서 변경 이벤트 메시지(JSON 구조)를 받을 수 있어요.

Oracle 예시 (Example of Oracle)

패키징 (Packaging)

Oracle 커넥터는 Oracle JDBC 드라이버를 포함하지 않으므로 커넥터와 함께 패키징해야 해요. 드라이버를 포함하지 않는 주요 이유는 다양한 버전과 Oracle 라이선스 때문이에요. Oracle DB 설치에 포함된 드라이버를 사용하거나 하나를 다운로드하는 것을 권장해요. 통합 테스트에는 커넥터 NAR 파일에 드라이버를 패키징하는 예시가 있어요.

구성 (Configuration)

Debezium은 LogMiner 또는 XStream API가 활성화된 Oracle DB를 요구해요. 이를 활성화하는 지원 옵션과 단계는 Oracle DB 버전마다 달라요. 문서에 설명되고 통합 테스트에서 사용된 단계는 사용 중인 Oracle DB의 버전/에디션에서 동작할 수도 있고 아닐 수도 있어요. 필요에 따라 Oracle DB 문서를 참고하세요.

다른 커넥터와 마찬가지로 JSON 또는 YAML로 커넥터를 구성할 수 있어요. YAML을 예시로 들면, 아래처럼 debezium-oracle-source-config.yaml 파일을 만들 수 있어요.

  • JSON
{
  "database.hostname": "localhost",
  "database.port": "1521",
  "database.user": "dbzuser",
  "database.password": "dbz",
  "database.dbname": "XE",
  "topic.prefix": "XE",
  "schema.exclude.list": "system,dbzuser",
  "snapshot.mode": "initial",
  "topic.namespace": "public/default",
  "task.class": "io.debezium.connector.oracle.OracleConnectorTask",
  "key.converter": "org.apache.kafka.connect.storage.StringConverter",
  "value.converter": "org.apache.kafka.connect.storage.StringConverter",
  "database.history": "org.apache.pulsar.io.debezium.PulsarDatabaseHistory",
  "database.tcpKeepAlive": "true",
  "decimal.handling.mode": "double",
  "database.history.pulsar.topic": "debezium-oracle-source-history-topic",
  "database.history.pulsar.service.url": "pulsar://127.0.0.1:6650"
}
  • YAML
tenant: "public"
namespace: "default"
name: "debezium-oracle-source"
topicName: "debezium-oracle-topic"
parallelism: 1
className: "org.apache.pulsar.io.debezium.oracle.DebeziumOracleSource"
database.dbname: "XE"
configs:
    database.hostname: "localhost"
    database.port: "1521"
    database.user: "dbzuser"
    database.password: "dbz"
    database.dbname: "XE"
    topic.prefix: "XE"
    schema.exclude.list: "system,dbzuser"
    snapshot.mode: "initial"
    topic.namespace: "public/default"
    task.class: "io.debezium.connector.oracle.OracleConnectorTask"
    key.converter: "org.apache.kafka.connect.storage.StringConverter"
    value.converter: "org.apache.kafka.connect.storage.StringConverter"
    database.history: "org.apache.pulsar.io.debezium.PulsarDatabaseHistory"
    database.tcpKeepAlive: "true"
    decimal.handling.mode: "double"
    database.history.pulsar.topic: "debezium-oracle-source-history-topic"
    database.history.pulsar.service.url: "pulsar://127.0.0.1:6650"

Debezium이 지원하는 전체 구성 프로퍼티 목록은 Debezium Connector for Oracle을 참고하세요.

Microsoft SQL 예시 (Example of Microsoft SQL)

구성 (Configuration)

Debezium은 CDC가 활성화된 SQL Server를 요구해요. 문서에 설명되고 통합 테스트에 사용된 단계들이에요. 더 자세한 내용은 Enable and disable change data capture in Microsoft SQL Server를 참고하세요.

다른 커넥터와 마찬가지로 JSON 또는 YAML로 커넥터를 구성할 수 있어요.

  • JSON
{
  "database.hostname": "localhost",
  "database.port": "1433",
  "database.user": "sa",
  "database.password": "MyP@ssw0rd!",
  "database.dbname": "MyTestDB",
  "topic.prefix": "mssql",
  "snapshot.mode": "schema_only",
  "topic.namespace": "public/default",
  "task.class": "io.debezium.connector.sqlserver.SqlServerConnectorTask",
  "key.converter": "org.apache.kafka.connect.storage.StringConverter",
  "value.converter": "org.apache.kafka.connect.storage.StringConverter",
  "database.history": "org.apache.pulsar.io.debezium.PulsarDatabaseHistory",
  "database.tcpKeepAlive": "true",
  "decimal.handling.mode": "double",
  "database.history.pulsar.topic": "debezium-mssql-source-history-topic",
  "database.history.pulsar.service.url": "pulsar://127.0.0.1:6650"
}
  • YAML
tenant: "public"
namespace: "default"
name: "debezium-mssql-source"
topicName: "debezium-mssql-topic"
parallelism: 1
className: "org.apache.pulsar.io.debezium.mssql.DebeziumMsSqlSource"
database.dbname: "mssql"
configs:
    database.hostname: "localhost"
    database.port: "1433"
    database.user: "sa"
    database.password: "MyP@ssw0rd!"
    database.dbname: "MyTestDB"
    topic.prefix: "mssql"
    snapshot.mode: "schema_only"
    topic.namespace: "public/default"
    task.class: "io.debezium.connector.sqlserver.SqlServerConnectorTask"
    key.converter: "org.apache.kafka.connect.storage.StringConverter"
    value.converter: "org.apache.kafka.connect.storage.StringConverter"
    database.history: "org.apache.pulsar.io.debezium.PulsarDatabaseHistory"
    database.tcpKeepAlive: "true"
    decimal.handling.mode: "double"
    database.history.pulsar.topic: "debezium-mssql-source-history-topic"
    database.history.pulsar.service.url: "pulsar://127.0.0.1:6650"

Debezium이 지원하는 전체 구성 프로퍼티 목록은 Debezium Connector for MS SQL을 참고하세요.

FAQ

Debezium postgres 커넥터가 스냅샷 생성 시 멈춤

#18 prio=5 os_prio=31 tid=0x00007fd83096f800 nid=0xa403 waiting on condition [0x000070000f534000]
    java.lang.Thread.State: WAITING (parking)
     at sun.misc.Unsafe.park(Native Method)
     ...

데이터 동기화 중 위와 같은 문제가 발생하면 이 글을 참고해 구성 파일에 다음 구성을 추가하세요.

max.queue.size=

더 알아보기 (Learn more)