압축 해제 후 생성된 디렉터리의 이름을 바꾸지 않고 plugin.path 로 지정할 디렉터리 바로 아래에 넣는다.
mkdir -p $KAFKA_HOME/plugins
cd /tmp
curl -LO https://repo1.maven.org/maven2/io/debezium/debezium-connector-mysql/3.6.3.Final/debezium-connector-mysql-3.6.3.Final-plugin.tar.gz
tar -xzf debezium-connector-mysql-3.6.3.Final-plugin.tar.gz -C $KAFKA_HOME/plugins/
ls $KAFKA_HOME/plugins/debezium-connector-mysql/
$KAFKA_HOME/config/connect-standalone.properties. 단일 노드 시험용이다. 운영은 connect-distributed.properties 로 분산 모드를 쓰고 REST(:8083/connectors)로 커넥터를 등록한다.
bootstrap.servers=192.168.103.113:9092
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=true
value.converter.schemas.enable=true
offset.storage.file.filename=/home/kafka/kafka/data/connect.offsets
# Flush much faster than normal, which is useful for testing/debugging
offset.flush.interval.ms=10000
# 플러그인 경로. 이 경로 바로 아래에 debezium-connector-mysql 디렉터리가 있어야 한다
plugin.path=/home/kafka/kafka/plugins
$KAFKA_HOME/config/connect-mysql.properties. Debezium 2.0 부터 database.server.name 은 topic.prefix, database.history.kafka.* 는 schema.history.internal.kafka.* 다[1].
name=mysql-cdc
connector.class=io.debezium.connector.mysql.MySqlConnector
tasks.max=1
database.hostname=hdfs04.encore.rnd
database.port=3306
database.user=debezium
database.password=${DEBEZIUM_PASSWORD}
# MySQL 의 server-id 와 겹치지 않는 값
database.server.id=184054
# 토픽 이름 접두사. <topic.prefix>.<db>.<table>
topic.prefix=hdfs04
# 잡을 데이터베이스·테이블
database.include.list=test
#table.include.list=test.test1
# 스키마 이력 토픽 (커넥터 내부용)
schema.history.internal.kafka.bootstrap.servers=192.168.103.113:9092
schema.history.internal.kafka.topic=schemahistory.hdfs04
위 버전과 저장소는 오래된 것이라 저장소가 존재하지 않을 수 있다.
name=mysql-cdc
connector.class=io.debezium.connector.mysql.MySqlConnector
tasks.max=1
database.hostname=hdfs04.encore.rnd
database.port=3306
database.user=debezium
database.password=${REDACTED}
database.server.id=1
database.server.name=hdfs04.encore.rnd
database.history.kafka.bootstrap.servers=192.168.103.113:9092
database.history.kafka.topic=testcdc
cd $KAFKA_HOME
bin/connect-standalone.sh config/connect-standalone.properties config/connect-mysql.properties
standalone 구동일 경우 다음과 같이 프로세스가 뜬다.
[kafka@zookeeper-kafka kafka]$ jps -l
38497 kafka.Kafka
87046 org.apache.kafka.connect.cli.ConnectStandalone
86267 org.apache.kafka.tools.ConsoleConsumer
토픽이 생겼는지 본다.
bin/kafka-topics.sh --bootstrap-server 192.168.103.113:9092 --list | grep hdfs04
bin/kafka-console-consumer.sh --bootstrap-server 192.168.103.113:9092 --topic hdfs04.test.test1 --from-beginning
test.test1 에 insert 한 뒤 delete 했을 때의 이벤트다(1.2.1.Final 기록). payload.op 가 c(create) · u(update) · d(delete) · r(스냅샷 read) 이고, delete 뒤에는 키만 있고 값이 null 인 tombstone 레코드가 한 건 더 온다.
{"schema":{"type":"struct","fields":[{"type":"struct","fields":[{"type":"string","optional":true,"field":"tcol1"}],"optional":true,"name":"hdfs04.encore.rnd.test.test1.Value","field":"before"},{"type":"struct","fields":[{"type":"string","optional":true,"field":"tcol1"}],"optional":true,"name":"hdfs04.encore.rnd.test.test1.Value","field":"after"},{"type":"struct","fields":[{"type":"string","optional":false,"field":"version"},{"type":"string","optional":false,"field":"connector"},{"type":"string","optional":false,"field":"name"},{"type":"int64","optional":false,"field":"ts_ms"},{"type":"string","optional":true,"name":"io.debezium.data.Enum","version":1,"parameters":{"allowed":"true,last,false"},"default":"false","field":"snapshot"},{"type":"string","optional":false,"field":"db"},{"type":"string","optional":true,"field":"table"},{"type":"int64","optional":false,"field":"server_id"},{"type":"string","optional":true,"field":"gtid"},{"type":"string","optional":false,"field":"file"},{"type":"int64","optional":false,"field":"pos"},{"type":"int32","optional":false,"field":"row"},{"type":"int64","optional":true,"field":"thread"},{"type":"string","optional":true,"field":"query"}],"optional":false,"name":"io.debezium.connector.mysql.Source","field":"source"},{"type":"string","optional":false,"field":"op"},{"type":"int64","optional":true,"field":"ts_ms"},{"type":"struct","fields":[{"type":"string","optional":false,"field":"id"},{"type":"int64","optional":false,"field":"total_order"},{"type":"int64","optional":false,"field":"data_collection_order"}],"optional":true,"field":"transaction"}],"optional":false,"name":"hdfs04.encore.rnd.test.test1.Envelope"},"payload":{"before":null,"after":{"tcol1":"1"},"source":{"version":"1.2.1.Final","connector":"mysql","name":"hdfs04.encore.rnd","ts_ms":1597280219000,"snapshot":"false","db":"test","table":"test1","server_id":1,"gtid":null,"file":"bin.000003","pos":4862,"row":0,"thread":63,"query":null},"op":"c","ts_ms":1597281006714,"transaction":null}}
{"schema":{"type":"struct","fields":[{"type":"struct","fields":[{"type":"string","optional":true,"field":"tcol1"}],"optional":true,"name":"hdfs04.encore.rnd.test.test1.Value","field":"before"},{"type":"struct","fields":[{"type":"string","optional":true,"field":"tcol1"}],"optional":true,"name":"hdfs04.encore.rnd.test.test1.Value","field":"after"},{"type":"struct","fields":[{"type":"string","optional":false,"field":"version"},{"type":"string","optional":false,"field":"connector"},{"type":"string","optional":false,"field":"name"},{"type":"int64","optional":false,"field":"ts_ms"},{"type":"string","optional":true,"name":"io.debezium.data.Enum","version":1,"parameters":{"allowed":"true,last,false"},"default":"false","field":"snapshot"},{"type":"string","optional":false,"field":"db"},{"type":"string","optional":true,"field":"table"},{"type":"int64","optional":false,"field":"server_id"},{"type":"string","optional":true,"field":"gtid"},{"type":"string","optional":false,"field":"file"},{"type":"int64","optional":false,"field":"pos"},{"type":"int32","optional":false,"field":"row"},{"type":"int64","optional":true,"field":"thread"},{"type":"string","optional":true,"field":"query"}],"optional":false,"name":"io.debezium.connector.mysql.Source","field":"source"},{"type":"string","optional":false,"field":"op"},{"type":"int64","optional":true,"field":"ts_ms"},{"type":"struct","fields":[{"type":"string","optional":false,"field":"id"},{"type":"int64","optional":false,"field":"total_order"},{"type":"int64","optional":false,"field":"data_collection_order"}],"optional":true,"field":"transaction"}],"optional":false,"name":"hdfs04.encore.rnd.test.test1.Envelope"},"payload":{"before":{"tcol1":"1"},"after":null,"source":{"version":"1.2.1.Final","connector":"mysql","name":"hdfs04.encore.rnd","ts_ms":1597280252000,"snapshot":"false","db":"test","table":"test1","server_id":1,"gtid":null,"file":"bin.000003","pos":5036,"row":0,"thread":63,"query":null},"op":"d","ts_ms":1597281040257,"transaction":null}}
null
Debezium connector for MySQL — Connector properties — 2026-09-20 확인. https://debezium.io/documentation/reference/stable/connectors/mysql.html ↩︎