本文主要是说明canal的部署和同步数据到es中
环境准备
这里需要将环境部署成以下的版本
版本
- canal-server 1.1.15
- canal-adapter 1.1.15
- mysql 5.7
- es 7.10.1
- kibana 7.10.1
开启mysql binlog,这里自行搜索一下怎么修改mysql配置
show variables like 'log_bin';
显示log_bin的值为ON则已经开启
建从库权限的账号参考语句
CREATE USER canal IDENTIFIED BY '123456';
GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'canal'@'%';
FLUSH PRIVILEGES;
创建库/表
CREATE DATABASE gordon_db CHARACTER SET = utf8 COLLATE = utf8_unicode_ci;
use gordon_db;
CREATE TABLE `dw_dim_mohurd_apt` (
`id` int(10),
`certno` varchar(32),
`name` varchar(128)
);
在kibana Dev Tools中执行
PUT mohurd_apt
{
"mappings" : {
"properties" : {
"certno" : {
"type" : "text",
"fields" : {
"keyword" : {
"type" : "keyword",
"ignore_above" : 256
}
}
},
"id" : {
"type" : "long"
},
"name" : {
"type" : "text",
"fields" : {
"keyword" : {
"type" : "keyword",
"ignore_above" : 256
}
}
}
}
}
}
安装
canal-server
依次执行以下命令下载,解压和修改配置文件
wget https://hub.fastgit.org/alibaba/canal/releases/download/canal-1.1.5/canal.deployer-1.1.5.tar.gz
tar -zxvf canal.deployer-1.1.5.tar.gz -C /app/canal-server
cd /app/canal-server
修改conf/example/instance.properties
#################################################
## mysql serverId , v1.0.26+ will autoGen
# canal.instance.mysql.slaveId=0
# enable gtid use true/false
canal.instance.gtidon=false
# position info
canal.instance.master.address=192.168.11.60:3306 #修改数据库ip端口
canal.instance.master.journal.name=
canal.instance.master.position=
canal.instance.master.timestamp=
canal.instance.master.gtid=
# rds oss binlog
canal.instance.rds.accesskey=
canal.instance.rds.secretkey=
canal.instance.rds.instanceId=
# table meta tsdb info
canal.instance.tsdb.enable=true
#canal.instance.tsdb.url=jdbc:mysql://127.0.0.1:3306/canal_tsdb
#canal.instance.tsdb.dbUsername=canal
#canal.instance.tsdb.dbPassword=canal
#canal.instance.standby.address =
#canal.instance.standby.journal.name =
#canal.instance.standby.position =
#canal.instance.standby.timestamp =
#canal.instance.standby.gtid=
# username/password
canal.instance.dbUsername=canal
canal.instance.dbPassword=123456
canal.instance.connectionCharset = UTF-8
# enable druid Decrypt database password
canal.instance.enableDruid=false
#canal.instance.pwdPublicKey=MFwwDQYJKoZIhvcNAQEBBQADSwAwSAJBALK4BUxdDltRRE5/zXpVEVPUgunvscYFtEip3pmLlhrWpacX7y7GCMo2/JM6LeHmiiNdH1FWgGCpUfircSwlWKUCAwEAAQ==
# table regex
canal.instance.filter.regex=.\\..*
# table black regex
canal.instance.filter.black.regex=mysql\\.slave_.*
# table field filter(format: schema1.tableName1:field1/field2,schema2.tableName2:field1/field2)
#canal.instance.filter.field=test1.t_product:id/subject/keywords,test2.t_company:id/name/contact/ch
# table field black filter(format: schema1.tableName1:field1/field2,schema2.tableName2:field1/field2)
#canal.instance.filter.black.field=test1.t_product:subject/product_image,test2.t_company:id/name/contact/ch
# mq config
canal.mq.topic=example
# dynamic topic route by schema or table regex
#canal.mq.dynamicTopic=mytest1.user,mytest2\\..*,.*\\..*
canal.mq.partition=0
# hash partition config
#canal.mq.partitionsNum=3
#canal.mq.partitionHash=test.table:id^name,.*\\..*
#canal.mq.dynamicTopicPartitionNum=test.*:4,mycanal:6
#################################################
启动canal-server
sh bin/startup.sh
canal-adapter
依次运行下面命令
wget https://hub.fastgit.org/alibaba/canal/releases/download/canal-1.1.5/canal.adapter-1.1.5.tar.gz
tar -zxvf canal.adapter-1.1.5.tar.gz -C /app/canal-adapter
cd /app/canal-adapter
修改配置conf/application.yml
server:
port: 8081
spring:
jackson:
date-format: yyyy-MM-dd HH:mm:ss
time-zone: GMT+8
default-property-inclusion: non_null
canal.conf:
mode: tcp #tcp kafka rocketMQ rabbitMQ
flatMessage: true
zookeeperHosts:
syncBatchSize: 1000
retries: 0
timeout:
accessKey:
secretKey:
consumerProperties:
# canal tcp consumer
canal.tcp.server.host: 192.168.11.30:11111
canal.tcp.zookeeper.hosts:
canal.tcp.batch.size: 500
canal.tcp.username:
canal.tcp.password:
# kafka consumer
kafka.bootstrap.servers: 127.0.0.1:9092
kafka.enable.auto.commit: false
kafka.auto.commit.interval.ms: 1000
kafka.auto.offset.reset: latest
kafka.request.timeout.ms: 40000
kafka.session.timeout.ms: 30000
kafka.isolation.level: read_committed
kafka.max.poll.records: 1000
# rocketMQ consumer
rocketmq.namespace:
rocketmq.namesrv.addr: 127.0.0.1:9876
rocketmq.batch.size: 1000
rocketmq.enable.message.trace: false
rocketmq.customized.trace.topic:
rocketmq.access.channel:
rocketmq.subscribe.filter:
# rabbitMQ consumer
rabbitmq.host:
rabbitmq.virtual.host:
rabbitmq.username:
rabbitmq.password:
rabbitmq.resource.ownerId:
# srcDataSources:
# defaultDS:
# url: jdbc:mysql://127.0.0.1:3306/mytest?useUnicode=true
# username: root
# password: 121212
srcDataSources:
defaultDS:
url: jdbc:mysql://192.168.11.60:3306/comp?useUnicode=true
username: canal
password: 123456
maxActive: 20
canalAdapters:
- instance: example # canal instance Name or mq topic name
groups:
- groupId: g1
outerAdapters:
- name: es7
key: eskey #这里是全量导入时的taskid
hosts: 192.168.11.73:9200
properties:
mode: rest # restful 模式
cluster.name: elasticsearch-cluster # 集群名字
修改 conf/es7/mohurd_apt.yml
dataSourceKey: defaultDS #和上面canal.conf.srcDataSources.defaultDS要一样
outerAdapterKey: eskey #和上面canal.conf.canalAdapters.instance.groups.outerAdapters.key要一样
destination: example #和上面canal.conf.canalAdapters.instance要一样
groupId: g1
esMapping:
_index: mohurd_apt #索引名称
_id: _id
sql: "select id as _id,id,certno,name
from dw_dim_mohurd_apt
" # 查询的sql返回的结构和es的filed对应
commitBatch: 3000
启动 /bin/startup.sh
全量同步
curl http://192.168.10.69:8081/etl/es7/eskey/mohurd_apt.yml -X POST
使用docker部署canal-server,canal-adapter
docker run --name canal-server -p 11111:11111 -m 1024m
-v /opt/docker/canal-server/conf:/home/admin/canal-server/conf
-v /opt/docker/canal-server/logs:/home/admin/canal-server/logs
-d canal/canal-server:v1.1.5
docker run --name canal-adapter -p 8081:8081 -v /opt/docker/canal-adapter/conf:/opt/canal-adapter/conf -v /opt/docker/canal-adapter/logs:/opt/canal-adapter/logs -d slpcat/canal-adapter:v1.1.5
然后对应修改宿主文件下的conf。