Canal source 커넥터
Canal source 커넥터
Canal source 커넥터는 MySQL에서 메시지를 가져와 Pulsar 토픽으로 보내요. Alibaba의 Canal을 사용해 MySQL binlog를 실시간으로 구독해서 Pulsar에 흘려보내는 구조예요.
모든 Pulsar 커넥터는 다운로드 페이지에서 내려받을 수 있어요. 이 문서에서는 커넥터 설정 속성과 함께, MySQL→Canal→Pulsar를 연결하는 전체 사용 예시를 단계별로 보여줍니다.
출처: 문서
본문
알아두기: 모든 Pulsar 커넥터는 다운로드 페이지에서 다운로드할 수 있어요.
Canal source 커넥터는 MySQL에서 메시지를 가져와 Pulsar 토픽으로 보내요.
구성 (Configuration)
Canal source 커넥터의 구성은 다음 속성을 가져요.
속성 (Property)
| 이름 | 필수 | 기본값 | 설명 |
|---|---|---|---|
| username | true | 없음 | Canal 서버 계정 (MySQL 아님). |
| password | true | 없음 | Canal 서버 비밀번호 (MySQL 아님). |
| destination | true | 없음 | Canal source 커넥터가 연결하는 소스 목적지(destination). |
| singleHostname | false | 없음 | Canal 서버 주소. |
| singlePort | false | 없음 | Canal 서버 포트. |
| cluster | true | false | Canal 서버 구성에 따른 클러스터 모드 활성화 여부. |
| zkServers | true | 없음 | Canal source 커넥터가 실제 데이터베이스 호스트를 알아내기 위해 통신하는 Zookeeper의 주소와 포트. |
| batchSize | false | 1000 | Canal에서 가져올 배치 크기. |
cluster 속성은 다음과 같이 동작해요.
true: 클러스터 모드.true로 설정하면zkServers와 통신해 실제 데이터베이스 호스트를 알아내요.false: standalone 모드.false로 설정하면singleHostname과singlePort로 지정된 데이터베이스에 연결해요.
예시 (Example)
Canal 커넥터를 사용하기 전에 다음 방법 중 하나로 구성 파일을 만들 수 있어요.
JSON:
{
"zkServers": "127.0.0.1:2181",
"batchSize": "5120",
"destination": "example",
"username": "",
"password": "",
"cluster": false,
"singleHostname": "127.0.0.1",
"singlePort": "11111",
}
YAML:
YAML 파일을 만들고 아래 내용을 복사해 넣어요.
configs:
zkServers: "127.0.0.1:2181"
batchSize: 5120
destination: "example"
username: ""
password: ""
cluster: false
singleHostname: "127.0.0.1"
singlePort: 11111
사용 예시 (Usage)
위와 같은 구성 파일로 MySQL 데이터를 저장하는 예시예요.
- MySQL 서버를 시작해요.
docker pull mysql:5.7
docker run -d -it --rm --name pulsar-mysql -p 3306:3306 -e MYSQL_ROOT_PASSWORD=canal -e MYSQL_USER=mysqluser -e MYSQL_PASSWORD=mysqlpw mysql:5.7
- 구성 파일
mysqld.cnf를 만들어요.
[mysqld]
pid-file = /var/run/mysqld/mysqld.pid
socket = /var/run/mysqld/mysqld.sock
datadir = /var/lib/mysql
#log-error = /var/log/mysql/error.log
# By default we only accept connections from localhost
#bind-address = 127.0.0.1
# Disabling symbolic-links is recommended to prevent assorted security risks
symbolic-links=0
log-bin=mysql-bin
binlog-format=ROW
server_id=1
- 구성 파일
mysqld.cnf를 MySQL 서버에 복사해요.
docker cp mysqld.cnf pulsar-mysql:/etc/mysql/mysql.conf.d/
- MySQL 서버를 재시작해요.
docker restart pulsar-mysql
- MySQL 서버에 테스트 데이터베이스를 만들어요.
docker exec -it pulsar-mysql /bin/bash
mysql -h 127.0.0.1 -uroot -pcanal -e 'create database test;'
- Canal 서버를 시작하고 MySQL 서버에 연결해요.
docker pull canal/canal-server:v1.1.2
docker run -d -it --link pulsar-mysql -e canal.auto.scan=false -e canal.destinations=test -e canal.instance.master.address=pulsar-mysql:3306 -e canal.instance.dbUsername=root -e canal.instance.dbPassword=canal -e canal.instance.connectionCharset=UTF-8 -e canal.instance.tsdb.enable=true -e canal.instance.gtidon=false --name=pulsar-canal-server -p 8000:8000 -p 2222:2222 -p 11111:11111 -p 11112:11112 -m 4096m canal/canal-server:v1.1.2
- Pulsar standalone을 시작해요.
docker pull apachepulsar/pulsar:5.0.0-M2
docker run --user 0 -d -it --link pulsar-canal-server -p 6650:6650 -p 8080:8080 -v $PWD/data:/pulsar/data --name pulsar-standalone apachepulsar/pulsar:5.0.0-M2 bin/pulsar standalone
- 구성 파일
canal-mysql-source-config.yaml을 수정해요.
configs:
zkServers: ""
batchSize: "5120"
destination: "test"
username: ""
password: ""
cluster: false
singleHostname: "pulsar-canal-server"
singlePort: "11111"
- 컨슈머 파일
pulsar-client.py를 만들어요.
import pulsar
client = pulsar.Client('pulsar://localhost:6650')
consumer = client.subscribe('my-topic',
subscription_name='my-sub')
while True:
msg = consumer.receive()
print("Received message: '%s'" % msg.data())
consumer.acknowledge(msg)
client.close()
- 구성 파일
canal-mysql-source-config.yaml과 컨슈머 파일pulsar-client.py를 Pulsar 서버에 복사해요.
docker cp canal-mysql-source-config.yaml pulsar-standalone:/pulsar/conf/
docker cp pulsar-client.py pulsar-standalone:/pulsar/
- Canal 커넥터를 다운로드하고 시작해요.
docker exec -it pulsar-standalone /bin/bash
curl -LO --output-dir connectors "https://www.apache.org/dyn/closer.lua/pulsar/pulsar-5.0.0-M2/connectors/pulsar-io-canal-5.0.0-M2.nar?action=download"
./bin/pulsar-admin source localrun \
--archive $PWD/connectors/pulsar-io-canal-5.0.0-M2.nar \
--classname org.apache.pulsar.io.canal.CanalStringSource \
--tenant public \
--namespace default \
--name canal \
--destination-topic-name my-topic \
--source-config-file /pulsar/conf/canal-mysql-source-config.yaml \
--parallelism 1
- MySQL에서 데이터를 소비해요.
docker exec -it pulsar-standalone /bin/bash
python pulsar-client.py
- 다른 창을 열어 MySQL 서버에 로그인해요.
docker exec -it pulsar-mysql /bin/bash
mysql -h 127.0.0.1 -uroot -pcanal
- MySQL 서버에 테이블을 만들고 데이터를 삽입·삭제·갱신해요.
mysql> use test;
mysql> show tables;
mysql> CREATE TABLE IF NOT EXISTS `test_table`(`test_id` INT UNSIGNED AUTO_INCREMENT,`test_title` VARCHAR(100) NOT NULL,
`test_author` VARCHAR(40) NOT NULL,
`test_date` DATE,PRIMARY KEY ( `test_id` ))ENGINE=InnoDB DEFAULT CHARSET=utf8;
mysql> INSERT INTO test_table (test_title, test_author, test_date) VALUES("a", "b", NOW());
mysql> UPDATE test_table SET test_title='c' WHERE test_title='a';
mysql> DELETE FROM test_table WHERE test_title='c';
더 알아보기 (Learn more)
- source 커넥터 실행 방법은 소스·sink 문서를 참고해요.
- 다른 커넥터 목록은 Pulsar IO 커넥터 문서에서 확인할 수 있어요.
- Canal 프로젝트가 궁금하다면 Alibaba Canal 문서를 참고해요.